diff --git a/cmd/mesh-controller/asker.go b/cmd/mesh-controller/asker.go index 5896226e..b83e439e 100644 --- a/cmd/mesh-controller/asker.go +++ b/cmd/mesh-controller/asker.go @@ -43,7 +43,9 @@ const askerName = broker.ControllerSeat // How long an ask lasts: a day when an answer approves, a week when every answer only acknowledges. const ( - askApproveFor = 24 * time.Hour + // askApproveFor is a day less a margin, so an ask is never refused at the router for lasting a day and + // a moment (the SDK's bound is a day). + askApproveFor = 24*time.Hour - 10*time.Minute askAcknowledgeFor = 7 * 24 * time.Hour // askEvery is how often what is open is asked about again, beside every change. askEvery = time.Minute @@ -53,6 +55,9 @@ const ( // askAgainAfterAnswer is how long a condition the operator answered is not asked about again with the // same answers: what was chosen takes a while to clear it, and asking again at once would ask twice. askAgainAfterAnswer = time.Hour + // askMostOpen is how many asks the controller holds open at once (the router refuses a fourth): the + // most urgent conditions first, then the oldest. + askMostOpen = asks.MostOpen ) // What became of an ask, as the controller keeps it. @@ -63,10 +68,13 @@ const ( // asked is one ask the controller made, as it keeps it. type asked struct { - ID string `json:"id"` - Condition string `json:"condition"` - Ask asks.Ask `json:"ask"` - Actions []conditions.Action `json:"actions"` + ID string `json:"id"` + Condition string `json:"condition"` + // Channels is what the channels were when it was asked (asker.channels): an ask the router refused is not + // asked again until the condition's answers or the channels change. + Channels string `json:"channels,omitempty"` + Ask asks.Ask `json:"ask"` + Actions []conditions.Action `json:"actions"` // Options are the actions by option id. Options map[string]int `json:"options"` State string `json:"state"` @@ -83,6 +91,9 @@ type askedStore interface { Get(ctx context.Context, id string) (*asked, error) Put(ctx context.Context, a asked) 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) } // asker is the controller asking the operator and acting on the answer. @@ -98,8 +109,14 @@ type asker struct { record func(ctx context.Context, act link.HandAct) error // routerRecord reads the router's record of an ask for a warrant missed; nil reads nothing. routerRecord func(ctx context.Context, id string) (*asks.Warrant, error) - now func() time.Time - logf func(string, ...any) + // routerHere says whether a router holds the seat and takes asks under the asker's name; nil is yes. + routerHere func(ctx context.Context) (bool, error) + // channels is what the channels are now, as a fingerprint: who holds which kind, promising what. + channels func(ctx context.Context) string + now func() time.Time + logf func(string, ...any) + + saidNoRouter bool mu sync.Mutex nudged chan struct{} @@ -156,6 +173,25 @@ func sameAsked(a []conditions.Action, b []conditions.Action) bool { // reconcile brings what is asked in line with what is open. func (a *asker) reconcile(ctx context.Context) error { now := a.now() + if a.routerHere != nil { + here, err := a.routerHere(ctx) + if err != nil { + return err + } + if !here { + if !a.saidNoRouter { + a.logf("no router takes asks under the controller's name (a module declaring %s with its ask "+ + "named by its caller, assigned): the operator is asked nothing until one is", broker.AsksSeat) + a.saidNoRouter = true + } + return nil + } + a.saidNoRouter = false + } + channels := "" + if a.channels != nil { + channels = a.channels(ctx) + } open, err := a.open(ctx) if err != nil { return err @@ -195,20 +231,47 @@ func (a *asker) reconcile(ctx context.Context) error { } } } - // What the operator answered lately, by condition: not asked again at once. - answered := map[string]asked{} + // What the operator answered lately, by condition: not asked again at once; and what the router refused, + // newest first: not asked again until the answers or the channels change. + answered, refused := map[string]asked{}, map[string]asked{} for _, r := range all { if r.State == string(asks.OutcomeChosen) && now.Sub(r.Ended) < askAgainAfterAnswer { answered[r.Condition] = r } + if r.State == string(asks.OutcomeRefused) { + if prior, has := refused[r.Condition]; !has || r.Opened.After(prior.Opened) { + refused[r.Condition] = r + } + } } wanted := map[string]bool{} - sort.Slice(open, func(i, j int) bool { return open[i].Key < open[j].Key }) + // The most urgent first, then the oldest: those are asked when no more than askMostOpen may be. + sort.SliceStable(open, func(i, j int) bool { + ui, uj := open[i].Severity == conditions.Urgent, open[j].Severity == conditions.Urgent + if ui != uj { + return ui + } + if !open[i].Raised.Equal(open[j].Raised) { + return open[i].Raised.Before(open[j].Raised) + } + return open[i].Key < open[j].Key + }) + 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++ + } + } 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 { + 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 { continue @@ -226,13 +289,19 @@ func (a *asker) reconcile(ctx context.Context) error { if err := a.store.Put(ctx, r); err != nil { return err } + openNow-- default: continue } } - if err := a.ask(ctx, c); err != nil { - a.logf("the operator could not be asked about %s: %v", c.Key, err) + 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] { @@ -280,9 +349,18 @@ func doesWords(act conditions.Action) string { return strings.ToLower(act.Label) } +// askText is a condition's words as an ask says them: without where an answer is given when no channel can +// give it (FromMeshMCPServer), since the ask is answered on a channel and the router says where else. +func askText(s string) string { + for _, with := range []string{", " + FromMeshMCPServer, " " + FromMeshMCPServer} { + s = strings.ReplaceAll(s, with, ".") + } + 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) { - q := asks.Ask{ID: id, Headline: c.Headline, Explanation: c.Explanation, Who: asks.Operator, + 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, Urgent: c.Severity == conditions.Urgent} options := map[string]int{} @@ -311,7 +389,7 @@ func newAskID() string { } // ask publishes one ask about a condition, and keeps it. -func (a *asker) ask(ctx context.Context, c conditions.Condition) error { +func (a *asker) ask(ctx context.Context, c conditions.Condition, channels string) error { now := a.now() id := newAskID() q, options := askOf(id, c, now) @@ -327,7 +405,7 @@ func (a *asker) ask(ctx context.Context, c conditions.Condition) error { } 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}) + State: askOpen, Opened: now, Channels: channels}) } // cancel takes an ask back. @@ -375,6 +453,11 @@ func (a *asker) Decided(ctx context.Context, body []byte) error { a.logf("the ask %s about %s ended %s; nothing is done", r.ID, r.Condition, w.Outcome) return a.store.Put(ctx, *r) } + if r.State != askOpen { + // Cancelled, replaced or expired in the controller's own record: no answer to it is acted on. + a.logf("REFUSED a warrant for the ask %s, which is %s in the controller's own record", r.ID, r.State) + return nil + } option, err := w.For(askerName, r.Ask) if err != nil { a.logf("REFUSED a warrant for the ask %s: %v", r.ID, err) @@ -401,10 +484,18 @@ func (a *asker) Decided(ctx context.Context, body []byte) error { a.logf("%s, for %s, which ended meanwhile: nothing is done", w.Says(), r.Condition) return a.store.Put(ctx, *r) } - r.Acted = "acting" - if err := a.store.Put(ctx, *r); err != 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) + if err != nil { return err } + if !claimed { + a.logf("the warrant for the ask %s was already taken by another delivery; nothing more is done", r.ID) + return nil + } + r.Acted = "acting" + why := fmt.Sprintf("%s (ask %s)", w.Says(), r.ID) args := map[string]string{} for k, v := range act.Arguments { diff --git a/cmd/mesh-controller/asker_test.go b/cmd/mesh-controller/asker_test.go index e3100eb4..1bc43ae1 100644 --- a/cmd/mesh-controller/asker_test.go +++ b/cmd/mesh-controller/asker_test.go @@ -27,6 +27,15 @@ func (m memAskedStore) Get(_ context.Context, id string) (*asked, error) { 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) { + r, ok := m[id] + if !ok || r.State != askOpen || r.Acted != "" { + 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) { var out []asked for _, r := range m { @@ -340,3 +349,128 @@ func TestAWarrantMissedWhileAwayIsReadFromTheRoutersRecord(t *testing.T) { t.Errorf("asked again after the answer: %d", n) } } + +// After review (2026-10-08): a refused ask is not asked again until its answers or the channels change. +func TestAnAskTheRouterRefusedWaitsUntilSomethingChanges(t *testing.T) { + r := newAskerRig(t) + channels := "channel/telegram=telegram@anchor[choice]own:true" + r.a.channels = func(context.Context) string { return channels } + r.open = []conditions.Condition{heldCondition()} + _ = r.a.reconcile(context.Background()) + first := r.asksSent(t)[0] + refusal, _ := json.Marshal(asks.Warrant{Ask: first.ID, Asker: "mesh-controller", Outcome: asks.OutcomeRefused, + Words: "no channel can carry any of its answers now", At: r.now}) + if err := r.a.Decided(context.Background(), refusal); err != nil { + t.Fatal(err) + } + if got := r.store[first.ID]; got.State != string(asks.OutcomeRefused) || !strings.Contains(got.Acted, "nothing") { + t.Fatalf("the refusal was kept as %+v", got) + } + for i := 0; i < 3; i++ { + r.now = r.now.Add(askEvery) + _ = r.a.reconcile(context.Background()) + } + if n := len(r.asksSent(t)); n != 1 { + t.Fatalf("asked again %d time(s) though nothing changed", n-1) + } + channels = "channel/telegram=telegram@anchor[choice,verified-sender]own:true" + _ = r.a.reconcile(context.Background()) + if n := len(r.asksSent(t)); n != 2 { + t.Errorf("not asked again once the channels changed: %d", n) + } +} + +// After review: at most three asks open at once, the most urgent first, then the oldest. +func TestAtMostThreeAsksAreOpenTheMostUrgentFirst(t *testing.T) { + r := newAskerRig(t) + var open []conditions.Condition + for i := 0; i < 4; i++ { + c := unitsCondition() + c.Key = "machine.m" + string(rune('a'+i)) + ".units" + c.Actions = []conditions.Action{conditions.SilenceAction(c.Key)} + c.Raised = r.now.Add(-time.Duration(10-i) * time.Hour) + open = append(open, c) + } + urgent := heldCondition() + urgent.Severity, urgent.Raised = conditions.Urgent, r.now.Add(-time.Minute) + r.open = append(open, urgent) + _ = r.a.reconcile(context.Background()) + sent := r.asksSent(t) + if len(sent) != askMostOpen || sent[0].About != urgent.Key || sent[1].About != "machine.ma.units" || sent[2].About != "machine.mb.units" { + var about []string + for _, q := range sent { + about = append(about, q.About) + } + t.Fatalf("asked %v", about) + } +} + +// After review: nothing is asked while no router takes asks under the controller's name, and that is said once. +func TestNothingIsAskedWithoutARouter(t *testing.T) { + r := newAskerRig(t) + var said []string + r.a.logf = func(f string, a ...any) { said = append(said, f) } + r.a.routerHere = func(context.Context) (bool, error) { return false, nil } + r.open = []conditions.Condition{heldCondition()} + _ = r.a.reconcile(context.Background()) + _ = r.a.reconcile(context.Background()) + if len(r.asksSent(t)) != 0 { + t.Error("asked with no router") + } + n := 0 + for _, s := range said { + if strings.Contains(s, "no router takes asks") { + n++ + } + } + if n != 1 { + t.Errorf("said %d times", n) + } +} + +// After review: the condition's words keep where an answer is given without a channel; the ask's text does not. +func TestTheAskDropsWhereItIsAnsweredAndTheConditionKeepsIt(t *testing.T) { + c := heldCondition() + if !strings.Contains(c.Explanation, FromMeshMCPServer) { + t.Fatalf("the condition lost where it is answered: %q", c.Explanation) + } + q, _ := askOf("x", c, 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) + } + if askApproveFor >= 24*time.Hour { + t.Errorf("an approving ask lasts %s, which the SDK may refuse at its bound", askApproveFor) + } +} + +// After review (security finding 9): a warrant is acted on only for an ask open in the controller's own record, +// once the claim stands, and never when given after the ask expired. +func TestAWarrantIsActedOnlyForAnOpenAskItClaimsBeforeItExpired(t *testing.T) { + r := newAskerRig(t) + r.open = []conditions.Condition{heldCondition()} + _ = r.a.reconcile(context.Background()) + w := r.warrantFor(t, heldCondition().Key, "Release") + late := w + late.At = r.store[w.Ask].Ask.Expires.Add(time.Minute) + body, _ := json.Marshal(late) + _ = r.a.Decided(context.Background(), body) + if len(r.called) != 0 { + t.Fatalf("acted on a warrant given after the ask expired: %v", r.called) + } + // Claimed already by another delivery: nothing done here. + kept := r.store[w.Ask] + kept.Acted = "acting" + r.store[w.Ask] = kept + body, _ = json.Marshal(w) + _ = r.a.Decided(context.Background(), body) + if len(r.called) != 0 { + t.Fatalf("acted though the claim was another's: %v", r.called) + } + // Cancelled in its own record: refused. + kept.Acted, kept.State = "", askCancelled + r.store[w.Ask] = kept + _ = r.a.Decided(context.Background(), body) + if len(r.called) != 0 { + t.Errorf("acted on a cancelled ask: %v", r.called) + } +} diff --git a/cmd/mesh-controller/asker_wire.go b/cmd/mesh-controller/asker_wire.go index 283bbb44..9f0ae5f5 100644 --- a/cmd/mesh-controller/asker_wire.go +++ b/cmd/mesh-controller/asker_wire.go @@ -9,6 +9,7 @@ import ( "encoding/json" "errors" "fmt" + "sort" "strings" "time" @@ -93,6 +94,41 @@ 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 { @@ -174,6 +210,53 @@ func asksRecords(ctx context.Context, inv *inventory.Inventory) (string, error) 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 + } +} + +// 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) @@ -198,6 +281,8 @@ func startAsking(ctx context.Context, open *stores, server *link.Server, conn *n return err }, routerRecord: routerRecordOf(conn, open.inventory), + routerHere: routerHereIn(open.inventory), + channels: channelsIn(open.inventory), now: time.Now, logf: func(format string, args ...any) { fmt.Printf(format+"\n", args...) }, } diff --git a/cmd/mesh-controller/plain_words.go b/cmd/mesh-controller/plain_words.go index e5254f4b..efbd976a 100644 --- a/cmd/mesh-controller/plain_words.go +++ b/cmd/mesh-controller/plain_words.go @@ -680,11 +680,11 @@ func walkWaitingWords(w waitFacts, in time.Duration, severity conditions.Severit } // waitingNeeds is what the operator does about a walk waiting past its urgent bound: nothing before it. -// Start and Stop are asked of the operator (novox/hq ADR 0259), so the words do not say where: the router -// says where each can be answered. +// Start and Stop are also asked of the operator (novox/hq ADR 0259); the condition's own words keep saying +// where they are given without a channel, and the ask's text drops that (askText). func waitingNeeds(severity conditions.Severity) string { if severity == conditions.Urgent { - return "start it, or stop it." + return "start it, or stop it, " + FromMeshMCPServer } return "" } @@ -717,8 +717,8 @@ func moduleNeeds(node string, rs []inventory.ResourceHealth) string { } } if unit != "" { - // Asked of the operator (moduleActions): the router says where it can be answered. - return fmt.Sprintf("restart its service %s on %s.", unit, node) + // Also asked of the operator (moduleActions); the ask's text drops where (askText). + return fmt.Sprintf("restart its service %s on %s %s", unit, node, FromMeshMCPServer) } return "" } @@ -805,11 +805,11 @@ func stalledWords(l stalledLine, o conditions.Observation) (headline, explanatio Arguments: map[string]string{"id": l.ID, "why": ""}} switch held { case "held": - needs, actions = "release it, or stop it.", []conditions.Action{release, stop} + needs, actions = "release it, or stop it, "+FromMeshMCPServer, []conditions.Action{release, stop} case "ready", "checked": needs = "merge its pull request, or close it." default: - needs, actions = "stop it.", []conditions.Action{stop} + needs, actions = "stop it "+FromMeshMCPServer, []conditions.Action{stop} } } return fmt.Sprintf("Delivery of %s %s %s", name, held, long), diff --git a/cmd/mesh-controller/plain_words_test.go b/cmd/mesh-controller/plain_words_test.go index a1c80343..c82c00e4 100644 --- a/cmd/mesh-controller/plain_words_test.go +++ b/cmd/mesh-controller/plain_words_test.go @@ -70,7 +70,7 @@ func TestADeliveryWaitingNeedsNothingUntilItsBoundThenOffersStartAndStop(t *test f.waits[0].since = now.Add(-5 * time.Hour) got = watchWaits(f) plainExample(t, got[0], "openrazer delivery waiting to start", - "Needs you: start it, or stop it. The change to openrazer is merged and built, and mesh-delivery (the "+ + "Needs you: start it, or stop it, from the mesh MCP server; this notification cannot do it. The change to openrazer is merged and built, and mesh-delivery (the "+ "module that decides when a delivery goes out) has not let it start for 5 hours, so mesh-delivery may "+ "be stuck.", "Start", "Stop") for i, want := range []string{"go", "stop"} { @@ -97,7 +97,7 @@ func TestAModuleUnhealthyAsksForARestartInWords(t *testing.T) { Resource: "openrazer-daemon", Target: "openrazer-daemon.service", Account: "jochen", Reason: "failed in the account's own service manager (exit-code)", Since: time.Now()}}) plainExample(t, o, "openrazer not working on g14", - "Needs you: restart its service openrazer-daemon on g14. "+ + "Needs you: restart its service openrazer-daemon on g14 from the mesh MCP server; this notification cannot do it. "+ "openrazer on g14 is not healthy: its service openrazer-daemon stopped with an error. It clears as soon "+ "as it runs again.", "Restart") if a := o.Actions[0]; a.Verb != "node-service-manager.restart" || a.Machine != "g14" || a.Level != conditions.LevelApprove || @@ -165,7 +165,7 @@ func TestADeliveryHeldAsksForReleaseOrStopInWords(t *testing.T) { got := stalledObservations([]stalledLine{{ID: "novox/hq@055550802096", State: "held", For: "36h2m6s", Bound: "24h0m0s", H2: "none: the state is the operator's", Says: "it waits for the operator"}}) plainExample(t, got[0], "Delivery of hq held for 36 hours", - "Needs you: release it, or stop it. A delivery of hq has been held for 36 hours, past its limit.", + "Needs you: release it, or stop it, from the mesh MCP server; this notification cannot do it. A delivery of hq has been held for 36 hours, past its limit.", "Release", "Stop") for i, verb := range []string{"mesh-delivery.release", "mesh-delivery.stop"} { if a := got[0].Actions[i]; a.Verb != verb || a.Arguments["id"] != "novox/hq@055550802096" || a.Level != conditions.LevelApprove { diff --git a/go.mod b/go.mod index 21220b46..8ebef80d 100644 --- a/go.mod +++ b/go.mod @@ -3,7 +3,9 @@ module github.com/novox/mesh-controller go 1.26.0 require ( + git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261008162031-55090da7e08f 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 github.com/novox/mesh-host v0.0.0 golang.org/x/crypto v0.57.0 @@ -11,7 +13,6 @@ require ( ) require ( - git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261008145004-62367ce15ad6 // indirect github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op // indirect github.com/google/go-tpm v0.9.8 // indirect github.com/jackc/pgpassfile v1.0.0 // indirect @@ -20,10 +21,8 @@ require ( github.com/klauspost/compress v1.20.0 // indirect github.com/minio/highwayhash v1.0.4 // indirect github.com/nats-io/jwt/v2 v2.8.1 // indirect - github.com/nats-io/nats-server/v2 v2.11.17 // indirect github.com/nats-io/nkeys v0.4.16 // indirect github.com/nats-io/nuid v1.0.1 // indirect - go.uber.org/automaxprocs v1.6.0 // indirect golang.org/x/sync v0.23.0 // indirect golang.org/x/sys v0.48.0 // indirect golang.org/x/text v0.42.0 // indirect diff --git a/go.sum b/go.sum index 49aa5e9d..70849338 100644 --- a/go.sum +++ b/go.sum @@ -1,13 +1,11 @@ -git.novox.be/novox/mesh-host v0.0.0-20261006095519-3e80b7ae325e h1:g9h4QRaAMg5yaJLwqtb0FoOs23DVGUYpW6qvnQ3oY5A= -git.novox.be/novox/mesh-host v0.0.0-20261006095519-3e80b7ae325e/go.mod h1:VlilMCRZ5yyNXg7SNigNBLr0Gt32jrGw5KSNq5JAVYs= -git.novox.be/novox/mesh-host v0.0.0-20261007120832-bdd44154ccac h1:KvnKtJ2rWeIE/t4GweK+JL0OjKSNxsrVP3/nMdpii8o= -git.novox.be/novox/mesh-host v0.0.0-20261007120832-bdd44154ccac/go.mod h1:VlilMCRZ5yyNXg7SNigNBLr0Gt32jrGw5KSNq5JAVYs= 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= 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= @@ -44,8 +42,6 @@ github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UV github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= -go.uber.org/automaxprocs v1.6.0 h1:O3y2/QNTOdbF+e/dpXNNW7Rx2hZ4sTIPyybbxyNqTUs= -go.uber.org/automaxprocs v1.6.0/go.mod h1:ifeIMSnPZuznNm6jmdzmU3/bfk01Fe2fotchwEFJ8r8= golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M= golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA= golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= diff --git a/internal/broker/controller_buckets_test.go b/internal/broker/controller_buckets_test.go index 96965471..804f1c5e 100644 --- a/internal/broker/controller_buckets_test.go +++ b/internal/broker/controller_buckets_test.go @@ -1,9 +1,15 @@ package broker import ( + "context" "slices" "strings" "testing" + "time" + + "github.com/nats-io/nats.go/jetstream" + + "github.com/novox/mesh-controller/internal/testbus" ) // **The controller may write every bucket it writes** (novox/hq to-be 45 ยง1, issue 269). Writing a @@ -59,3 +65,34 @@ func TestTheWatchedSignalsMayBeSaidAndHeard(t *testing.T) { t.Error("the controller may not ask who answers, or hears every API call") } } + +// The controller's record of what it asked the operator is bounded (correctness review of 2026-10-08): one +// value a key, a month's age, and a size it cannot outgrow. +func TestWhatTheControllerAskedIsBounded(t *testing.T) { + js, err := Dial(testbus.URL(t)) + if err != nil { + t.Fatal(err) + } + defer js.Close() + if err := js.EnsureControllerBuckets(); err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + kv, err := jetstream.New(js.Conn()) + if err != nil { + t.Fatal(err) + } + bucket, err := kv.KeyValue(ctx, AskedBucket) + if err != nil { + t.Fatal(err) + } + status, err := bucket.Status(ctx) + if err != nil { + t.Fatal(err) + } + info := status.(*jetstream.KeyValueBucketStatus).StreamInfo() + if status.History() != 1 || status.TTL() != AskedKeptFor || info.Config.MaxBytes <= 0 || info.Config.MaxBytes > 64<<20 { + t.Errorf("history %d, age %s, bytes %d", status.History(), status.TTL(), info.Config.MaxBytes) + } +} 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 52dfe165..491d2ba2 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 @@ -289,6 +289,11 @@ func (w Warrant) For(asker string, a Ask) (Option, error) { if w.Level != o.Level { return Option{}, fmt.Errorf("the option %s is %s, and the warrant says %s", o.ID, o.Level, w.Level) } + // A choice made after the ask expired is no answer to it, whatever the router said. + if !w.At.IsZero() && w.At.After(a.Expires) { + return Option{}, fmt.Errorf("the warrant was given at %s, after the ask %s expired at %s", + w.At.UTC().Format(time.RFC3339), a.ID, a.Expires.UTC().Format(time.RFC3339)) + } return o, nil } diff --git a/vendor/modules.txt b/vendor/modules.txt index 406b79c6..9bf6f00c 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -1,4 +1,4 @@ -# git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261008145004-62367ce15ad6 +# git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261008162031-55090da7e08f ## explicit; go 1.22 git.novox.be/novox/mesh-sdk/go/asks # github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op @@ -82,8 +82,6 @@ github.com/nats-io/nuid ## explicit; go 1.26.0 github.com/novox/mesh-host/internal/declaration github.com/novox/mesh-host/validate -# go.uber.org/automaxprocs v1.6.0 -## explicit; go 1.20 # golang.org/x/crypto v0.57.0 ## explicit; go 1.26.0 golang.org/x/crypto/acme