Files
mesh-controller/cmd/mesh-controller/asker_wire.go
T
jschoubben 0c8c9ffae9
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
Ask the operator only once the bus holds the controller's grant to ask (hq issue 353)
The grant is composed from the router's assignment and reaches the bus when its machine is next
pushed. Between assign and push the record said a router was here and the bus refused every ask
(seven refusals on 2026-10-09, 17:54 to 17:56). The asker now judges, as a push does, whether the
bus's machine was last sent the user list composed now; while it was not, nothing is published, it
is said once, and the conditions that need the operator are raised as undelivered, naming the push
that carries the grant.
2026-10-09 18:13:58 +02:00

338 lines
10 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"
"sort"
"strings"
"sync"
"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)
}
// Create keeps a new ask under its id, and only where none is kept: never over another.
func (b busAsked) Create(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.Create(ctx, r.ID, body)
return err
}
// Change applies change to the ask kept under id by compare-and-set on its key's revision (the review of
// 2026-10-09, L2): read, changed, and written only over the revision read; when another write came between,
// read again and asked again, at most askChangeTries times. change says whether to write at all.
func (b busAsked) Change(ctx context.Context, id string, change func(*asked) bool) (bool, error) {
kv, err := b.kv(ctx)
if err != nil {
return false, err
}
for try := 0; try < askChangeTries; try++ {
e, err := kv.Get(ctx, id)
if errors.Is(err, jetstream.ErrKeyNotFound) {
return false, nil
}
if err != nil {
return false, err
}
var r asked
if err := json.Unmarshal(e.Value(), &r); err != nil {
return false, err
}
if !change(&r) {
return false, nil
}
body, err := json.Marshal(r)
if err != nil {
return false, err
}
if _, err := kv.Update(ctx, id, body, e.Revision()); err != nil {
var api *jetstream.APIError
if errors.Is(err, jetstream.ErrKeyExists) || (errors.As(err, &api) && api.ErrorCode == jetstream.JSErrCodeStreamWrongLastSequence) {
continue
}
return false, err
}
return true, nil
}
return false, fmt.Errorf("the ask %s changed under every one of %d tries", id, askChangeTries)
}
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
}
// routerHereIn says whether a module declaring the asks seat, with its ask named by its caller, is assigned:
// without it nothing takes an ask, and asking would only fill a queue nobody reads.
func routerHereIn(inv *inventory.Inventory) func(ctx context.Context) (bool, error) {
return func(ctx context.Context) (bool, error) {
entries, err := inv.Catalogued(ctx)
if err != nil {
return false, err
}
for _, e := range entries {
for _, s := range e.Manifest.DefinesSeats {
if s.Name == broker.AsksSeat && s.NamedByCaller("ask") && len(e.On) > 0 {
return true, nil
}
}
}
return false, nil
}
}
// grantHeldIn says whether the bus holds the controller's grant to ask (novox/hq issue 353): the user list the
// machine holding the bus was last sent is the one the mesh composes now (brokerBehind, the same judgement a
// push makes to send that machine first). While it is behind, the controller's ask is refused by the bus,
// whatever the record says of the router, so nothing is asked and the operator is told to push that machine.
// Judged at most every grantLookEvery: composing the list resolves the bus's machine whole.
func grantHeldIn(open *stores) func(ctx context.Context) (bool, string, error) {
var mu sync.Mutex
var at time.Time
var held bool
var why string
return func(ctx context.Context) (bool, string, error) {
mu.Lock()
defer mu.Unlock()
if !at.IsZero() && time.Since(at) < grantLookEvery {
return held, why, nil
}
machine, behind, err := brokerBehind(ctx, open, nil)
if err != nil {
return false, "", err
}
at, held, why = time.Now(), !behind, ""
if behind {
why = fmt.Sprintf("the bus's user list on %s is behind what the mesh composes, so the bus has not been "+
"given the controller's grant to ask; `push %s` carries it", machine, machine)
}
return held, why, nil
}
}
// grantLookEvery is how often the bus's user list is judged against the one its machine was last sent.
const grantLookEvery = 30 * time.Second
// channelsIn is what the channels are now, as a fingerprint: each module claiming a kind of the channel
// bench, where, promising what, and whether of its own account. An ask the router refused is asked again
// once this changes.
func channelsIn(inv *inventory.Inventory) func(ctx context.Context) string {
return func(ctx context.Context) string {
entries, err := inv.Catalogued(ctx)
if err != nil {
return ""
}
var parts []string
for _, e := range entries {
for _, c := range e.Manifest.Claims {
if c.Kind == "" || !catalogue.KindedBenches[c.Name] {
continue
}
on := append([]string(nil), e.On...)
sort.Strings(on)
caps := append([]string(nil), c.Capabilities...)
sort.Strings(caps)
parts = append(parts, fmt.Sprintf("%s/%s=%s@%s[%s]own:%t", c.Name, c.Kind, e.Manifest.Module,
strings.Join(on, ","), strings.Join(caps, ","), e.Manifest.RunsAs != ""))
}
}
sort.Strings(parts)
return strings.Join(parts, ";")
}
}
// 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),
routerHere: routerHereIn(open.inventory),
grantHeld: grantHeldIn(open),
channels: channelsIn(open.inventory),
raise: func(ctx context.Context, obs []conditions.Observation) error {
return keeper.Reconcile(ctx, sourceAsker, obs)
},
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)
}