diff --git a/cmd/mesh-controller/asker.go b/cmd/mesh-controller/asker.go index 5202fc19..3b28ba5e 100644 --- a/cmd/mesh-controller/asker.go +++ b/cmd/mesh-controller/asker.go @@ -84,21 +84,74 @@ type asked struct { // Acted is what the controller did on the warrant: empty before it did anything, "acting" while it acts, // then "done", "failed: …" or "nothing: …". Anything but empty is never acted on again. Acted string `json:"acted,omitempty"` - // Drill is an ask started at the controller's terminal (drill.go): about no condition, its answers + // Part is which ask of its condition this is (askPart): empty for the one that carries the condition's + // answers, or the authorising ones where it has both; "acknowledge" for its acknowledging answers asked + // apart (the review of 2026-10-09, M1). + Part string `json:"part,omitempty"` + // Rehearsal is an ask started at the controller's terminal (rehearse.go): about no condition, its answers // perform nothing, and the reconciling of conditions leaves it alone. - Drill bool `json:"drill,omitempty"` + Rehearsal bool `json:"rehearsal,omitempty"` } -// askedStore keeps the asks (broker.AskedBucket). +// partKey is an ask's place among what is asked: its condition and its part. +func partKey(condition, part string) string { return condition + "#" + part } + +// partAcknowledge is the part of a condition asked apart for its acknowledging answers. +const partAcknowledge = "acknowledge" + +// askPart is one ask a condition is asked with: its part, what it is about, and its answers. +type askPart struct { + name string + about string + actions []conditions.Action +} + +// levelOf is an action's level as an option offers it: one that says none is never taken for less than +// approve. +func levelOf(act conditions.Action) asks.Level { + if act.Level == "" { + return asks.Approve + } + return asks.Level(act.Level) +} + +// partsOf is the asks a condition is asked with (the review of 2026-10-09, M1): one, when its answers are all +// of one kind; else its authorising answers (Release, Stop, Restart) in one ask, about the condition, and its +// acknowledging ones (Silence) in another. **An acknowledgement never shares an ask with an approval**: a +// channel that only acknowledges would otherwise answer the ask, and end the approval with it. +func partsOf(c conditions.Condition) []askPart { + var ack, auth []conditions.Action + for _, act := range c.Actions { + if levelOf(act) == asks.Acknowledge { + ack = append(ack, act) + } else { + auth = append(auth, act) + } + } + if len(ack) == 0 || len(auth) == 0 { + return []askPart{{about: c.Key, actions: c.Actions}} + } + return []askPart{{about: c.Key, actions: auth}, + {name: partAcknowledge, about: c.Key + "." + partAcknowledge, actions: ack}} +} + +// askedStore keeps the asks (broker.AskedBucket). **Every write after the first is a compare-and-set** (the +// review of 2026-10-09, L2): an ask is created once, and changed only over the revision it was read at, the +// change decided again on what is read — so two controllers, or two deliveries of one warrant, never write +// over each other, and of two that would act only the one whose write stands does. type askedStore interface { Get(ctx context.Context, id string) (*asked, error) - Put(ctx context.Context, a asked) error + // Create keeps a new ask, and refuses one already kept under its id. + Create(ctx context.Context, a asked) error + // Change applies change to the ask kept under id, by compare-and-set, and says whether its write stood. + // change says whether to write at all; on a write that came between, it is asked again on what is read. + Change(ctx context.Context, id string, change func(*asked) bool) (bool, error) All(ctx context.Context) ([]asked, error) - // Claim marks an open ask acting, by compare-and-set, and says whether this write stood: of two - // deliveries of one warrant, or two controllers, only the one whose write stands acts. - Claim(ctx context.Context, id string, w asks.Warrant) (bool, error) } +// askChangeTries is how often a change is read and tried again when another write came between. +const askChangeTries = 5 + // asker is the controller asking the operator and acting on the answer. type asker struct { open func(ctx context.Context) ([]conditions.Condition, error) @@ -217,11 +270,12 @@ func (a *asker) reconcile(ctx context.Context) error { if err != nil { return err } - byCondition := map[string]asked{} + byCondition := map[string]asked{} // by partKey for _, r := range all { - if r.State == askOpen && !r.Drill { - if prior, held := byCondition[r.Condition]; !held || r.Opened.After(prior.Opened) { - byCondition[r.Condition] = r + if r.State == askOpen && !r.Rehearsal { + k := partKey(r.Condition, r.Part) + if prior, held := byCondition[k]; !held || r.Opened.After(prior.Opened) { + byCondition[k] = r } } } @@ -243,8 +297,8 @@ func (a *asker) reconcile(ctx context.Context) error { } byCondition = map[string]asked{} for _, r := range all { - if r.State == askOpen && !r.Drill { - byCondition[r.Condition] = r + if r.State == askOpen && !r.Rehearsal { + byCondition[partKey(r.Condition, r.Part)] = r } } } @@ -252,12 +306,13 @@ func (a *asker) reconcile(ctx context.Context) error { // newest first: not asked again until the answers or the channels change. answered, refused := map[string]asked{}, map[string]asked{} for _, r := range all { + k := partKey(r.Condition, r.Part) if r.State == string(asks.OutcomeChosen) && now.Sub(r.Ended) < askAgainAfterAnswer { - answered[r.Condition] = r + answered[k] = r } if r.State == string(asks.OutcomeRefused) { - if prior, has := refused[r.Condition]; !has || r.Opened.After(prior.Opened) { - refused[r.Condition] = r + if prior, has := refused[k]; !has || r.Opened.After(prior.Opened) { + refused[k] = r } } } @@ -277,54 +332,70 @@ func (a *asker) reconcile(ctx context.Context) error { }) openNow := 0 for _, c := range open { - if r, held := byCondition[c.Key]; held && wants(c, now) && sameAsked(r.Actions, c.Actions) && now.Before(r.Ask.Expires) { - openNow++ + if !wants(c, now) { + continue + } + for _, p := range partsOf(c) { + if r, held := byCondition[partKey(c.Key, p.name)]; held && sameAsked(r.Actions, p.actions) && now.Before(r.Ask.Expires) { + openNow++ + } } } for _, c := range open { if !wants(c, now) { continue } - wanted[c.Key] = true - if r, was := refused[c.Key]; was && sameAsked(r.Actions, c.Actions) && r.Channels == channels { - if _, held := byCondition[c.Key]; !held { - unasked = append(unasked, c) - if r.Warrant != nil && r.Warrant.Words != "" { - refusedWords = append(refusedWords, r.Warrant.Words) + saidUnasked := false + for _, p := range partsOf(c) { + key := partKey(c.Key, p.name) + wanted[key] = true + if r, was := refused[key]; was && sameAsked(r.Actions, p.actions) && r.Channels == channels { + if _, held := byCondition[key]; !held { + if !saidUnasked { + unasked, saidUnasked = append(unasked, c), true + } + if r.Warrant != nil && r.Warrant.Words != "" { + refusedWords = append(refusedWords, r.Warrant.Words) + } + continue // refused, and nothing it was refused for has changed } - continue // refused, and nothing it was refused for has changed } - } - if r, done := answered[c.Key]; done && sameAsked(r.Actions, c.Actions) { - if _, held := byCondition[c.Key]; !held { + if r, done := answered[key]; done && sameAsked(r.Actions, p.actions) { + if _, held := byCondition[key]; !held { + continue + } + } + if r, held := byCondition[key]; held { + switch { + case !sameAsked(r.Actions, p.actions): + if err := a.cancel(ctx, r, "its answers changed"); err != nil { + return err + } + case !now.Before(r.Ask.Expires): + // Expired unanswered: the router says so too; asked again below while it lasts. + if _, err := a.store.Change(ctx, r.ID, func(x *asked) bool { + if x.State != askOpen { + return false + } + x.State, x.Ended = string(asks.OutcomeExpired), now + return true + }); err != nil { + return err + } + openNow-- + default: + continue + } + } + if openNow >= askMostOpen { + continue // asked when one of the open ones ends, most urgent first + } + if err := a.ask(ctx, c, p, channels); err != nil { + a.logf("the operator could not be asked about %s: %v", c.Key, err) continue } + openNow++ } - if r, held := byCondition[c.Key]; held { - switch { - case !sameAsked(r.Actions, c.Actions): - if err := a.cancel(ctx, r, "its answers changed"); err != nil { - return err - } - case !now.Before(r.Ask.Expires): - // Expired unanswered: the router says so too; asked again below while it lasts. - r.State, r.Ended = string(asks.OutcomeExpired), now - if err := a.store.Put(ctx, r); err != nil { - return err - } - openNow-- - default: - continue - } - } - if openNow >= askMostOpen { - continue // asked when one of the open ones ends, most urgent first - } - if err := a.ask(ctx, c, channels); err != nil { - a.logf("the operator could not be asked about %s: %v", c.Key, err) - continue - } - openNow++ } for key, r := range byCondition { if !wanted[key] { @@ -420,18 +491,15 @@ func askText(s string) string { return strings.ReplaceAll(s, "..", ".") } -// askOf is the ask a condition is asked with. -func askOf(id string, c conditions.Condition, now time.Time) (asks.Ask, map[string]int) { +// askOf is the ask one part of a condition is asked with. +func askOf(id string, c conditions.Condition, p askPart, now time.Time) (asks.Ask, map[string]int) { q := asks.Ask{ID: id, Headline: c.Headline, Explanation: askText(c.Explanation), Who: asks.Operator, - OnExpiry: "nothing is done, and you are asked again while it lasts", About: c.Key, + OnExpiry: "nothing is done, and you are asked again while it lasts", About: p.about, Urgent: c.Severity == conditions.Urgent} options := map[string]int{} approves := false - for i, act := range c.Actions { - level := asks.Level(act.Level) - if level == "" { - level = asks.Approve // an action that says nothing of its level is never taken for less - } + for i, act := range p.actions { + level := levelOf(act) // an action that says nothing of its level is never taken for less than approve approves = approves || level != asks.Acknowledge oid := optionID(act.Label) options[oid] = i @@ -448,10 +516,14 @@ func askOf(id string, c conditions.Condition, now time.Time) (asks.Ask, map[stri return q, options } -// boundAct is what an option's Binds digests: the act exactly as the controller will perform it, with its -// level, and never its label or words. -func boundAct(act conditions.Action) map[string]any { - return map[string]any{"verb": act.Verb, "machine": act.Machine, "arguments": act.Arguments, "level": act.Level} +// boundAct is what an option's Binds digests: the act exactly as the controller will perform it — its verb, +// machine, level, and each argument as "arg." — and never its label or words. +func boundAct(act conditions.Action) asks.Act { + out := asks.Act{"verb": act.Verb, "machine": act.Machine, "level": act.Level} + for k, v := range act.Arguments { + out["arg."+k] = v + } + return out } func newAskID() string { @@ -460,11 +532,15 @@ func newAskID() string { return "c" + hex.EncodeToString(b[:]) } -// ask publishes one ask about a condition, and keeps it. -func (a *asker) ask(ctx context.Context, c conditions.Condition, channels string) error { +// askUnsent is an ask kept and never published: asked again at the next look. +const askUnsent = "unsent" + +// ask publishes one ask about a part of a condition, kept before it is published (the review of 2026-10-09, +// L3): a warrant for it then always finds it, and one whose publishing failed is marked so and asked again. +func (a *asker) ask(ctx context.Context, c conditions.Condition, p askPart, channels string) error { now := a.now() id := newAskID() - q, options := askOf(id, c, now) + q, options := askOf(id, c, p, now) if err := q.Check(now); err != nil { return err } @@ -472,24 +548,47 @@ func (a *asker) ask(ctx context.Context, c conditions.Condition, channels string if err != nil { return err } + if err := a.store.Create(ctx, asked{ID: id, Condition: c.Key, Part: p.name, Ask: q, Actions: p.actions, + Options: options, State: askOpen, Opened: now, Channels: channels}); err != nil { + return fmt.Errorf("the ask could not be kept, so it was not asked: %w", err) + } if err := a.publish(ctx, asks.AskSubject(askerName), body, "ask."+id); err != nil { + if _, cerr := a.store.Change(ctx, id, func(x *asked) bool { + if x.State != askOpen || x.Acted != "" { + return false + } + x.State, x.Ended, x.Acted = askUnsent, a.now(), "nothing: it could not be published: "+err.Error() + return true + }); cerr != nil { + a.logf("the ask %s could not be published, and could not be marked so: %v", id, cerr) + } return err } a.logf("asked the operator about %s (%s): %d answer(s)", c.Key, id, len(q.Options)) - return a.store.Put(ctx, asked{ID: id, Condition: c.Key, Ask: q, Actions: c.Actions, Options: options, - State: askOpen, Opened: now, Channels: channels}) + return nil } -// cancel takes an ask back. +// cancel takes an ask back: kept cancelled first, so a warrant that comes after is refused, then said to the +// router; a cancel the router did not hear leaves the ask to expire there, and nothing is done on it here. func (a *asker) cancel(ctx context.Context, r asked, why string) error { + stood, err := a.store.Change(ctx, r.ID, func(x *asked) bool { + if x.State != askOpen { + return false + } + x.State, x.Ended = askCancelled, a.now() + return true + }) + if err != nil || !stood { + return err + } body, _ := json.Marshal(map[string]string{"id": r.ID}) if err := a.publish(ctx, asks.CancelSubject(askerName), body, "cancel."+r.ID); err != nil { - a.logf("the ask %s about %s could not be taken back (%v); taken back at the next look", r.ID, r.Condition, err) + a.logf("the ask %s about %s is taken back here, and the router could not be told (%v): it expires there, "+ + "and no answer to it is acted on", r.ID, r.Condition, err) return nil } a.logf("took back the ask %s about %s: %s", r.ID, r.Condition, why) - r.State, r.Ended = askCancelled, a.now() - return a.store.Put(ctx, r) + return nil } // Decided takes the router's word on one of the controller's asks (link.Decider). An error is returned only @@ -517,13 +616,21 @@ func (a *asker) Decided(ctx context.Context, body []byte) error { } now := a.now() if w.Outcome != asks.OutcomeChosen { - r.State, r.Ended, r.Warrant = string(w.Outcome), now, &w - r.Acted = "nothing: the ask " + string(w.Outcome) + acted := "nothing: the ask " + string(w.Outcome) if w.Words != "" { - r.Acted += ": " + w.Words + acted += ": " + w.Words + } + if _, err := a.store.Change(ctx, r.ID, func(x *asked) bool { + if x.Acted != "" { + return false + } + x.State, x.Ended, x.Warrant, x.Acted = string(w.Outcome), now, &w, acted + return true + }); err != nil { + return err } a.logf("the ask %s about %s ended %s; nothing is done", r.ID, r.Condition, w.Outcome) - return a.store.Put(ctx, *r) + return nil } if r.State != askOpen { // Cancelled, replaced or expired in the controller's own record: no answer to it is acted on. @@ -547,24 +654,39 @@ func (a *asker) Decided(ctx context.Context, body []byte) error { a.logf("REFUSED a warrant for the ask %s: %v", r.ID, err) return nil } - r.State, r.Warrant = string(asks.OutcomeChosen), &w open, err := a.open(ctx) if err != nil { return err } - stillOpen := r.Drill // a drill is about no condition + stillOpen := r.Rehearsal // a rehearsal is about no condition for _, c := range open { stillOpen = stillOpen || c.Key == r.Condition } if !stillOpen { // The asker checks the state is still what it asked about before it acts (to-be 46 §10, step 7). - r.Ended, r.Acted = now, "nothing: the condition ended before the answer" + if _, err := a.store.Change(ctx, r.ID, func(x *asked) bool { + if x.State != askOpen || x.Acted != "" { + return false + } + x.State, x.Warrant, x.Ended, x.Acted = string(asks.OutcomeChosen), &w, now, + "nothing: the condition ended before the answer" + return true + }); err != nil { + return err + } a.logf("%s, for %s, which ended meanwhile: nothing is done", w.Says(), r.Condition) - return a.store.Put(ctx, *r) + return nil } // Claimed before acting, by compare-and-set: only the delivery whose write stands acts (security review - // of 2026-10-08, finding 9). - claimed, err := a.store.Claim(ctx, r.ID, w) + // of 2026-10-08, finding 9). Not by the warrant's message id, which another publisher could take first: + // the controller's own record decides. + claimed, err := a.store.Change(ctx, r.ID, func(x *asked) bool { + if x.State != askOpen || x.Acted != "" { + return false + } + x.State, x.Warrant, x.Acted = string(asks.OutcomeChosen), &w, "acting" + return true + }) if err != nil { return err } @@ -584,20 +706,26 @@ func (a *asker) Decided(ctx context.Context, body []byte) error { } var acted error switch { - case r.Drill && act.Verb == drillVerb: - // A drill's answer performs nothing: it is recorded below as the operator's decision. + case r.Rehearsal && act.Verb == rehearsalVerb: + // A rehearsal's answer performs nothing: it is recorded below as the operator's decision. case act.Arguments["silence"] != "": acted = a.silence(ctx, act.Arguments["silence"], conditions.MaxSilence, byWords(w), why) default: acted = a.call(ctx, act, args) } - r.Ended = a.now() - r.Acted = "done" + ended, outcome := a.now(), "done" if acted != nil { - r.Acted = "failed: " + acted.Error() + outcome = "failed: " + acted.Error() } - if err := a.store.Put(ctx, *r); err != nil { - return err + r.Ended, r.Acted = ended, outcome + if _, err := a.store.Change(ctx, r.ID, func(x *asked) bool { + if x.Acted != "acting" { + return false + } + x.Ended, x.Acted = ended, outcome + return true + }); err != nil { + a.logf("%s was acted on (%s), and how it ended could NOT be kept: %v", r.ID, outcome, err) } verbArgs := []string{act.Verb} if act.Machine != "" { diff --git a/cmd/mesh-controller/asker_bus_test.go b/cmd/mesh-controller/asker_bus_test.go new file mode 100644 index 00000000..0e22c949 --- /dev/null +++ b/cmd/mesh-controller/asker_bus_test.go @@ -0,0 +1,151 @@ +package main + +import ( + "context" + "encoding/json" + "sync" + "testing" + "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/conditions" + "github.com/novox/mesh-controller/internal/link" + "github.com/novox/mesh-controller/internal/testbus" +) + +// busAsker is an asker on a real bus's `asked` bucket, counting what it performs: two of them are two +// controllers sharing one record. +type busAskerRig struct { + mu sync.Mutex + called int + acts int + open []conditions.Condition + sent [][]byte +} + +func (rig *busAskerRig) asker(t *testing.T, conn *nats.Conn, now time.Time) *asker { + return &asker{ + open: func(context.Context) ([]conditions.Condition, error) { + rig.mu.Lock() + defer rig.mu.Unlock() + return rig.open, nil + }, + silence: func(context.Context, string, time.Duration, string, string) error { return nil }, + store: busAsked{conn: conn}, + publish: func(_ context.Context, subject string, body []byte, _ string) error { + rig.mu.Lock() + defer rig.mu.Unlock() + if subject == asks.AskSubject(askerName) { + rig.sent = append(rig.sent, body) + } + return nil + }, + call: func(context.Context, conditions.Action, map[string]string) error { + time.Sleep(20 * time.Millisecond) // long enough for the other delivery to arrive meanwhile + rig.mu.Lock() + defer rig.mu.Unlock() + rig.called++ + return nil + }, + record: func(context.Context, link.HandAct) error { + rig.mu.Lock() + defer rig.mu.Unlock() + rig.acts++ + return nil + }, + now: func() time.Time { return now }, + logf: t.Logf, + } +} + +func askedBus(t *testing.T) *nats.Conn { + t.Helper() + conn, err := nats.Connect(testbus.URL(t)) + if err != nil { + t.Fatal(err) + } + t.Cleanup(conn.Close) + js, err := jetstream.New(conn) + if err != nil { + t.Fatal(err) + } + if _, err := js.CreateKeyValue(context.Background(), jetstream.KeyValueConfig{Bucket: broker.AskedBucket}); err != nil { + t.Fatal(err) + } + return conn +} + +// The review of 2026-10-09 (L7): two deliveries of one warrant, to two controllers at once, perform its act +// exactly once and record it once — the record's compare-and-set decides, never the warrant's message id. +func TestTwoAnswersAtOnceActOnce(t *testing.T) { + conn := askedBus(t) + now := time.Date(2026, 10, 9, 14, 0, 0, 0, time.UTC) + rig := &busAskerRig{open: []conditions.Condition{heldCondition()}} + first, second := rig.asker(t, conn, now), rig.asker(t, conn, now) + if err := first.reconcile(context.Background()); err != nil { + t.Fatal(err) + } + if len(rig.sent) != 1 { + t.Fatalf("asked %d times", len(rig.sent)) + } + var q asks.Ask + _ = json.Unmarshal(rig.sent[0], &q) + release, _ := q.Option("release") + w := asks.Warrant{Ask: q.ID, Asker: askerName, About: q.About, Outcome: asks.OutcomeChosen, Option: release.ID, + Label: release.Label, Level: release.Level, Channel: "telegram", Proofs: []string{"P1"}, At: now, + AskDigest: q.Digest(), By: &asks.Person{Who: asks.Operator, Kind: "telegram", Identity: "42", Verified: "user id verified"}} + body, _ := json.Marshal(w) + var wg sync.WaitGroup + for _, a := range []*asker{first, second, first, second} { + wg.Add(1) + go func(a *asker) { + defer wg.Done() + if err := a.Decided(context.Background(), body); err != nil { + t.Error(err) + } + }(a) + } + wg.Wait() + if rig.called != 1 || rig.acts != 1 { + t.Fatalf("performed %d time(s), recorded %d time(s)", rig.called, rig.acts) + } + got, err := busAsked{conn: conn}.Get(context.Background(), q.ID) + if err != nil || got == nil || got.Acted != "done" { + t.Fatalf("kept as %+v (%v)", got, err) + } +} + +// The review of 2026-10-09 (L2): a write decided on a record read earlier never lands over one made since. A +// cancel read before the answer was acted on leaves the act's record as it is. +func TestAStaleCancelDoesNotWriteOverAnAct(t *testing.T) { + conn := askedBus(t) + now := time.Date(2026, 10, 9, 14, 0, 0, 0, time.UTC) + rig := &busAskerRig{open: []conditions.Condition{heldCondition()}} + a := rig.asker(t, conn, now) + if err := a.reconcile(context.Background()); err != nil { + t.Fatal(err) + } + var q asks.Ask + _ = json.Unmarshal(rig.sent[0], &q) + stale, _ := busAsked{conn: conn}.Get(context.Background(), q.ID) + release, _ := q.Option("release") + w := asks.Warrant{Ask: q.ID, Asker: askerName, About: q.About, Outcome: asks.OutcomeChosen, Option: release.ID, + Label: release.Label, Level: release.Level, Channel: "telegram", At: now, AskDigest: q.Digest(), + By: &asks.Person{Who: asks.Operator, Kind: "telegram", Identity: "42", Verified: "user id verified"}} + body, _ := json.Marshal(w) + if err := a.Decided(context.Background(), body); err != nil { + t.Fatal(err) + } + if err := a.cancel(context.Background(), *stale, "the condition ended"); err != nil { + t.Fatal(err) + } + got, _ := busAsked{conn: conn}.Get(context.Background(), q.ID) + if got.State != string(asks.OutcomeChosen) || got.Acted != "done" { + t.Errorf("a stale cancel wrote over the act: %+v", got) + } +} diff --git a/cmd/mesh-controller/asker_test.go b/cmd/mesh-controller/asker_test.go index 5e40c9c4..024fd1fd 100644 --- a/cmd/mesh-controller/asker_test.go +++ b/cmd/mesh-controller/asker_test.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "strings" + "sync" "testing" "time" @@ -19,24 +20,40 @@ import ( type memAskedStore map[string]asked +// memAskedMu guards every memAskedStore: Change is a compare-and-set as the bus's is. +var memAskedMu sync.Mutex + func (m memAskedStore) Get(_ context.Context, id string) (*asked, error) { + memAskedMu.Lock() + defer memAskedMu.Unlock() r, ok := m[id] if !ok { return nil, nil } return &r, nil } -func (m memAskedStore) Put(_ context.Context, r asked) error { m[r.ID] = r; return nil } -func (m memAskedStore) Claim(_ context.Context, id string, w asks.Warrant) (bool, error) { +func (m memAskedStore) Create(_ context.Context, r asked) error { + memAskedMu.Lock() + defer memAskedMu.Unlock() + if _, kept := m[r.ID]; kept { + return errors.New("an ask is kept under that id") + } + m[r.ID] = r + return nil +} +func (m memAskedStore) Change(_ context.Context, id string, change func(*asked) bool) (bool, error) { + memAskedMu.Lock() + defer memAskedMu.Unlock() r, ok := m[id] - if !ok || r.State != askOpen || r.Acted != "" { + if !ok || !change(&r) { return false, nil } - r.State, r.Warrant, r.Acted = string(asks.OutcomeChosen), &w, "acting" m[id] = r return true, nil } func (m memAskedStore) All(context.Context) ([]asked, error) { + memAskedMu.Lock() + defer memAskedMu.Unlock() var out []asked for _, r := range m { out = append(out, r) @@ -437,7 +454,7 @@ func TestTheAskDropsWhereItIsAnsweredAndTheConditionKeepsIt(t *testing.T) { if !strings.Contains(c.Explanation, FromMeshMCPServer) { t.Fatalf("the condition lost where it is answered: %q", c.Explanation) } - q, _ := askOf("x", c, time.Now()) + q, _ := askOf("x", c, partsOf(c)[0], time.Now()) if strings.Contains(q.Explanation, "mesh MCP server") || !strings.HasPrefix(q.Explanation, "Needs you: release it, or stop it.") { t.Errorf("the ask says %q", q.Explanation) } @@ -507,7 +524,7 @@ func TestAWarrantPerformsOnlyTheActItsOptionBound(t *testing.T) { }) } // Every option of an ask binds its act. - q, _ := askOf("x", heldCondition(), time.Now()) + q, _ := askOf("x", heldCondition(), partsOf(heldCondition())[0], time.Now()) for _, o := range q.Options { if o.Binds == "" { t.Errorf("the option %s binds nothing", o.ID) @@ -563,3 +580,65 @@ func TestAnAskThatCannotBeDeliveredIsSaid(t *testing.T) { t.Errorf("nothing needs asking, and still said: %+v", got) } } + +// The review of 2026-10-09 (M1): an acknowledging answer never shares an ask with an authorising one. A +// condition offering Restart and Silence is asked twice — Restart alone, about the condition, and Silence +// alone, apart — so Silence chosen on a channel that only acknowledges leaves the Restart ask open. +func TestAnAcknowledgementNeverSharesAnAskWithAnApproval(t *testing.T) { + r := newAskerRig(t) + key := "module.shanks.plex.down" + c := conditions.Condition{Key: key, Kind: "module-down", Severity: conditions.Urgent, Headline: "Plex down on shanks", + Explanation: "Needs you: restart it, or silence this.", Needs: "restart it, or silence this.", + Actions: []conditions.Action{ + {Label: "Restart", Verb: "node-service-manager.restart", Machine: "shanks", Level: conditions.LevelApprove, + Arguments: map[string]string{"unit": "plex"}}, + conditions.SilenceAction(key)}} + r.open = []conditions.Condition{c} + if err := r.a.reconcile(context.Background()); err != nil { + t.Fatal(err) + } + sent := r.asksSent(t) + if len(sent) != 2 { + t.Fatalf("asked %d time(s): %+v", len(sent), sent) + } + for _, q := range sent { + if err := q.Check(r.now); err != nil { + t.Errorf("%s: %v", q.About, err) + } + levels := map[asks.Level]bool{} + for _, o := range q.Options { + levels[o.Level] = true + } + if len(levels) != 1 { + t.Errorf("the ask about %s mixes levels: %+v", q.About, q.Options) + } + } + byAbout := map[string]asks.Ask{} + for _, q := range sent { + byAbout[q.About] = q + } + if q := byAbout[key]; len(q.Options) != 1 || q.Options[0].Label != "Restart" { + t.Errorf("the condition's own ask: %+v", q) + } + if q := byAbout[key+".acknowledge"]; len(q.Options) != 1 || q.Options[0].Level != asks.Acknowledge { + t.Errorf("the acknowledging ask: %+v", q) + } + // Silence chosen: performed, and the Restart ask stays open, never asked twice. + answerWith(t, r, r.warrantFor(t, key, "Silence for a week")) + if len(r.silenced) != 1 || len(r.called) != 0 { + t.Fatalf("silenced %v called %v", r.silenced, r.called) + } + _ = r.a.reconcile(context.Background()) + open := 0 + for _, a := range r.store { + if a.State == askOpen && a.Condition == key { + open++ + if a.Part != "" || a.Ask.Options[0].Label != "Restart" { + t.Errorf("the open ask is %+v", a) + } + } + } + if open != 1 || len(r.asksSent(t)) != 2 { + t.Errorf("after the silence: %d open, %d asked", open, len(r.asksSent(t))) + } +} diff --git a/cmd/mesh-controller/asker_wire.go b/cmd/mesh-controller/asker_wire.go index dd7ba87a..454986c8 100644 --- a/cmd/mesh-controller/asker_wire.go +++ b/cmd/mesh-controller/asker_wire.go @@ -57,7 +57,8 @@ func (b busAsked) Get(ctx context.Context, id string) (*asked, error) { return &r, json.Unmarshal(e.Value(), &r) } -func (b busAsked) Put(ctx context.Context, r asked) error { +// 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 @@ -66,10 +67,49 @@ func (b busAsked) Put(ctx context.Context, r asked) error { if err != nil { return err } - _, err = kv.Put(ctx, r.ID, body) + _, 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 { @@ -94,41 +134,6 @@ func (b busAsked) All(ctx context.Context) ([]asked, error) { return out, nil } -// Claim marks an open ask acting, by compare-and-set on its key's revision: only the write that stands acts. -func (b busAsked) Claim(ctx context.Context, id string, w asks.Warrant) (bool, error) { - kv, err := b.kv(ctx) - if err != nil { - return false, err - } - 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 r.State != askOpen || r.Acted != "" { - return false, nil - } - r.State, r.Warrant, r.Acted = string(asks.OutcomeChosen), &w, "acting" - 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) { - return false, nil - } - return false, err - } - return true, 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 { diff --git a/cmd/mesh-controller/drill.go b/cmd/mesh-controller/drill.go deleted file mode 100644 index 1db79cf8..00000000 --- a/cmd/mesh-controller/drill.go +++ /dev/null @@ -1,110 +0,0 @@ -package main - -// The drill of the operator's answers (novox/hq ADR 0259, the live acceptance after rollout): an ask the -// operator starts at the controller's terminal, answered on the phone, whose approval changes nothing and is -// recorded as a person's decision like any other. -// -// mesh-controller drill [--for 15m] -// -// It asks with two answers — Approve (an approval, so only a channel that proves who answered carries it) and -// Decline (an acknowledgement) — each bound to the drill's own act. The serving controller acts on the warrant -// as on any other: it claims the ask once, checks the act is the one bound, performs nothing, and records the -// hand-act `warrant` with who answered, through which channel, and the proofs. `hand-acts` then shows it. -// -// **The terminal's alone**: a command a verb runs (MESH_VERB set) is refused, so no agent starts a drill — a -// drill is a question the operator expects, and one an agent could start would teach them to approve what they -// did not ask for. - -import ( - "context" - "encoding/json" - "errors" - "flag" - "fmt" - "os" - "time" - - "github.com/nats-io/nats.go" - - "git.novox.be/novox/mesh-sdk/go/asks" - - "github.com/novox/mesh-controller/internal/conditions" -) - -// drillVerb is the act a drill's answers bind: nothing is called. -const drillVerb = "drill" - -// drillActions are the drill's two answers. -func drillActions() []conditions.Action { - return []conditions.Action{ - {Label: "Approve", Verb: drillVerb, Level: conditions.LevelApprove, Arguments: map[string]string{"drill": "approve"}}, - {Label: "Decline", Verb: drillVerb, Level: conditions.LevelAcknowledge, Arguments: map[string]string{"drill": "decline"}}, - } -} - -// drillAsk is the drill's ask, as the router is sent it. -func drillAsk(id string, now time.Time, lasts time.Duration) (asks.Ask, map[string]int) { - q := asks.Ask{ID: id, Headline: "Drill: approve this test question?", Who: asks.Operator, - Explanation: "Needs you: approve or decline. You started this drill at the controller's terminal. Approving " + - "changes nothing on the mesh; it is recorded as your decision, so you can check the record.", - OnExpiry: "nothing is done", Expires: now.Add(lasts), About: "drill." + id} - options := map[string]int{} - for i, act := range drillActions() { - binds, _ := asks.ActDigest(boundAct(act)) - oid := optionID(act.Label) - options[oid] = i - q.Options = append(q.Options, asks.Option{ID: oid, Label: act.Label, Does: doesDrill(act), Level: asks.Level(act.Level), - Binds: binds}) - } - return q, options -} - -func doesDrill(act conditions.Action) string { - if act.Arguments["drill"] == "approve" { - return "nothing changes; your approval is recorded" - } - return "nothing changes; your answer is recorded" -} - -func drillCommand(ctx context.Context, args []string) error { - if os.Getenv(verbVar) != "" { - return errors.New("drill is the controller's terminal's alone: a verb may not start one, so no agent asks " + - "the operator a question they did not start (novox/hq ADR 0259)") - } - set := flag.NewFlagSet("drill", flag.ContinueOnError) - lasts := set.Duration("for", 15*time.Minute, "how long the question waits for an answer") - if err := set.Parse(args); err != nil { - return err - } - if *lasts < time.Minute || *lasts > askApproveFor { - return fmt.Errorf("a drill waits between a minute and %s", askApproveFor) - } - js, err := aBus() - if err != nil { - return err - } - defer js.Close() - now := time.Now() - id := newAskID() - q, options := drillAsk(id, now, *lasts) - if err := q.Check(now); err != nil { - return err - } - store := busAsked{conn: js.Conn()} - // Kept before it is published, as the asker keeps every ask, so a warrant always finds it. - if err := store.Put(ctx, asked{ID: id, Condition: q.About, Ask: q, Actions: drillActions(), Options: options, - State: askOpen, Opened: now, Drill: true}); err != nil { - return fmt.Errorf("the drill could not be kept in the controller's asks: %w", err) - } - body, err := json.Marshal(q) - if err != nil { - return err - } - if _, err := js.Context().Publish(asks.AskSubject(askerName), body, nats.MsgId("ask."+id), nats.Context(ctx)); err != nil { - return fmt.Errorf("the drill could not be asked: %w", err) - } - fmt.Printf("drill %s asked: answer it on your phone before %s. Then `mesh-controller hand-acts` shows the "+ - "answer as a warrant, with who answered, through which channel and the proofs; nothing else changes.\n", - id, q.Expires.Local().Format("15:04")) - return nil -} diff --git a/cmd/mesh-controller/drill_test.go b/cmd/mesh-controller/drill_test.go deleted file mode 100644 index 074abd4b..00000000 --- a/cmd/mesh-controller/drill_test.go +++ /dev/null @@ -1,60 +0,0 @@ -package main - -import ( - "context" - "strings" - "testing" - "time" - - "git.novox.be/novox/mesh-sdk/go/asks" -) - -// A drill (the live acceptance of novox/hq ADR 0259): its approval is a warrant like any other — claimed once, -// its act checked against what the option bound, recorded as the operator's decision with who, how and the -// proofs — and it performs nothing. The reconciling of conditions leaves it open. -func TestADrillsApprovalIsRecordedAndPerformsNothing(t *testing.T) { - r := newAskerRig(t) - q, options := drillAsk("cdrill", r.now, askerDrillFor) - if err := q.Check(r.now); err != nil { - t.Fatalf("the drill's ask is refused: %v", err) - } - r.store["cdrill"] = asked{ID: "cdrill", Condition: q.About, Ask: q, Actions: drillActions(), Options: options, - State: askOpen, Opened: r.now, Drill: true} - if err := r.a.reconcile(context.Background()); err != nil { - t.Fatal(err) - } - if got := r.store["cdrill"]; got.State != askOpen { - t.Fatalf("the reconciling of conditions ended the drill: %+v", got) - } - approve, _ := q.Option("approve") - w := asks.Warrant{Ask: "cdrill", Asker: "mesh-controller", Outcome: asks.OutcomeChosen, Option: approve.ID, - Label: approve.Label, Level: approve.Level, Channel: "telegram", Proofs: []string{"P1"}, At: r.now, - AskDigest: q.Digest(), By: &asks.Person{Who: asks.Operator, Kind: "telegram", Identity: "42", Verified: "user id verified"}} - answerWith(t, r, w) - answerWith(t, r, w) // heard again - if len(r.called)+len(r.silenced) != 0 { - t.Errorf("a drill performed something: %v %v", r.called, r.silenced) - } - if len(r.acts) != 1 { - t.Fatalf("hand-acts %+v", r.acts) - } - act := r.acts[0] - if act.Verb != handActWarrant || act.By != "the operator, as telegram identity 42" || act.Ask != "cdrill" || - strings.Join(act.Args, " ") != "drill drill=approve" || act.Outcome != "done" || strings.Join(act.Proofs, ",") != "P1" { - t.Errorf("the drill's record: %+v", act) - } - if !personsDecision(act) { - t.Error("a drill's answer counts as a repair") - } -} - -// Only the terminal starts a drill: a verb's process is refused before anything is asked. -func TestADrillIsTheTerminalsAlone(t *testing.T) { - t.Setenv(verbVar, "mesh-controller.command") - if err := drillCommand(context.Background(), nil); err == nil || !strings.Contains(err.Error(), "terminal") { - t.Fatalf("a verb started a drill: %v", err) - } -} - -// askerDrillFor is how long the test's drill waits. -const askerDrillFor = 15 * time.Minute diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index 6a085307..81e86c87 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -78,8 +78,8 @@ func run() error { return rotateCommand(ctx, args[1:]) case "ask": return askCommand(ctx, args[1:]) - case "drill": - return drillCommand(ctx, args[1:]) + case "rehearse": + return rehearseCommand(ctx, args[1:]) case "builds": return buildsCommand(ctx, args[1:]) // The build queue, controlled by hand (novox/hq ADR 0219). diff --git a/cmd/mesh-controller/rehearse.go b/cmd/mesh-controller/rehearse.go new file mode 100644 index 00000000..5ec46044 --- /dev/null +++ b/cmd/mesh-controller/rehearse.go @@ -0,0 +1,111 @@ +package main + +// The rehearsal of the operator's answers — not a drill, which in the glossary is something broken on purpose (novox/hq ADR 0259, the live acceptance after rollout): an ask the +// operator starts at the controller's terminal, answered on the phone, whose approval changes nothing and is +// recorded as a person's decision like any other. +// +// mesh-controller rehearse [--for 15m] +// +// It asks with two answers, Approve and Decline, each bound to the rehearsal's own act and **both at the level +// approve** (the review of 2026-10-09, M1: an acknowledgement never shares an ask with an approval), so only a +// channel that proves who answered carries either — the rehearsal is of exactly that. The serving controller acts on the warrant +// as on any other: it claims the ask once, checks the act is the one bound, performs nothing, and records the +// hand-act `warrant` with who answered, through which channel, and the proofs. `hand-acts` then shows it. +// +// **The terminal's alone**: a command a verb runs (MESH_VERB set) is refused, so no agent starts a rehearsal — a +// rehearsal is a question the operator expects, and one an agent could start would teach them to approve what they +// did not ask for. + +import ( + "context" + "encoding/json" + "errors" + "flag" + "fmt" + "os" + "time" + + "github.com/nats-io/nats.go" + + "git.novox.be/novox/mesh-sdk/go/asks" + + "github.com/novox/mesh-controller/internal/conditions" +) + +// rehearsalVerb is the act a rehearsal's answers bind: nothing is called. +const rehearsalVerb = "rehearsal" + +// rehearsalActions are the rehearsal's two answers. +func rehearsalActions() []conditions.Action { + return []conditions.Action{ + {Label: "Approve", Verb: rehearsalVerb, Level: conditions.LevelApprove, Arguments: map[string]string{"rehearsal": "approve"}}, + {Label: "Decline", Verb: rehearsalVerb, Level: conditions.LevelApprove, Arguments: map[string]string{"rehearsal": "decline"}}, + } +} + +// rehearsalAsk is the rehearsal's ask, as the router is sent it. +func rehearsalAsk(id string, now time.Time, lasts time.Duration) (asks.Ask, map[string]int) { + q := asks.Ask{ID: id, Headline: "Rehearsal: approve this test question?", Who: asks.Operator, + Explanation: "Needs you: approve or decline. You started this rehearsal at the controller's terminal. Approving " + + "changes nothing on the mesh; it is recorded as your decision, so you can check the record.", + OnExpiry: "nothing is done", Expires: now.Add(lasts), About: "rehearsal." + id} + options := map[string]int{} + for i, act := range rehearsalActions() { + binds, _ := asks.ActDigest(boundAct(act)) + oid := optionID(act.Label) + options[oid] = i + q.Options = append(q.Options, asks.Option{ID: oid, Label: act.Label, Does: doesRehearsal(act), Level: asks.Level(act.Level), + Binds: binds}) + } + return q, options +} + +func doesRehearsal(act conditions.Action) string { + if act.Arguments["rehearsal"] == "approve" { + return "nothing changes; your approval is recorded" + } + return "nothing changes; your answer is recorded" +} + +func rehearseCommand(ctx context.Context, args []string) error { + if os.Getenv(verbVar) != "" { + return errors.New("rehearse is the controller's terminal's alone: a verb may not start one, so no agent asks " + + "the operator a question they did not start (novox/hq ADR 0259)") + } + set := flag.NewFlagSet("rehearse", flag.ContinueOnError) + lasts := set.Duration("for", 15*time.Minute, "how long the question waits for an answer") + if err := set.Parse(args); err != nil { + return err + } + if *lasts < time.Minute || *lasts > askApproveFor { + return fmt.Errorf("a rehearsal waits between a minute and %s", askApproveFor) + } + js, err := aBus() + if err != nil { + return err + } + defer js.Close() + now := time.Now() + id := newAskID() + q, options := rehearsalAsk(id, now, *lasts) + if err := q.Check(now); err != nil { + return err + } + store := busAsked{conn: js.Conn()} + // Kept before it is published, as the asker keeps every ask, so a warrant always finds it. + if err := store.Create(ctx, asked{ID: id, Condition: q.About, Ask: q, Actions: rehearsalActions(), Options: options, + State: askOpen, Opened: now, Rehearsal: true}); err != nil { + return fmt.Errorf("the rehearsal could not be kept in the controller's asks: %w", err) + } + body, err := json.Marshal(q) + if err != nil { + return err + } + if _, err := js.Context().Publish(asks.AskSubject(askerName), body, nats.MsgId("ask."+id), nats.Context(ctx)); err != nil { + return fmt.Errorf("the rehearsal could not be asked: %w", err) + } + fmt.Printf("rehearsal %s asked: answer it on your phone before %s. Then `mesh-controller hand-acts` shows the "+ + "answer as a warrant, with who answered, through which channel and the proofs; nothing else changes.\n", + id, q.Expires.Local().Format("15:04")) + return nil +} diff --git a/cmd/mesh-controller/rehearse_test.go b/cmd/mesh-controller/rehearse_test.go new file mode 100644 index 00000000..4cfdeb20 --- /dev/null +++ b/cmd/mesh-controller/rehearse_test.go @@ -0,0 +1,60 @@ +package main + +import ( + "context" + "strings" + "testing" + "time" + + "git.novox.be/novox/mesh-sdk/go/asks" +) + +// A rehearsal (the live acceptance of novox/hq ADR 0259): its approval is a warrant like any other — claimed once, +// its act checked against what the option bound, recorded as the operator's decision with who, how and the +// proofs — and it performs nothing. The reconciling of conditions leaves it open. +func TestARehearsalsApprovalIsRecordedAndPerformsNothing(t *testing.T) { + r := newAskerRig(t) + q, options := rehearsalAsk("crehearsal", r.now, askerRehearsalFor) + if err := q.Check(r.now); err != nil { + t.Fatalf("the rehearsal's ask is refused: %v", err) + } + r.store["crehearsal"] = asked{ID: "crehearsal", Condition: q.About, Ask: q, Actions: rehearsalActions(), Options: options, + State: askOpen, Opened: r.now, Rehearsal: true} + if err := r.a.reconcile(context.Background()); err != nil { + t.Fatal(err) + } + if got := r.store["crehearsal"]; got.State != askOpen { + t.Fatalf("the reconciling of conditions ended the rehearsal: %+v", got) + } + approve, _ := q.Option("approve") + w := asks.Warrant{Ask: "crehearsal", Asker: "mesh-controller", Outcome: asks.OutcomeChosen, Option: approve.ID, + Label: approve.Label, Level: approve.Level, Channel: "telegram", Proofs: []string{"P1"}, At: r.now, + AskDigest: q.Digest(), By: &asks.Person{Who: asks.Operator, Kind: "telegram", Identity: "42", Verified: "user id verified"}} + answerWith(t, r, w) + answerWith(t, r, w) // heard again + if len(r.called)+len(r.silenced) != 0 { + t.Errorf("a rehearsal performed something: %v %v", r.called, r.silenced) + } + if len(r.acts) != 1 { + t.Fatalf("hand-acts %+v", r.acts) + } + act := r.acts[0] + if act.Verb != handActWarrant || act.By != "the operator, as telegram identity 42" || act.Ask != "crehearsal" || + strings.Join(act.Args, " ") != "rehearsal rehearsal=approve" || act.Outcome != "done" || strings.Join(act.Proofs, ",") != "P1" { + t.Errorf("the rehearsal's record: %+v", act) + } + if !personsDecision(act) { + t.Error("a rehearsal's answer counts as a repair") + } +} + +// Only the terminal starts a rehearsal: a verb's process is refused before anything is asked. +func TestARehearsalIsTheTerminalsAlone(t *testing.T) { + t.Setenv(verbVar, "mesh-controller.command") + if err := rehearseCommand(context.Background(), nil); err == nil || !strings.Contains(err.Error(), "terminal") { + t.Fatalf("a verb started a rehearsal: %v", err) + } +} + +// askerRehearsalFor is how long the test's rehearsal waits. +const askerRehearsalFor = 15 * time.Minute diff --git a/go.mod b/go.mod index 5da9ba9f..a6b3c201 100644 --- a/go.mod +++ b/go.mod @@ -3,7 +3,7 @@ module github.com/novox/mesh-controller go 1.26.0 require ( - git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261009081503-d4077b473ea8 + git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261009095928-76902998cd39 github.com/jackc/pgx/v5 v5.10.0 github.com/nats-io/nats-server/v2 v2.11.17 github.com/nats-io/nats.go v1.54.0 diff --git a/go.sum b/go.sum index eadf28cd..eb765ca2 100644 --- a/go.sum +++ b/go.sum @@ -1,13 +1,7 @@ -git.novox.be/novox/mesh-host v0.0.0-20261007162834-56e2ebec4bac h1:yLtFS0pDCCqIE9Zx8hgXEFG9fUWzf8L9WQoKV+Amk1E= -git.novox.be/novox/mesh-host v0.0.0-20261007162834-56e2ebec4bac/go.mod h1:VlilMCRZ5yyNXg7SNigNBLr0Gt32jrGw5KSNq5JAVYs= git.novox.be/novox/mesh-host v0.0.0-20261009081005-b28d7bbcbff4 h1:f4rBnKSemuN0Z9dTtRJMigIGfEs6ltFPOILJGHGab74= git.novox.be/novox/mesh-host v0.0.0-20261009081005-b28d7bbcbff4/go.mod h1:72ZATZjxMLaJfWdvlSDJrygIoBzCmKIjCDMhEXxVzTo= -git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261008145004-62367ce15ad6 h1:JT7xM1bnLNInW7/oImV2OlXTrcQ4/GSM0Y8tAb+AhmY= -git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261008145004-62367ce15ad6/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY= -git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261008162031-55090da7e08f h1:BNvyWq899GwP7F3sY4ACieB5a5fnFAq+sJ9lP6HQ5qI= -git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261008162031-55090da7e08f/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY= -git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261009081503-d4077b473ea8 h1:soqhLNpEXThdq6PdiPy6ExxjJ+yjhh1N1n9E3j1CtrM= -git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261009081503-d4077b473ea8/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY= +git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261009095928-76902998cd39 h1:WHW6CgbuTxP7M+qRBOgzsiG9vT49xdkZ/rarc9/vKMA= +git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261009095928-76902998cd39/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY= github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op h1:Z/MZK75wC/NSrkgqeNIa7jexam9uWzhLmFTSCPI/kn0= github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op/go.mod h1:FQyySiasQQM8735Ddel3MRojmy4dA1IqCeyJ5jmPMbI= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= diff --git a/internal/inventory/asks_lab_test.go b/internal/inventory/asks_lab_test.go index 62593d54..a96a9a73 100644 --- a/internal/inventory/asks_lab_test.go +++ b/internal/inventory/asks_lab_test.go @@ -98,7 +98,13 @@ func TestTheAsksLabBus(t *testing.T) { Emits: own.Emits, Serves: own.Serves} } records := broker.Records{Nodes: []string{labMachine}, Assigned: map[string][]broker.Declared{}, - People: map[string][]string{}, Interchangeable: map[string]bool{}} + People: map[string][]string{}, Interchangeable: map[string]bool{}, RootFree: map[string]bool{}} + // Whether the lab's machine is root-free is the lab's to say (MESH_LAB_ASKS_ROOT_FREE=true): it has no + // node-engine to judge it. Unsaid, it is not, and no kind is composed with verified-sender — as a push + // composes on a machine that is not (novox/hq ADR 0259 §8). + if os.Getenv("MESH_LAB_ASKS_ROOT_FREE") == "true" { + records.RootFree[labMachine] = true + } var buckets []broker.Bucket var trafficSeats []broker.Seat for _, m := range manifests { diff --git a/vendor/git.novox.be/novox/mesh-sdk/go/asks/asks.go b/vendor/git.novox.be/novox/mesh-sdk/go/asks/asks.go index ca2e8006..efe2a661 100644 --- a/vendor/git.novox.be/novox/mesh-sdk/go/asks/asks.go +++ b/vendor/git.novox.be/novox/mesh-sdk/go/asks/asks.go @@ -8,10 +8,11 @@ package asks import ( "crypto/sha256" "encoding/hex" - "encoding/json" "errors" "fmt" "regexp" + "sort" + "strconv" "strings" "time" ) @@ -107,25 +108,67 @@ type Option struct { Binds string `json:"binds,omitempty"` } -// ActDigest is the digest an asker puts in Option.Binds: SHA-256 over the act's JSON (Go's encoding sorts a -// map's keys, so the same act gives the same digest), written "sha256:". -func ActDigest(act any) (string, error) { - raw, err := json.Marshal(act) - if err != nil { - return "", fmt.Errorf("the act cannot be digested: %w", err) +// An Act is what an option binds: named fields, each a string — the verb, the machine, the level, and each +// argument under a name of its own ("arg.delivery"). Flat on purpose: its digest is over these names and +// values alone, never over how a language happens to encode a struct. +type Act map[string]string + +// ActDigest is the digest an asker puts in Option.Binds: SHA-256 over the act's canonical encoding (canonical), +// written "sha256:". The same names and values give the same digest in any language, whatever order +// they were set in; a field renamed, added or emptied gives another. +func ActDigest(act Act) (string, error) { + if len(act) == 0 { + return "", fmt.Errorf("the act cannot be digested: it names nothing") } - sum := sha256.Sum256(raw) + keys := make([]string, 0, len(act)) + for k := range act { + if k == "" { + return "", fmt.Errorf("the act cannot be digested: a field has no name") + } + keys = append(keys, k) + } + sort.Strings(keys) + var b strings.Builder + b.WriteString("novox.act.v1\n") + for _, k := range keys { + canonical(&b, k) + canonical(&b, act[k]) + } + sum := sha256.Sum256([]byte(b.String())) return "sha256:" + hex.EncodeToString(sum[:]), nil } +// canonical writes one value as its length in bytes, a colon, the bytes and a newline: no value can be read +// as another's end or start, so two different sequences of values never encode the same. +func canonical(b *strings.Builder, v string) { + b.WriteString(strconv.Itoa(len(v))) + b.WriteByte(':') + b.WriteString(v) + b.WriteByte('\n') +} + // Digest is the digest of the ask exactly as its asker published it: its id, words, options with what each // binds, who answers, and its expiry. The router puts it in the warrant (Warrant.AskDigest), and an asker // acts only on a warrant whose digest is that of the ask it keeps — so a warrant answers one ask, as the // person was shown it, and nothing published under the same id before or after. +// +// It is over the ask's named fields in a fixed order, each written canonically, and the expiry as UTC +// RFC 3339 to the nanosecond — never over a language's encoding of the struct, so a field added to Ask +// later changes no digest until it is added here, on purpose. func (a Ask) Digest() string { - a.Expires = a.Expires.UTC() - raw, _ := json.Marshal(a) - sum := sha256.Sum256(raw) + var b strings.Builder + b.WriteString("novox.ask.v1\n") + for _, v := range []string{a.ID, a.Headline, a.Explanation, a.Who, + a.Expires.UTC().Format(time.RFC3339Nano), a.OnExpiry, a.About, strconv.FormatBool(a.Urgent), + strconv.Itoa(len(a.Options))} { + canonical(&b, v) + } + for _, o := range a.Options { + for _, v := range []string{o.ID, o.Label, o.Does, string(o.Level), o.Binds} { + canonical(&b, v) + } + } + sum := sha256.Sum256([]byte(b.String())) return "sha256:" + hex.EncodeToString(sum[:]) } @@ -341,7 +384,7 @@ func (w Warrant) For(asker string, a Ask) (Option, error) { // Performs checks that the act an asker is about to perform is the one the chosen option bound when it // asked: the act's digest equals the option's Binds. An acknowledge option that bound nothing passes. -func (o Option) Performs(act any) error { +func (o Option) Performs(act Act) error { if o.Binds == "" && o.Level == Acknowledge { return nil } diff --git a/vendor/modules.txt b/vendor/modules.txt index f3c9f514..af68ed2c 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -1,4 +1,4 @@ -# git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261009081503-d4077b473ea8 +# git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261009095928-76902998cd39 ## explicit; go 1.22 git.novox.be/novox/mesh-sdk/go/asks # github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op