211 lines
5.9 KiB
Go
211 lines
5.9 KiB
Go
package main
|
|
|
|
// The asker on the bus: its asks in the controller's bucket `asked`, its asks and cancels published on the
|
|
// seat under the controller's name, the verbs a warrant chooses called with the controller's grant, and the
|
|
// router's record of its asks read under its name (novox/hq ADR 0259).
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
"github.com/nats-io/nats.go/jetstream"
|
|
|
|
"git.novox.be/novox/mesh-sdk/go/asks"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"github.com/novox/mesh-controller/internal/catalogue"
|
|
"github.com/novox/mesh-controller/internal/conditions"
|
|
"github.com/novox/mesh-controller/internal/inventory"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// askerFrom is the serving controller's asker; nil in any other process.
|
|
var askerFrom *asker
|
|
|
|
// askWithin is how long a verb a warrant chose is given to answer.
|
|
const askWithin = time.Minute
|
|
|
|
type busAsked struct{ conn *nats.Conn }
|
|
|
|
func (b busAsked) kv(ctx context.Context) (jetstream.KeyValue, error) {
|
|
js, err := jetstream.New(b.conn)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return js.KeyValue(ctx, broker.AskedBucket)
|
|
}
|
|
|
|
func (b busAsked) Get(ctx context.Context, id string) (*asked, error) {
|
|
kv, err := b.kv(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
e, err := kv.Get(ctx, id)
|
|
if errors.Is(err, jetstream.ErrKeyNotFound) {
|
|
return nil, nil
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var r asked
|
|
return &r, json.Unmarshal(e.Value(), &r)
|
|
}
|
|
|
|
func (b busAsked) Put(ctx context.Context, r asked) error {
|
|
kv, err := b.kv(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
body, err := json.Marshal(r)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
_, err = kv.Put(ctx, r.ID, body)
|
|
return err
|
|
}
|
|
|
|
func (b busAsked) All(ctx context.Context) ([]asked, error) {
|
|
kv, err := b.kv(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
lister, err := kv.ListKeys(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer func() { _ = lister.Stop() }()
|
|
var out []asked
|
|
for k := range lister.Keys() {
|
|
e, err := kv.Get(ctx, k)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
var r asked
|
|
if json.Unmarshal(e.Value(), &r) == nil {
|
|
out = append(out, r)
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// callAction performs an action's verb as the controller, through the grant that names it.
|
|
func callAction(conn *nats.Conn) func(ctx context.Context, a conditions.Action, args map[string]string) error {
|
|
return func(ctx context.Context, a conditions.Action, args map[string]string) error {
|
|
seat, verb, ok := strings.Cut(a.Verb, ".")
|
|
if !ok {
|
|
return fmt.Errorf("%q names no seat and verb", a.Verb)
|
|
}
|
|
body := map[string]any{}
|
|
for k, v := range args {
|
|
body[k] = v
|
|
}
|
|
if seat == catalogue.DeliverySeat {
|
|
_, err := askDeliveryOwner(ctx, conn, verb, body)
|
|
return err
|
|
}
|
|
granted := false
|
|
for _, v := range broker.VerbsTheControllerActsOnAWarrant {
|
|
granted = granted || (v.Seat == seat && v.Verb == verb)
|
|
}
|
|
if !granted {
|
|
return fmt.Errorf("%s: %w", a.Verb, errNotGranted)
|
|
}
|
|
var answer link.Answer
|
|
var err error
|
|
if a.Machine != "" {
|
|
answer, err = link.AskSeatTool(ctx, conn, seat, verb, a.Machine, body, askWithin)
|
|
} else {
|
|
answer, err = link.AskMeshSeatTool(ctx, conn, seat, verb, body, askWithin)
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if answer.Error != "" {
|
|
return fmt.Errorf("%s refused: %s", a.Verb, answer.Error)
|
|
}
|
|
return nil
|
|
}
|
|
}
|
|
|
|
// routerRecordOf reads the router's record of one of the controller's asks, under its name, and answers
|
|
// how it ended when it did: the bucket is the one the asks seat's declarer names as its records.
|
|
func routerRecordOf(conn *nats.Conn, inv *inventory.Inventory) func(ctx context.Context, id string) (*asks.Warrant, error) {
|
|
return func(ctx context.Context, id string) (*asks.Warrant, error) {
|
|
bucket, err := asksRecords(ctx, inv)
|
|
if err != nil || bucket == "" {
|
|
return nil, err
|
|
}
|
|
reply, err := conn.RequestWithContext(ctx, "$JS.API.DIRECT.GET.KV_"+bucket+".$KV."+bucket+"."+askerName+"."+id, nil)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if reply.Header.Get("Status") != "" {
|
|
return nil, nil // none, or not readable: the event says it
|
|
}
|
|
var rec struct {
|
|
State string `json:"state"`
|
|
Warrant *asks.Warrant `json:"warrant"`
|
|
}
|
|
if json.Unmarshal(reply.Data, &rec) != nil || rec.State == "open" || rec.Warrant == nil {
|
|
return nil, nil
|
|
}
|
|
return rec.Warrant, nil
|
|
}
|
|
}
|
|
|
|
// asksRecords is the bucket the asks seat's declarer keeps its record of asks in.
|
|
func asksRecords(ctx context.Context, inv *inventory.Inventory) (string, error) {
|
|
declared, err := inv.Catalogue(ctx)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
for _, m := range declared {
|
|
for _, s := range m.DefinesSeats {
|
|
if s.Name == broker.AsksSeat && len(s.Records) > 0 {
|
|
return broker.BucketName(m.Module, s.Records[0]), nil
|
|
}
|
|
}
|
|
}
|
|
return "", nil
|
|
}
|
|
|
|
// startAsking makes the serving controller's asker and hands it the router's words.
|
|
func startAsking(ctx context.Context, open *stores, server *link.Server, conn *nats.Conn, keeper *conditions.Keeper) {
|
|
js, err := jetstream.New(conn)
|
|
if err != nil {
|
|
fmt.Printf("the operator cannot be asked: %v\n", err)
|
|
return
|
|
}
|
|
a := &asker{
|
|
open: keeper.Open,
|
|
silence: func(ctx context.Context, key string, d time.Duration, by, why string) error {
|
|
_, err := keeper.Silence(ctx, key, d, by, why)
|
|
return err
|
|
},
|
|
store: busAsked{conn: conn},
|
|
publish: func(ctx context.Context, subject string, body []byte, id string) error {
|
|
_, err := js.Publish(ctx, subject, body, jetstream.WithMsgID(id))
|
|
return err
|
|
},
|
|
call: callAction(conn),
|
|
record: func(ctx context.Context, act link.HandAct) error {
|
|
_, err := link.RecordHandAct(ctx, conn, act)
|
|
return err
|
|
},
|
|
routerRecord: routerRecordOf(conn, open.inventory),
|
|
now: time.Now,
|
|
logf: func(format string, args ...any) { fmt.Printf(format+"\n", args...) },
|
|
}
|
|
if err := server.Decides(a); err != nil {
|
|
fmt.Printf("the operator's answers cannot be heard, so nothing is asked: %v\n", err)
|
|
return
|
|
}
|
|
askerFrom = a
|
|
go a.keep(ctx)
|
|
}
|