diff --git a/cmd/mesh-controller/asker.go b/cmd/mesh-controller/asker.go new file mode 100644 index 00000000..5896226e --- /dev/null +++ b/cmd/mesh-controller/asker.go @@ -0,0 +1,478 @@ +package main + +// The controller asks, and acts on the operator's warrant (novox/hq ADR 0259 §6). It holds no channel, no +// identity and no factor: it asks the router like any other module, and performs the answer chosen with its +// own grant. +// +// - **For every open, unsilenced condition that needs the operator and names its answers**, one ask is +// published on the `operator-channel` seat under the controller's own name: the condition's words, its +// actions as options at their levels (Silence acknowledges; Release, Stop, Start and Restart approve), +// answered by the operator, expiring after a day (a week when every option only acknowledges). A +// condition that clears, is silenced, or changes its answers has its ask cancelled; an ask that expired +// unanswered is asked again while the condition lasts. Each ask is kept in the controller's bucket +// `asked`, so a restart neither asks twice nor forgets. +// - **On a warrant**, heard on the seat's event under the controller's own name (which only the router may +// say), the controller acts once per ask: only for an ask it holds, only for the option it offered at +// that option's level, and only while the condition is still open. It performs the action as itself — +// a silence through its own conditions, any other through the verb the action names — with the warrant's +// words as its why, and records it in the hand-act log as the operator's decision, naming the channel, +// the ask and the proofs. An ask that ended without a choice is recorded and nothing is done. +// - **A warrant it missed** while away is read from the router's record of its asks, under its own name. + +import ( + "context" + "crypto/rand" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "sort" + "strings" + "sync" + "time" + + "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" +) + +// The asker's name on the seat: the controller's module. +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 + askAcknowledgeFor = 7 * 24 * time.Hour + // askEvery is how often what is open is asked about again, beside every change. + askEvery = time.Minute + // askCatchUpAfter is how old an open ask is before the router's record of it is read: a warrant heard + // on the event needs no reading. + askCatchUpAfter = 2 * time.Minute + // 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 +) + +// What became of an ask, as the controller keeps it. +const ( + askOpen = "open" + askCancelled = "cancelled" +) + +// 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"` + // Options are the actions by option id. + Options map[string]int `json:"options"` + State string `json:"state"` + Opened time.Time `json:"opened"` + Ended time.Time `json:"ended,omitempty"` + Warrant *asks.Warrant `json:"warrant,omitempty"` + // 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"` +} + +// askedStore keeps the asks (broker.AskedBucket). +type askedStore interface { + Get(ctx context.Context, id string) (*asked, error) + Put(ctx context.Context, a asked) error + All(ctx context.Context) ([]asked, error) +} + +// asker is the controller asking the operator and acting on the answer. +type asker struct { + open func(ctx context.Context) ([]conditions.Condition, error) + silence func(ctx context.Context, key string, d time.Duration, by, why string) error + store askedStore + // publish puts a message on a subject's stream, de-duplicated by id. + publish func(ctx context.Context, subject string, body []byte, id string) error + // call performs an action's verb with its arguments, as the controller. + call func(ctx context.Context, a conditions.Action, args map[string]string) error + // record writes the hand-act log. + 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) + + mu sync.Mutex + nudged chan struct{} +} + +func (a *asker) nudge() { + if a == nil { + return + } + a.mu.Lock() + if a.nudged == nil { + a.nudged = make(chan struct{}, 1) + } + ch := a.nudged + a.mu.Unlock() + select { + case ch <- struct{}{}: + default: + } +} + +// keep asks until ctx ends: now, on every change of a condition, and every askEvery. +func (a *asker) keep(ctx context.Context) { + a.nudge() + tick := time.NewTicker(askEvery) + defer tick.Stop() + a.mu.Lock() + nudged := a.nudged + a.mu.Unlock() + for { + select { + case <-ctx.Done(): + return + case <-tick.C: + case <-nudged: + } + if err := a.reconcile(ctx); err != nil { + a.logf("what the operator is asked could not be brought up to date: %v", err) + } + } +} + +// wants says whether a condition is one to ask about now. +func wants(c conditions.Condition, now time.Time) bool { + return len(c.Actions) > 0 && c.Needs != "" && !c.SilencedAt(now) +} + +func sameAsked(a []conditions.Action, b []conditions.Action) bool { + x, _ := json.Marshal(a) + y, _ := json.Marshal(b) + return string(x) == string(y) +} + +// reconcile brings what is asked in line with what is open. +func (a *asker) reconcile(ctx context.Context) error { + now := a.now() + open, err := a.open(ctx) + if err != nil { + return err + } + all, err := a.store.All(ctx) + if err != nil { + return err + } + byCondition := map[string]asked{} + for _, r := range all { + if r.State == askOpen { + if prior, held := byCondition[r.Condition]; !held || r.Opened.After(prior.Opened) { + byCondition[r.Condition] = r + } + } + } + // A warrant missed while away, read from the router's record. + if a.routerRecord != nil { + for _, r := range byCondition { + if now.Sub(r.Opened) < askCatchUpAfter { + continue + } + if w, err := a.routerRecord(ctx, r.ID); err == nil && w != nil { + body, _ := json.Marshal(w) + if err := a.Decided(ctx, body); err != nil { + return err + } + } + } + if all, err = a.store.All(ctx); err != nil { + return err + } + byCondition = map[string]asked{} + for _, r := range all { + if r.State == askOpen { + byCondition[r.Condition] = r + } + } + } + // What the operator answered lately, by condition: not asked again at once. + answered := map[string]asked{} + for _, r := range all { + if r.State == string(asks.OutcomeChosen) && now.Sub(r.Ended) < askAgainAfterAnswer { + answered[r.Condition] = r + } + } + wanted := map[string]bool{} + sort.Slice(open, func(i, j int) bool { return open[i].Key < open[j].Key }) + for _, c := range open { + if !wants(c, now) { + continue + } + wanted[c.Key] = true + if r, done := answered[c.Key]; done && sameAsked(r.Actions, c.Actions) { + if _, held := byCondition[c.Key]; !held { + continue + } + } + 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 + } + default: + continue + } + } + if err := a.ask(ctx, c); err != nil { + a.logf("the operator could not be asked about %s: %v", c.Key, err) + } + } + for key, r := range byCondition { + if !wanted[key] { + if err := a.cancel(ctx, r, "the condition ended, was silenced or needs nothing now"); err != nil { + return err + } + } + } + return nil +} + +// optionID is an action's label as an option's id: "Silence for a week" is silence-for-a-week. +func optionID(label string) string { + var b strings.Builder + dash := false + for _, r := range strings.ToLower(label) { + switch { + case r >= 'a' && r <= 'z', r >= '0' && r <= '9': + b.WriteRune(r) + dash = false + case !dash && b.Len() > 0: + b.WriteByte('-') + dash = true + } + } + return strings.TrimSuffix(b.String(), "-") +} + +// doesWords is what an action does, in the words an option says it with. +func doesWords(act conditions.Action) string { + switch { + case act.Arguments["silence"] != "": + return "nothing more is said of it for a week" + case act.Verb == "mesh-delivery.release": + return "the delivery goes on" + case act.Verb == "mesh-delivery.stop": + return "the delivery ends" + case act.Verb == broker.ControllerSeat+".plans" && act.Arguments["go"] != "": + return "the delivery starts" + case act.Verb == broker.ControllerSeat+".plans" && act.Arguments["stop"] != "": + return "the delivery is stopped" + case strings.HasSuffix(act.Verb, ".restart"): + return "its service is restarted on " + act.Machine + } + return strings.ToLower(act.Label) +} + +// 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, + OnExpiry: "nothing is done, and you are asked again while it lasts", About: c.Key, + 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 + } + approves = approves || level != asks.Acknowledge + oid := optionID(act.Label) + options[oid] = i + q.Options = append(q.Options, asks.Option{ID: oid, Label: act.Label, Does: doesWords(act), Level: level}) + } + q.Expires = now.Add(askAcknowledgeFor) + if approves { + q.Expires = now.Add(askApproveFor) + } + return q, options +} + +func newAskID() string { + var b [8]byte + _, _ = rand.Read(b[:]) + 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) error { + now := a.now() + id := newAskID() + q, options := askOf(id, c, now) + if err := q.Check(now); err != nil { + return err + } + body, err := json.Marshal(q) + if err != nil { + return err + } + if err := a.publish(ctx, asks.AskSubject(askerName), body, "ask."+id); err != nil { + 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}) +} + +// cancel takes an ask back. +func (a *asker) cancel(ctx context.Context, r asked, why string) error { + 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) + 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) +} + +// Decided takes the router's word on one of the controller's asks (link.Decider). An error is returned only +// when what was decided could not be kept, so the word is held and heard again. +func (a *asker) Decided(ctx context.Context, body []byte) error { + var w asks.Warrant + if err := json.Unmarshal(body, &w); err != nil { + a.logf("the router's word on an ask could not be read; ignored: %v", err) + return nil + } + if w.Asker != askerName { + a.logf("REFUSED a warrant for %s's ask %s: the controller acts only on its own", w.Asker, w.Ask) + return nil + } + r, err := a.store.Get(ctx, w.Ask) + if err != nil { + return err + } + if r == nil { + a.logf("REFUSED a warrant for the ask %s, which the controller does not hold", w.Ask) + return nil + } + if r.Acted != "" { + return nil // heard again: acted on once + } + 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) + if w.Words != "" { + r.Acted += ": " + w.Words + } + a.logf("the ask %s about %s ended %s; nothing is done", r.ID, r.Condition, w.Outcome) + return a.store.Put(ctx, *r) + } + option, err := w.For(askerName, r.Ask) + if err != nil { + a.logf("REFUSED a warrant for the ask %s: %v", r.ID, err) + return nil + } + index, offered := r.Options[option.ID] + if !offered || index >= len(r.Actions) { + a.logf("REFUSED a warrant for the ask %s: it chose %s, which no action stands for", r.ID, option.ID) + return nil + } + act := r.Actions[index] + r.State, r.Warrant = string(asks.OutcomeChosen), &w + open, err := a.open(ctx) + if err != nil { + return err + } + stillOpen := false + 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" + 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 { + return err + } + why := fmt.Sprintf("%s (ask %s)", w.Says(), r.ID) + args := map[string]string{} + for k, v := range act.Arguments { + args[k] = v + } + if v, takes := args["why"]; takes && v == "" { + args["why"] = why + } + var acted error + if act.Arguments["silence"] != "" { + acted = a.silence(ctx, act.Arguments["silence"], conditions.MaxSilence, byWords(w), why) + } else { + acted = a.call(ctx, act, args) + } + r.Ended = a.now() + r.Acted = "done" + if acted != nil { + r.Acted = "failed: " + acted.Error() + } + if err := a.store.Put(ctx, *r); err != nil { + return err + } + verbArgs := []string{act.Verb} + if act.Machine != "" { + verbArgs = append(verbArgs, "on "+act.Machine) + } + keys := make([]string, 0, len(args)) + for k := range args { + keys = append(keys, k) + } + sort.Strings(keys) + for _, k := range keys { + if k != "why" { + verbArgs = append(verbArgs, k+"="+args[k]) + } + } + if err := a.record(ctx, link.HandAct{Verb: handActWarrant, Args: verbArgs, Why: why, By: byWords(w), + Cause: conditions.CauseOperatorAnswer, Condition: r.Condition, Via: viaWords(w), Ask: r.ID, + Proofs: w.Proofs, RequestedBy: r.Condition, Outcome: r.Acted}); err != nil { + a.logf("%s was done, and could NOT be recorded in the hand-act log: %v", why, err) + } + a.logf("%s: %s", why, r.Acted) + return nil +} + +// handActWarrant is the verb an act the operator chose on a warrant is recorded under: a person's decision, +// never a repair (handActVerbs). +const handActWarrant = "warrant" + +// byWords is who chose, as the hand-act log says it: "the operator, as telegram identity 42". +func byWords(w asks.Warrant) string { + if w.By == nil { + return "the operator" + } + return fmt.Sprintf("the %s, as %s identity %s", w.By.Who, w.By.Kind, w.By.Identity) +} + +// viaWords is the channel an answer came through: its module and kind, and how the sender was known. +func viaWords(w asks.Warrant) string { + if w.By == nil { + return w.Channel + } + via := w.Channel + " (" + w.By.Kind + ")" + if w.By.Verified != "" { + via += ", " + w.By.Verified + } + return via +} + +// errNotGranted is an action whose verb the controller's grant does not name. +var errNotGranted = errors.New("the controller's grant does not name this verb") diff --git a/cmd/mesh-controller/asker_test.go b/cmd/mesh-controller/asker_test.go new file mode 100644 index 00000000..e3100eb4 --- /dev/null +++ b/cmd/mesh-controller/asker_test.go @@ -0,0 +1,342 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "strings" + "testing" + "time" + + "git.novox.be/novox/mesh-sdk/go/asks" + + "github.com/novox/mesh-controller/internal/conditions" + "github.com/novox/mesh-controller/internal/link" +) + +// novox/hq ADR 0259 §6: the controller asks the operator for the answers its conditions name, and performs +// the one chosen on the router's warrant — once, for its own ask, the option offered, at its level. + +type memAskedStore map[string]asked + +func (m memAskedStore) Get(_ context.Context, id string) (*asked, error) { + 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) All(context.Context) ([]asked, error) { + var out []asked + for _, r := range m { + out = append(out, r) + } + return out, nil +} + +type published struct { + subject, id string + body []byte +} + +type askerRig struct { + a *asker + open []conditions.Condition + store memAskedStore + sent []published + called []string + silenced []string + acts []link.HandAct + now time.Time +} + +func newAskerRig(t *testing.T) *askerRig { + r := &askerRig{store: memAskedStore{}, now: time.Date(2026, 10, 8, 14, 0, 0, 0, time.UTC)} + r.a = &asker{ + open: func(context.Context) ([]conditions.Condition, error) { return r.open, nil }, + silence: func(_ context.Context, key string, d time.Duration, by, why string) error { + r.silenced = append(r.silenced, key+" for "+d.String()+" by "+by+" because "+why) + return nil + }, + store: r.store, + publish: func(_ context.Context, subject string, body []byte, id string) error { + r.sent = append(r.sent, published{subject, id, body}) + return nil + }, + call: func(_ context.Context, a conditions.Action, args map[string]string) error { + raw, _ := json.Marshal(args) + r.called = append(r.called, a.Verb+"@"+a.Machine+" "+string(raw)) + return nil + }, + record: func(_ context.Context, act link.HandAct) error { r.acts = append(r.acts, act); return nil }, + now: func() time.Time { return r.now }, + logf: t.Logf, + } + return r +} + +func heldCondition() conditions.Condition { + o := stalledObservations([]stalledLine{{ID: "novox/hq@055550802096", State: "held", For: "36h2m6s", + Bound: "24h0m0s", H2: "none: the state is the operator's"}})[0] + return conditions.Condition{Key: o.Key(), Kind: o.Kind, Severity: conditions.Warning, Headline: o.Headline, + Explanation: conditions.Verdict(o.Needs, o.Explanation), Needs: o.Needs, Actions: o.Actions} +} + +func unitsCondition() conditions.Condition { + key := "machine.shanks.units" + return conditions.Condition{Key: key, Kind: "machine-units", Severity: conditions.Warning, + Headline: "3 failed services on shanks", Explanation: "Needs you: mend or remove them on shanks, or silence this.", + Needs: "mend or remove them on shanks, or silence this.", Actions: []conditions.Action{conditions.SilenceAction(key)}} +} + +func (r *askerRig) asksSent(t *testing.T) []asks.Ask { + t.Helper() + var out []asks.Ask + for _, p := range r.sent { + if p.subject != asks.AskSubject("mesh-controller") { + continue + } + var q asks.Ask + if err := json.Unmarshal(p.body, &q); err != nil { + t.Fatal(err) + } + out = append(out, q) + } + return out +} + +func TestAnAskIsMadeForEachConditionThatNamesItsAnswers(t *testing.T) { + r := newAskerRig(t) + quiet := conditions.Condition{Key: "machine.ace.silent", Headline: "ace silent", Explanation: "Nothing for you to do. x"} + r.open = []conditions.Condition{heldCondition(), unitsCondition(), quiet} + if err := r.a.reconcile(context.Background()); err != nil { + t.Fatal(err) + } + sent := r.asksSent(t) + if len(sent) != 2 { + t.Fatalf("asked %d times: %+v", len(sent), sent) + } + byAbout := map[string]asks.Ask{} + for _, q := range sent { + byAbout[q.About] = q + if err := q.Check(r.now); err != nil { + t.Errorf("%s: %v", q.About, err) + } + } + held := byAbout[heldCondition().Key] + if len(held.Options) != 2 || held.Options[0].Label != "Release" || held.Options[0].Level != asks.Approve || + held.Options[1].ID != "stop" || held.Expires != r.now.Add(askApproveFor) || held.Who != asks.Operator || + held.OnExpiry == "" { + t.Errorf("the held delivery is asked %+v", held) + } + units := byAbout["machine.shanks.units"] + if len(units.Options) != 1 || units.Options[0].Level != asks.Acknowledge || units.Expires != r.now.Add(askAcknowledgeFor) { + t.Errorf("the failed units are asked %+v", units) + } + // No second ask while one is open. + r.now = r.now.Add(time.Minute) + _ = r.a.reconcile(context.Background()) + if n := len(r.asksSent(t)); n != 2 { + t.Errorf("asked again while open: %d", n) + } +} + +func TestAnAskIsTakenBackWhenItsConditionEndsAndAskedAgainAfterItExpires(t *testing.T) { + r := newAskerRig(t) + r.open = []conditions.Condition{heldCondition(), unitsCondition()} + _ = r.a.reconcile(context.Background()) + // The units are silenced, the held delivery lasts past its ask's day. + units := unitsCondition() + units.Silenced = &conditions.Silence{Until: r.now.Add(48 * time.Hour)} + r.open = []conditions.Condition{heldCondition(), units} + r.now = r.now.Add(askApproveFor) + if err := r.a.reconcile(context.Background()); err != nil { + t.Fatal(err) + } + var cancels int + for _, p := range r.sent { + if p.subject == asks.CancelSubject("mesh-controller") { + cancels++ + } + } + if cancels != 1 { + t.Errorf("cancels %d, want the silenced one's", cancels) + } + if sent := r.asksSent(t); len(sent) != 3 || sent[2].About != heldCondition().Key { + t.Errorf("the expired ask was not asked again: %+v", sent) + } +} + +// warrantFor is the router's warrant for the open ask about a condition, choosing an option by label. +func (r *askerRig) warrantFor(t *testing.T, condition, label string) asks.Warrant { + t.Helper() + for _, a := range r.store { + if a.Condition != condition || a.State != askOpen { + continue + } + for _, o := range a.Ask.Options { + if o.Label == label { + return asks.Warrant{Ask: a.ID, Asker: "mesh-controller", About: condition, Outcome: asks.OutcomeChosen, + Option: o.ID, Label: o.Label, Level: o.Level, Channel: "telegram", Proofs: []string{"P1"}, At: r.now, + By: &asks.Person{Who: asks.Operator, Kind: "telegram", Identity: "42", Verified: "user id verified"}} + } + } + } + t.Fatalf("no open ask about %s offers %s", condition, label) + return asks.Warrant{} +} + +func answerWith(t *testing.T, r *askerRig, w asks.Warrant) { + t.Helper() + body, _ := json.Marshal(w) + if err := r.a.Decided(context.Background(), body); err != nil { + t.Fatal(err) + } +} + +func TestAWarrantIsActedOnOnce(t *testing.T) { + r := newAskerRig(t) + r.open = []conditions.Condition{heldCondition()} + _ = r.a.reconcile(context.Background()) + w := r.warrantFor(t, heldCondition().Key, "Release") + answerWith(t, r, w) + answerWith(t, r, w) // heard again + if len(r.called) != 1 { + t.Fatalf("called %v", r.called) + } + want := `mesh-delivery.release@ {"id":"novox/hq@055550802096","why":"the operator, via telegram (user id verified), chose Release (ask ` + w.Ask + `)"}` + if r.called[0] != want { + t.Errorf("called\n %s\nwant\n %s", r.called[0], want) + } + 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.Via != "telegram (telegram), user id verified" || act.Ask != w.Ask || strings.Join(act.Proofs, ",") != "P1" || + act.Cause != conditions.CauseOperatorAnswer || act.Condition != heldCondition().Key || act.Outcome != "done" { + t.Errorf("the hand-act %+v", act) + } + if !personsDecision(act) { + t.Error("an act on a warrant counts as a repair") + } + if got := r.store[w.Ask]; got.State != string(asks.OutcomeChosen) || got.Acted != "done" { + t.Errorf("kept %+v", got) + } +} + +func TestAWarrantThatIsNotForItsOwnAskIsRefused(t *testing.T) { + for name, change := range map[string]func(*asks.Warrant){ + "another asker": func(w *asks.Warrant) { w.Asker = "mesh-delivery" }, + "an ask not held": func(w *asks.Warrant) { w.Ask = "c0000000000000000" }, + "an option not offered": func(w *asks.Warrant) { w.Option = "delete" }, + "another level": func(w *asks.Warrant) { w.Level = asks.Acknowledge }, + "no person": func(w *asks.Warrant) { w.By = nil }, + } { + t.Run(name, func(t *testing.T) { + r := newAskerRig(t) + r.open = []conditions.Condition{heldCondition()} + _ = r.a.reconcile(context.Background()) + w := r.warrantFor(t, heldCondition().Key, "Stop") + change(&w) + answerWith(t, r, w) + if len(r.called)+len(r.acts)+len(r.silenced) != 0 { + t.Errorf("acted on it: %v %v %v", r.called, r.acts, r.silenced) + } + }) + } +} + +func TestEachAnswerCallsExactlyItsVerb(t *testing.T) { + plan := "plan-1791454185265004861" + waiting := conditions.Condition{Key: "plan." + plan + ".waiting", Severity: conditions.Urgent, + Headline: "openrazer delivery waiting to start", Needs: "start it, or stop it.", + Explanation: "Needs you: start it, or stop it.", Actions: waitingActions(plan, conditions.Urgent)} + module := conditions.Condition{Key: "module.openrazer.g14.unhealthy", Severity: conditions.Warning, + Headline: "openrazer not working on g14", Needs: "restart its service openrazer-daemon on g14.", + Explanation: "Needs you: restart it.", Actions: []conditions.Action{{Label: "Restart", + Verb: "node-service-manager.restart", Machine: "g14", Level: conditions.LevelApprove, + Arguments: map[string]string{"unit": "openrazer-daemon.service", "scope": "user"}}}} + for _, tc := range []struct { + c conditions.Condition + label string + want string + }{ + {waiting, "Start", `mesh-controller.plans@ {"cause":"operator-answer","go":"` + plan + `","why":"`}, + {waiting, "Stop", `mesh-controller.plans@ {"cause":"operator-answer","stop":"` + plan + `","why":"`}, + {module, "Restart", `node-service-manager.restart@g14 {"scope":"user","unit":"openrazer-daemon.service"}`}, + } { + r := newAskerRig(t) + r.open = []conditions.Condition{tc.c} + _ = r.a.reconcile(context.Background()) + answerWith(t, r, r.warrantFor(t, tc.c.Key, tc.label)) + if len(r.called) != 1 || !strings.HasPrefix(r.called[0], tc.want) { + t.Errorf("%s: called %v, want %s…", tc.label, r.called, tc.want) + } + } +} + +func TestASilenceChosenIsTheControllersOwnAndAnAnswerToAnAsk(t *testing.T) { + r := newAskerRig(t) + r.open = []conditions.Condition{unitsCondition()} + _ = r.a.reconcile(context.Background()) + w := r.warrantFor(t, "machine.shanks.units", "Silence for a week") + w.Level, w.Proofs = asks.Acknowledge, nil + w.By = &asks.Person{Who: asks.Operator, Kind: "desktop", Identity: "g14", + Verified: "a desk click: whoever was at the operator's session on g14"} + w.Channel = "desk-channel" + answerWith(t, r, w) + if len(r.called) != 0 || len(r.silenced) != 1 || !strings.HasPrefix(r.silenced[0], "machine.shanks.units for 168h0m0s by the operator, as desktop identity g14") { + t.Fatalf("silenced %v, called %v", r.silenced, r.called) + } + if len(r.acts) != 1 || r.acts[0].Cause != conditions.CauseOperatorAnswer || len(r.acts[0].Proofs) != 0 { + t.Errorf("%+v", r.acts) + } +} + +func TestAnAskThatEndedWithoutAChoiceDoesNothing(t *testing.T) { + r := newAskerRig(t) + r.open = []conditions.Condition{heldCondition()} + _ = r.a.reconcile(context.Background()) + w := r.warrantFor(t, heldCondition().Key, "Release") + w.Outcome, w.Option, w.Label, w.Level, w.By, w.Words = asks.OutcomeExpired, "", "", "", nil, "nobody answered in time" + answerWith(t, r, w) + if len(r.called)+len(r.acts) != 0 || r.store[w.Ask].State != string(asks.OutcomeExpired) || + !strings.HasPrefix(r.store[w.Ask].Acted, "nothing") { + t.Errorf("called %v acts %v kept %+v", r.called, r.acts, r.store[w.Ask]) + } + // And a choice for a condition that ended meanwhile does nothing either. + r2 := newAskerRig(t) + r2.open = []conditions.Condition{heldCondition()} + _ = r2.a.reconcile(context.Background()) + w2 := r2.warrantFor(t, heldCondition().Key, "Release") + r2.open = nil + answerWith(t, r2, w2) + if len(r2.called) != 0 || r2.store[w2.Ask].Acted != "nothing: the condition ended before the answer" { + t.Errorf("%v %+v", r2.called, r2.store[w2.Ask]) + } +} + +func TestAWarrantMissedWhileAwayIsReadFromTheRoutersRecord(t *testing.T) { + r := newAskerRig(t) + r.open = []conditions.Condition{heldCondition()} + _ = r.a.reconcile(context.Background()) + w := r.warrantFor(t, heldCondition().Key, "Stop") + r.a.routerRecord = func(_ context.Context, id string) (*asks.Warrant, error) { + if id != w.Ask { + return nil, errors.New("another ask") + } + return &w, nil + } + r.now = r.now.Add(askCatchUpAfter) + if err := r.a.reconcile(context.Background()); err != nil { + t.Fatal(err) + } + if len(r.called) != 1 || !strings.HasPrefix(r.called[0], "mesh-delivery.stop@") { + t.Errorf("called %v", r.called) + } + if n := len(r.asksSent(t)); n != 1 { + t.Errorf("asked again after the answer: %d", n) + } +} diff --git a/cmd/mesh-controller/asker_wire.go b/cmd/mesh-controller/asker_wire.go new file mode 100644 index 00000000..283bbb44 --- /dev/null +++ b/cmd/mesh-controller/asker_wire.go @@ -0,0 +1,210 @@ +package main + +// The asker on the bus: its asks in the controller's bucket `asked`, its asks and cancels published on the +// seat under the controller's name, the verbs a warrant chooses called with the controller's grant, and the +// router's record of its asks read under its name (novox/hq ADR 0259). + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "strings" + "time" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + + "git.novox.be/novox/mesh-sdk/go/asks" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/conditions" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// askerFrom is the serving controller's asker; nil in any other process. +var askerFrom *asker + +// askWithin is how long a verb a warrant chose is given to answer. +const askWithin = time.Minute + +type busAsked struct{ conn *nats.Conn } + +func (b busAsked) kv(ctx context.Context) (jetstream.KeyValue, error) { + js, err := jetstream.New(b.conn) + if err != nil { + return nil, err + } + return js.KeyValue(ctx, broker.AskedBucket) +} + +func (b busAsked) Get(ctx context.Context, id string) (*asked, error) { + kv, err := b.kv(ctx) + if err != nil { + return nil, err + } + e, err := kv.Get(ctx, id) + if errors.Is(err, jetstream.ErrKeyNotFound) { + return nil, nil + } + if err != nil { + return nil, err + } + var r asked + return &r, json.Unmarshal(e.Value(), &r) +} + +func (b busAsked) Put(ctx context.Context, r asked) error { + kv, err := b.kv(ctx) + if err != nil { + return err + } + body, err := json.Marshal(r) + if err != nil { + return err + } + _, err = kv.Put(ctx, r.ID, body) + return err +} + +func (b busAsked) All(ctx context.Context) ([]asked, error) { + kv, err := b.kv(ctx) + if err != nil { + return nil, err + } + lister, err := kv.ListKeys(ctx) + if err != nil { + return nil, err + } + defer func() { _ = lister.Stop() }() + var out []asked + for k := range lister.Keys() { + e, err := kv.Get(ctx, k) + if err != nil { + continue + } + var r asked + if json.Unmarshal(e.Value(), &r) == nil { + out = append(out, r) + } + } + return out, nil +} + +// callAction performs an action's verb as the controller, through the grant that names it. +func callAction(conn *nats.Conn) func(ctx context.Context, a conditions.Action, args map[string]string) error { + return func(ctx context.Context, a conditions.Action, args map[string]string) error { + seat, verb, ok := strings.Cut(a.Verb, ".") + if !ok { + return fmt.Errorf("%q names no seat and verb", a.Verb) + } + body := map[string]any{} + for k, v := range args { + body[k] = v + } + if seat == catalogue.DeliverySeat { + _, err := askDeliveryOwner(ctx, conn, verb, body) + return err + } + granted := false + for _, v := range broker.VerbsTheControllerActsOnAWarrant { + granted = granted || (v.Seat == seat && v.Verb == verb) + } + if !granted { + return fmt.Errorf("%s: %w", a.Verb, errNotGranted) + } + var answer link.Answer + var err error + if a.Machine != "" { + answer, err = link.AskSeatTool(ctx, conn, seat, verb, a.Machine, body, askWithin) + } else { + answer, err = link.AskMeshSeatTool(ctx, conn, seat, verb, body, askWithin) + } + if err != nil { + return err + } + if answer.Error != "" { + return fmt.Errorf("%s refused: %s", a.Verb, answer.Error) + } + return nil + } +} + +// routerRecordOf reads the router's record of one of the controller's asks, under its name, and answers +// how it ended when it did: the bucket is the one the asks seat's declarer names as its records. +func routerRecordOf(conn *nats.Conn, inv *inventory.Inventory) func(ctx context.Context, id string) (*asks.Warrant, error) { + return func(ctx context.Context, id string) (*asks.Warrant, error) { + bucket, err := asksRecords(ctx, inv) + if err != nil || bucket == "" { + return nil, err + } + reply, err := conn.RequestWithContext(ctx, "$JS.API.DIRECT.GET.KV_"+bucket+".$KV."+bucket+"."+askerName+"."+id, nil) + if err != nil { + return nil, err + } + if reply.Header.Get("Status") != "" { + return nil, nil // none, or not readable: the event says it + } + var rec struct { + State string `json:"state"` + Warrant *asks.Warrant `json:"warrant"` + } + if json.Unmarshal(reply.Data, &rec) != nil || rec.State == "open" || rec.Warrant == nil { + return nil, nil + } + return rec.Warrant, nil + } +} + +// asksRecords is the bucket the asks seat's declarer keeps its record of asks in. +func asksRecords(ctx context.Context, inv *inventory.Inventory) (string, error) { + declared, err := inv.Catalogue(ctx) + if err != nil { + return "", err + } + for _, m := range declared { + for _, s := range m.DefinesSeats { + if s.Name == broker.AsksSeat && len(s.Records) > 0 { + return broker.BucketName(m.Module, s.Records[0]), nil + } + } + } + return "", nil +} + +// startAsking makes the serving controller's asker and hands it the router's words. +func startAsking(ctx context.Context, open *stores, server *link.Server, conn *nats.Conn, keeper *conditions.Keeper) { + js, err := jetstream.New(conn) + if err != nil { + fmt.Printf("the operator cannot be asked: %v\n", err) + return + } + a := &asker{ + open: keeper.Open, + silence: func(ctx context.Context, key string, d time.Duration, by, why string) error { + _, err := keeper.Silence(ctx, key, d, by, why) + return err + }, + store: busAsked{conn: conn}, + publish: func(ctx context.Context, subject string, body []byte, id string) error { + _, err := js.Publish(ctx, subject, body, jetstream.WithMsgID(id)) + return err + }, + call: callAction(conn), + record: func(ctx context.Context, act link.HandAct) error { + _, err := link.RecordHandAct(ctx, conn, act) + return err + }, + routerRecord: routerRecordOf(conn, open.inventory), + now: time.Now, + logf: func(format string, args ...any) { fmt.Printf(format+"\n", args...) }, + } + if err := server.Decides(a); err != nil { + fmt.Printf("the operator's answers cannot be heard, so nothing is asked: %v\n", err) + return + } + askerFrom = a + go a.keep(ctx) +} diff --git a/cmd/mesh-controller/conditions.go b/cmd/mesh-controller/conditions.go index 458382fd..bf0711e9 100644 --- a/cmd/mesh-controller/conditions.go +++ b/cmd/mesh-controller/conditions.go @@ -46,7 +46,7 @@ func keeperOn(ctx context.Context, conn *nats.Conn) (*conditions.Keeper, error) Say: func(format string, args ...any) { fmt.Fprintf(os.Stderr, format+"\n", args...) }, // What status leads with changed: composed again soon (a nudge outside the serving controller // does nothing). - Changed: statusFrom.nudge, + Changed: func() { statusFrom.nudge(); askerFrom.nudge() }, // Written under the lease, carrying its epoch (novox/hq to-be 45 §6). Epoch: func() (uint64, error) { return theLease.epoch(context.WithoutCancel(ctx)) }}), nil } diff --git a/cmd/mesh-controller/delivery_conditions_test.go b/cmd/mesh-controller/delivery_conditions_test.go index 27be9213..916b3ca3 100644 --- a/cmd/mesh-controller/delivery_conditions_test.go +++ b/cmd/mesh-controller/delivery_conditions_test.go @@ -136,12 +136,13 @@ func TestTheDeliveryOwnerIsAskedOverTheBus(t *testing.T) { if err != nil { t.Fatal(err) } - for _, verb := range []string{"stalled", "close"} { + // And release and stop, which the operator's warrant chooses (novox/hq ADR 0259). + for _, verb := range []string{"stalled", "close", "release", "stop"} { if !slices.Contains(granted.Publish, link.SeatToolSubject(catalogue.DeliverySeat, verb)) { t.Errorf("the controller may not ask %s.%s", catalogue.DeliverySeat, verb) } } - if _, err := askDeliveryOwner(t.Context(), nil, "stop", nil); err == nil || !strings.Contains(err.Error(), "grant") { + if _, err := askDeliveryOwner(t.Context(), nil, "retire-history", nil); err == nil || !strings.Contains(err.Error(), "grant") { t.Fatalf("a verb the grant does not name was asked: %v", err) } conn, err := nats.Connect(testbus.URL(t)) diff --git a/cmd/mesh-controller/handacts.go b/cmd/mesh-controller/handacts.go index 19a662b8..0bf2c13e 100644 --- a/cmd/mesh-controller/handacts.go +++ b/cmd/mesh-controller/handacts.go @@ -60,11 +60,18 @@ var handActVerbs = []handActVerb{ // a person's word (ADR 0242), which the push itself reads from what it carried (recorded_push.go). {Verb: "push", Decision: "a recorded build moves only by a person's push: that push is the word its " + "upgrade policy asks for (ADR 0242)", DecidedWhen: pushedRecorded}, - {Verb: "plans stop"}, + // Stopping or starting a walk the operator chose on a warrant (novox/hq ADR 0259) is their decision. + {Verb: "plans stop", Decision: "the operator's answer to an ask is their decision, not a repair (ADR 0259)", + DecidedFor: []string{conditions.CauseOperatorAnswer}}, {Verb: "plans close"}, // A walk started by a person instead of its delivery's owner (novox/hq ADR 0239): the owner down, or - // not trusted with it — either is a repair the owner should have made. - {Verb: "plans go"}, + // not trusted with it — either is a repair the owner should have made. Unless the operator chose it on + // a warrant (ADR 0259). + {Verb: "plans go", Decision: "the operator's answer to an ask is their decision, not a repair (ADR 0259)", + DecidedFor: []string{conditions.CauseOperatorAnswer}}, + // An act the operator chose on a warrant (novox/hq ADR 0259): asked by the controller, answered on a + // channel that proved who answered, performed by the controller as itself. + {Verb: handActWarrant, Decision: "the operator chose it, answering what the controller asked (ADR 0259)"}, {Verb: "broker consumer-reset"}, // Silencing the same condition twice says the condition, or what it watches, wants mending — unless // it is the operator's answer on a notification: a decision to live with it (novox/hq ADR 0258). diff --git a/cmd/mesh-controller/module_health.go b/cmd/mesh-controller/module_health.go index f34f999d..8f3d1f45 100644 --- a/cmd/mesh-controller/module_health.go +++ b/cmd/mesh-controller/module_health.go @@ -383,6 +383,7 @@ func moduleUnhealthyObservation(module, node string, rs []inventory.ResourceHeal Explanation: fmt.Sprintf("%s on %s is not healthy: %s. It clears as soon as it runs again.", module, node, namesWords(plain, 3)), Needs: needs, + Actions: moduleActions(node, rs), Resolved: fmt.Sprintf("%s works again on %s", module, node)} } diff --git a/cmd/mesh-controller/plain_words.go b/cmd/mesh-controller/plain_words.go index 78d75c1e..e123d61f 100644 --- a/cmd/mesh-controller/plain_words.go +++ b/cmd/mesh-controller/plain_words.go @@ -648,13 +648,29 @@ 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. func waitingNeeds(severity conditions.Severity) string { if severity == conditions.Urgent { - return "start it, or stop it, " + FromMeshMCPServer + return "start it, or stop it." } return "" } +// waitingActions are the answers to a walk waiting past its urgent bound: start it, or stop it — the plan's +// own verbs, approved by the operator (novox/hq ADR 0259). None before the bound. +func waitingActions(plan string, severity conditions.Severity) []conditions.Action { + if severity != conditions.Urgent || plan == "" { + return nil + } + return []conditions.Action{ + {Label: "Start", Verb: "mesh-controller.plans", Level: conditions.LevelApprove, + Arguments: map[string]string{"go": plan, "why": "", "cause": conditions.CauseOperatorAnswer}}, + {Label: "Stop", Verb: "mesh-controller.plans", Level: conditions.LevelApprove, + Arguments: map[string]string{"stop": plan, "why": "", "cause": conditions.CauseOperatorAnswer}}, + } +} + // moduleNeeds is what the operator can do about a module unhealthy on a machine: log in again where its // account's groups wait for it (ADR 0252), restart a failed service, or nothing where the mesh restarts it. // No answer is offered for a restart: a desk click performs only an acknowledgement (ADR 0258). @@ -669,11 +685,35 @@ func moduleNeeds(node string, rs []inventory.ResourceHealth) string { } } if unit != "" { - return fmt.Sprintf("restart its service %s on %s %s", unit, node, FromMeshMCPServer) + // Asked of the operator (moduleActions): the router says where it can be answered. + return fmt.Sprintf("restart its service %s on %s.", unit, node) } return "" } +// moduleActions are the answers to a module unhealthy on a machine: restart its failed service there, +// approved by the operator (novox/hq ADR 0259) — none when the mesh restarts it, or a new login is what it +// waits for. +func moduleActions(node string, rs []inventory.ResourceHealth) []conditions.Action { + for _, r := range rs { + if strings.Contains(r.Reason, "relogin needed") { + return nil + } + } + for _, r := range rs { + if r.Kind != link.KindUnit || r.Target == "" { + continue + } + scope := "system" + if r.Account != "" { + scope = "user" + } + return []conditions.Action{{Label: "Restart", Verb: "node-service-manager.restart", Machine: node, + Level: conditions.LevelApprove, Arguments: map[string]string{"unit": r.Target, "scope": scope}}} + } + return nil +} + // FromMeshMCPServer ends what the operator needs when no notification can do it (ADR 0258), naming the mesh MCP // server (the glossary's word; "console" is retired): the answer is not an // acknowledgement, so it is given where the operator is known to be the one asking, until answers are @@ -725,15 +765,19 @@ func stalledWords(l stalledLine, o conditions.Observation) (headline, explanatio long = "for " + humanDuration(d) } if o.Resolver == conditions.ResolverOperator { - // Words only: releasing or stopping a delivery is not an acknowledgement, so no desk click - // performs it (ADR 0258). + // Asked of the operator, approved on a channel that proves who answered (novox/hq ADR 0259); the + // router says where each can be answered, so the words do not. + release := conditions.Action{Label: "Release", Verb: "mesh-delivery.release", Level: conditions.LevelApprove, + Arguments: map[string]string{"id": l.ID, "why": ""}} + stop := conditions.Action{Label: "Stop", Verb: "mesh-delivery.stop", Level: conditions.LevelApprove, + Arguments: map[string]string{"id": l.ID, "why": ""}} switch held { case "held": - needs = "release it, or stop it, " + FromMeshMCPServer + needs, actions = "release it, or stop it.", []conditions.Action{release, stop} case "ready", "checked": needs = "merge its pull request, or close it." default: - needs = "stop it " + FromMeshMCPServer + needs, actions = "stop it.", []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 c0626897..a1c80343 100644 --- a/cmd/mesh-controller/plain_words_test.go +++ b/cmd/mesh-controller/plain_words_test.go @@ -65,13 +65,21 @@ func TestADeliveryWaitingNeedsNothingUntilItsBoundThenOffersStartAndStop(t *test t.Errorf("the summary lost the way on for whoever looks closer: %q", got[0].Summary) } - // Past four hours it is urgent, and offers the controller's own answers. + // Past four hours it is urgent, and asks the operator to start or stop it (novox/hq ADR 0259): the plan's + // own verbs, approved, which the controller performs on the warrant. The router says where to answer. 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, from the mesh MCP server; this notification cannot do it. The change to openrazer is merged and built, and mesh-delivery (the "+ + "Needs you: start it, or stop 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.") + "be stuck.", "Start", "Stop") + for i, want := range []string{"go", "stop"} { + a := got[0].Actions[i] + if a.Verb != "mesh-controller.plans" || a.Arguments[want] != "plan-1791454185265004861" || + a.Level != conditions.LevelApprove || a.Arguments["cause"] != conditions.CauseOperatorAnswer { + t.Errorf("%s: %+v", a.Label, a) + } + } // Many modules are counted, not listed in the headline. f.waits[0].modules = []string{"a", "b", "c", "d"} @@ -82,16 +90,20 @@ func TestADeliveryWaitingNeedsNothingUntilItsBoundThenOffersStartAndStop(t *test } // **A module unhealthy**: "openrazer on g14 is not healthy: its unit openrazer-daemon.service failed in the -// account's own service manager (exit-code)". Restarting is not an acknowledgement, so it is said in words -// and offered as no answer (ADR 0258). +// account's own service manager (exit-code)". Restarting is not an acknowledgement: it is asked of the +// operator at the approve level (novox/hq ADR 0259), so a desk click never performs it (ADR 0258). func TestAModuleUnhealthyAsksForARestartInWords(t *testing.T) { o := moduleUnhealthyObservation("openrazer", "g14", []inventory.ResourceHealth{{Kind: link.KindUnit, - Resource: "openrazer-daemon", Target: "openrazer-daemon.service", + 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 from the mesh MCP server; this notification cannot do it. "+ + "Needs you: restart its service openrazer-daemon on g14. "+ "openrazer on g14 is not healthy: its service openrazer-daemon stopped with an error. It clears as soon "+ - "as it runs again.") + "as it runs again.", "Restart") + if a := o.Actions[0]; a.Verb != "node-service-manager.restart" || a.Machine != "g14" || a.Level != conditions.LevelApprove || + a.Arguments["unit"] != "openrazer-daemon.service" || a.Arguments["scope"] != "user" { + t.Errorf("restart: %+v", a) + } // An account waiting for a new login (ADR 0252) asks for the login, held to the plain rule. o = moduleUnhealthyObservation("openrazer", "g14", []inventory.ResourceHealth{{Kind: "account", Resource: "operator-in-group", Target: "jochen", Reason: "relogin needed: the account is in the group"}}) @@ -153,8 +165,13 @@ 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, from the mesh MCP server; this notification cannot do it. A delivery of hq has been held for 36 hours, past its limit.", - ) + "Needs you: release it, or stop 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 { + t.Errorf("%+v", a) + } + } } // **Every kind the controller raises has plain words**, and its words are plain for a subject of every diff --git a/cmd/mesh-controller/signals.go b/cmd/mesh-controller/signals.go index 875c4940..8be74926 100644 --- a/cmd/mesh-controller/signals.go +++ b/cmd/mesh-controller/signals.go @@ -361,6 +361,7 @@ func watchWaits(f *signalFacts) []conditions.Observation { Headline: deliveryName(w.modules, w.repository) + " waiting to start", Explanation: walkWaitingWords(w, in, severity), Needs: waitingNeeds(severity), + Actions: waitingActions(w.id, severity), Resolved: deliveryName(w.modules, w.repository) + " no longer waiting"}) } return out diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index 6f73c860..e95ab332 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -673,6 +673,9 @@ func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus li // under the lease and the brake, every act said. healers := newHealing(open, keeper, bus, server.JetStream()) go healers.keep(watching) + // And the asker (novox/hq ADR 0259): what needs the operator and names its answers is asked of them, + // and the answer chosen is performed on its warrant. + startAsking(watching, open, server, bus.Conn, keeper) go forgettingOldHeals(watching, open.inventory) fmt.Printf("watching the mesh: %d signal(s) every %s, %d probe(s) every %s; what is wrong is kept in %s "+ "and said as %s events\n", len(watchedRows()), watchEvery, len(runnableProbes()), doctorEvery, diff --git a/go.mod b/go.mod index a84b281a..15a83fb8 100644 --- a/go.mod +++ b/go.mod @@ -11,6 +11,7 @@ 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 diff --git a/go.sum b/go.sum index 04d5c9e7..7ece5e11 100644 --- a/go.sum +++ b/go.sum @@ -4,6 +4,8 @@ git.novox.be/novox/mesh-host v0.0.0-20261007120832-bdd44154ccac h1:KvnKtJ2rWeIE/ 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-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= 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/broker/controller_asks_test.go b/internal/broker/controller_asks_test.go new file mode 100644 index 00000000..777f7c8a --- /dev/null +++ b/internal/broker/controller_asks_test.go @@ -0,0 +1,53 @@ +package broker + +import ( + "slices" + "testing" +) + +// novox/hq ADR 0259 §6: the controller asks the operator through the router's seat as any user of it, under +// its own name, hears its own warrants, reads its own record, and calls the verbs a warrant chooses. +func TestTheControllerAsksUnderItsOwnNameAndCallsTheVerbsAWarrantChooses(t *testing.T) { + records := Records{Nodes: []string{"anchor"}, Assigned: map[string][]Declared{"anchor": { + {Module: "messenger", Holds: []Seat{operatorChannel()}}, + }}} + users, err := Users(records) + if err != nil { + t.Fatal(err) + } + got := perms(t, users[0]) + for _, s := range []string{ + "mesh.seat.operator-channel.accept.ask.mesh-controller", + "mesh.seat.operator-channel.accept.cancel.mesh-controller", + "$JS.API.DIRECT.GET.KV_messenger_asks.$KV.messenger_asks.mesh-controller.c1", + "mesh.seat.mesh-delivery.tool.release", "mesh.seat.mesh-delivery.tool.stop", + "mesh.seat.node-service-manager.tool.restart.g14", "mesh.seat.mesh-controller.tool.plans", + } { + if !allowed(got.Publish, s) { + t.Errorf("the controller may not publish %s", s) + } + } + for _, s := range []string{ + "mesh.seat.operator-channel.accept.ask.mesh-delivery", + "mesh.seat.operator-channel.event.decided.mesh-controller", + // (A direct get of another asker's record is not refused here: the controller holds the whole + // JetStream API, as the only writer of stream definitions.) + "mesh.seat.node-service-manager.tool.stop.g14", + } { + if allowed(got.Publish, s) { + t.Errorf("the controller may publish %s", s) + } + } + if !allowed(got.Subscribe, DecidedSubject) || allowed(got.Subscribe, "mesh.seat.operator-channel.event.decided.mesh-delivery") { + t.Error("the controller does not hear exactly its own warrants") + } + // Its events consumer carries them, so a controller that was away hears what was decided meanwhile. + if !slices.Contains(ControllerFollows, DecidedSubject) { + t.Error("the controller does not follow its warrants") + } + // Without a holder of the seat it is granted no ask at all. + alone, _ := Users(Records{Nodes: []string{"anchor"}, Assigned: map[string][]Declared{}}) + if allowed(perms(t, alone[0]).Publish, "mesh.seat.operator-channel.accept.ask.mesh-controller") { + t.Error("asked a seat nobody holds") + } +} diff --git a/internal/broker/controller_buckets.go b/internal/broker/controller_buckets.go index 12850d5c..69bc81ad 100644 --- a/internal/broker/controller_buckets.go +++ b/internal/broker/controller_buckets.go @@ -34,8 +34,15 @@ var ( // LeaseBucket holds the controller's lease (to-be 45 §6): one key, `holder`, which the instance // allowed to act writes by compare-and-set and renews; its revision when taken is the epoch. LeaseBucket = BucketName(ControllerSeat, "lease") + // AskedBucket keeps what the controller asked the operator about its conditions (novox/hq ADR 0259): + // each ask by its id, its options and the actions they stand for, how it ended and whether the + // controller acted on its warrant — so a restart neither asks twice nor acts twice. + AskedBucket = BucketName(ControllerSeat, "asked") ) +// AskedKeptFor is how long an ask is kept after it was made: a month, as the router keeps its own. +const AskedKeptFor = 30 * 24 * time.Hour + // LeaseTTL is how long the lease's key lives unrenewed (to-be 45 §6): fifteen seconds, renewed // every five. The bucket's age, so the bus forgets a holder that stopped renewing. const LeaseTTL = 15 * time.Second @@ -59,12 +66,12 @@ const ( // IsControllerBucket says a bucket is the controller's own, not a module's state nothing declares. func IsControllerBucket(bucket string) bool { return bucket == CallsBucket || bucket == HandActsBucket || bucket == ConditionsBucket || - bucket == ConditionHistoryBucket || bucket == LeaseBucket + bucket == ConditionHistoryBucket || bucket == LeaseBucket || bucket == AskedBucket } // ControllerBuckets are the controller's own buckets, in the order they are asserted. func ControllerBuckets() []string { - return []string{LeaseBucket, CallsBucket, HandActsBucket, ConditionsBucket, ConditionHistoryBucket} + return []string{LeaseBucket, CallsBucket, HandActsBucket, ConditionsBucket, ConditionHistoryBucket, AskedBucket} } // ControllerBucketsAsserter is what raising the controller's buckets needs of a connection. @@ -149,6 +156,18 @@ func (j *JetStream) EnsureControllerBuckets() error { }); err != nil { return fmt.Errorf("asserting bucket %s: %w", ConditionHistoryBucket, err) } + if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{ + Bucket: AskedBucket, + Description: "what the controller asked the operator about its conditions, and what came of each (novox/hq " + + "ADR 0259): written by the controller alone; an ask acted on is acted on once", + History: 1, + TTL: AskedKeptFor, + MaxValueSize: 32 << 10, + MaxBytes: 32 << 20, + Storage: jetstream.FileStorage, + }); err != nil { + return fmt.Errorf("asserting bucket %s: %w", AskedBucket, err) + } return nil } diff --git a/internal/broker/nats.go b/internal/broker/nats.go index a92dd85a..49da6d59 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -186,8 +186,17 @@ var VerbsTheBusStepAsks = []SeatVerb{{Seat: "node-backup", Verb: "now"}} // VerbsTheControllerAsksTheDeliveryOwner are the mesh-delivery seat's verbs the controller calls (novox/hq // ADR 0239): its self-check reads `stalled`, and healer H2 takes the one transition the table allows // through `close`. A mesh seat's verb is flat: no machine in the subject. +// +// And, since novox/hq ADR 0259, `release` and `stop`: the controller asks the operator for them about a +// delivery held past its bound, and calls them on the operator's warrant, with its why. var VerbsTheControllerAsksTheDeliveryOwner = []SeatVerb{{Seat: "mesh-delivery", Verb: "stalled"}, - {Seat: "mesh-delivery", Verb: "close"}} + {Seat: "mesh-delivery", Verb: "close"}, {Seat: "mesh-delivery", Verb: "release"}, {Seat: "mesh-delivery", Verb: "stop"}} + +// VerbsTheControllerActsOnAWarrant are the other seat verbs the controller calls when the operator's warrant +// chooses them (novox/hq ADR 0259): a machine's service restarted, and a walk started or stopped through the +// controller's own `plans`. Named one by one; a node seat's on any machine, a mesh seat's flat. +var VerbsTheControllerActsOnAWarrant = []SeatVerb{{Seat: "node-service-manager", Verb: "restart"}, + {Seat: ControllerSeat, Verb: "plans"}} // perMachineEvents are a node-scoped seat's events about the holder itself, whose last token is the // holder's machine (novox/hq ADR 0219): `paused.`, the build agent saying whether it takes work. @@ -387,6 +396,20 @@ func PermissionsFor(p Principal) (Permissions, error) { for _, v := range VerbsTheControllerAsksTheDeliveryOwner { pub = append(pub, "mesh.seat."+v.Seat+".tool."+v.Verb) } + // And the verbs a warrant chooses (novox/hq ADR 0259): a node seat's on any machine, its own flat. + for _, v := range VerbsTheControllerActsOnAWarrant { + if v.Seat == ControllerSeat { + pub = append(pub, "mesh.seat."+v.Seat+".tool."+v.Verb) + continue + } + pub = append(pub, "mesh.seat."+v.Seat+".tool."+v.Verb+".*") + } + // And asking the operator (novox/hq ADR 0259): an ask and its cancel under its own name, its warrants + // heard under its own name, the record of its asks read under its own name — as any user of the seat, + // derived the same way, from the seat its holder declares. + tp, ts := SeatTrafficOf(ControllerSeat, nil, p.Uses, nil).grants() + pub = append(pub, tp...) + sub = append(sub, ts...) // And asks who answers (novox/hq to-be 45 §4, D3): the self-check finds every seat's holder by // the same discovery the console reads. The question only; the answers come to its own inbox. pub = append(pub, "$SRV.INFO") diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 4d0a91b5..080e8504 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -270,8 +270,19 @@ var ControllerFollows = []string{ // seat to check before it merges — every machine of the facts snapshot composed with the change. // Appended, because the index is a name. moduleEventSubject("gitea", "pull.updated"), + // **The operator's answers to what the controller asked** (novox/hq ADR 0259): the router's warrant, or + // the end of an ask without one, said to the controller alone under its own name. On the stream, so a + // controller that was away hears what was decided meanwhile. Appended, because the index is a name. + DecidedSubject, } +// AsksSeat is the seat an ask is made on and its warrant heard from (novox/hq ADR 0259): the router's. +const AsksSeat = "operator-channel" + +// DecidedSubject is where the router says the controller's warrants: the seat's event named by the +// controller as its caller. +var DecidedSubject = seatEventSubject(AsksSeat, "decided."+ControllerSeat) + // The provider standing events, by their local names. Written here as well as in the catalogue // (catalogue.ProvisionerEvents), which this package cannot import; a test keeps them agreeing. const ( diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 2891cc56..25ec96cc 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,8 +24,8 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] } - subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.gitea.event.pull.updated", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_asked.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-controller.tool.plans", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.release", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.mesh-delivery.tool.stop", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*", "mesh.seat.node-service-manager.tool.restart.*"] } + subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.gitea.event.pull.updated", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built", "mesh.seat.operator-channel.event.decided.mesh-controller"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { diff --git a/internal/broker/users.go b/internal/broker/users.go index 83924b60..d3ed05a1 100644 --- a/internal/broker/users.go +++ b/internal/broker/users.go @@ -78,7 +78,7 @@ type Records struct { // is a mesh that cannot be told anything, and there is no state of the records in which that is // correct. func Users(r Records) ([]Principal, error) { - out := []Principal{{Kind: KindController}} + out := []Principal{{Kind: KindController, Uses: asksSeatOf(r)}} for _, node := range sortedCopy(r.Nodes) { witness := false @@ -195,3 +195,21 @@ func sortedNames(in map[string][]string) []string { // controllerModule is the controller's module: the machine assigned it witnesses its upgrades. const controllerModule = "mesh-controller" + +// asksSeatOf is the seat an ask is made on, as its holder declares it (novox/hq ADR 0259): the controller +// asks the operator through it like any other user, and is granted what its declaration names for a caller. +// None while nothing holds it. +func asksSeatOf(r Records) []Seat { + for _, node := range sortedCopy(r.Nodes) { + for _, d := range r.Assigned[node] { + for _, s := range d.Holds { + if s.Name == AsksSeat && namesVerb(s.ByCaller, "ask") { + seat := s + seat.Kind, seat.Capabilities = "", nil + return []Seat{seat} + } + } + } + } + return nil +} diff --git a/internal/conditions/plain.go b/internal/conditions/plain.go index 1a2dca67..53bad38e 100644 --- a/internal/conditions/plain.go +++ b/internal/conditions/plain.go @@ -53,8 +53,19 @@ type Action struct { Verb string `json:"verb"` Machine string `json:"machine,omitempty"` Arguments map[string]string `json:"arguments,omitempty"` + // Level is how much proof its answer needs (novox/hq ADR 0234 §8, ADR 0259): LevelAcknowledge for what + // any granted principal may already do, LevelApprove for what only the operator's proven word does. + // The controller asks for every action, and performs the one chosen on the warrant the router issues. + Level string `json:"level,omitempty"` } +// The assurance levels an action's answer needs (novox/hq ADR 0234 §8): acknowledge, approve. Destroy is +// not asked for by any condition: nothing carries its second proof yet. +const ( + LevelAcknowledge = "acknowledge" + LevelApprove = "approve" +) + // The two verdicts an explanation opens with. const ( NothingToDo = "Nothing for you to do." @@ -66,7 +77,7 @@ const ( // only kind of answer a desk click performs until answers are authorised (novox/hq ADR 0258). Its cause // marks it as an answer, which the hand-act log does not count as a repair. func SilenceAction(key string) Action { - return Action{Label: "Silence for a week", Verb: "mesh-controller.conditions", + return Action{Label: "Silence for a week", Verb: "mesh-controller.conditions", Level: LevelAcknowledge, Arguments: map[string]string{"silence": key, "for": "7d", "why": "", "cause": CauseOperatorAnswer}} } diff --git a/internal/link/contracts.go b/internal/link/contracts.go index b3acb7c7..381fdf4f 100644 --- a/internal/link/contracts.go +++ b/internal/link/contracts.go @@ -47,6 +47,10 @@ var Contracts = map[string]Contract{ KindPullUpdated: {Unordered: "a pull request's head, asked to be checked: each head is its own commit, and its " + "verdict is set on that commit alone, so a head heard late is checked and judged as itself and never " + "stands for a newer one (novox/hq to-be 45 §9)"}, + KindDecided: {Unordered: "the router's word on one ask, by the ask's id: an ask is ended once, by compare-and-set " + + "at the router, and the controller acts on it once, recording that it did under the ask's id — so a word " + + "heard again, or late, does nothing more (novox/hq ADR 0259)", + Tests: []string{"TestAWarrantIsActedOnOnce"}}, KindCatchUp: {Unordered: "a catalogue asking what it missed: answered from the record, whenever asked"}, KindProvisioner: {Unordered: "a provider's newest word about a consumer, said again every fifteen minutes " + "while it holds (ADR 0224): the condition keeps the last observed, and S8 says when the words stop. " + diff --git a/internal/link/handacts.go b/internal/link/handacts.go index e4e916c6..dd42f6e0 100644 --- a/internal/link/handacts.go +++ b/internal/link/handacts.go @@ -49,6 +49,16 @@ type HandAct struct { Kind string `json:"kind,omitempty"` // Carried is what such a push moved, one "module from → to" per module, so the log says it. Carried []string `json:"carried,omitempty"` + // Via, Ask, Proofs and RequestedBy are an act the operator chose on a warrant (novox/hq ADR 0234 §8, ADR + // 0259): the channel it came through (module and kind, and how the sender was known), the ask's id, the + // proofs present (P1, P2, P3), and what asked (a condition's key). By then names the operator as that + // kind's identity. Absent from every other act. + Via string `json:"via,omitempty"` + Ask string `json:"ask,omitempty"` + Proofs []string `json:"proofs,omitempty"` + RequestedBy string `json:"requested-by,omitempty"` + // Outcome is what came of an act recorded after it was done: done, or the verb's refusal. + Outcome string `json:"outcome,omitempty"` } // KindRecordedBuilds is a push that only moved recorded builds: the person's word their `record` policy diff --git a/internal/link/receive.go b/internal/link/receive.go index 3e4cbdb7..0173777a 100644 --- a/internal/link/receive.go +++ b/internal/link/receive.go @@ -44,6 +44,9 @@ const ( // KindPullUpdated is the forge announcing a pull request's new head: checked before it merges // (novox/hq to-be 45 §9). KindPullUpdated = "pull-updated" + // KindDecided is the router's warrant for an ask the controller made, or that ask's end without one + // (novox/hq ADR 0259): said to the controller alone, under its own name. + KindDecided = "decided" ) // Control is one thing a node or a module said, as the controller must act on it. diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index f3d9d5e0..64671c2f 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -53,7 +53,7 @@ func Nats(js *broker.JetStream) Inbound { // whatever was asked for — and not at all when nothing was. func (n *natsInbound) Also(kind string) error { switch kind { - case KindModuleMoved, KindCatchUp, KindSourceMoved, KindProvisioner, KindPullUpdated: + case KindModuleMoved, KindCatchUp, KindSourceMoved, KindProvisioner, KindPullUpdated, KindDecided: n.follows[kind] = true return nil default: @@ -259,6 +259,8 @@ func kindOfSubject(subject string) (string, bool) { return KindSourceMoved, true case PullUpdatedSubject: return KindPullUpdated, true + case broker.DecidedSubject: + return KindDecided, true case BuildOutcome(), BuildOutcomeOf(TheBuildMachineBefore): // A build's outcome is the role's event now, so it arrives on the events stream rather than // the control branch — and is acted on by the same handler, because what the controller does diff --git a/internal/link/serve.go b/internal/link/serve.go index f07627e7..0565e6f5 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -70,7 +70,7 @@ type Checker interface { } // PullUpdatedSubject is where the forge's pull requests land: the controller's own follow of them. -var PullUpdatedSubject = broker.ControllerFollows[len(broker.ControllerFollows)-1] +var PullUpdatedSubject = "mesh.mod.gitea.event.pull.updated" // Server acts on what nodes and modules say. // @@ -96,6 +96,8 @@ type Server struct { checker Checker // healths keeps what machines say of their long-running resources (novox/hq ADR 0240). healths Healths + // decider acts on the operator's warrants for what the controller asked (novox/hq ADR 0259). + decider Decider log *log.Logger // giveUp is how long one message is held for the store; zero means GiveUpAfter. @@ -138,6 +140,21 @@ func (s *Server) Checks(c Checker) error { return nil } +// Decider is what the controller does with the router's word on an ask it made (novox/hq ADR 0259): act +// on a warrant once, or record how the ask ended without one. +type Decider interface { + Decided(ctx context.Context, body []byte) error +} + +// Decides says what to do about the router's word on the controller's asks, and asks for it delivered. +func (s *Server) Decides(d Decider) error { + if err := s.inbound.Also(KindDecided); err != nil { + return err + } + s.decider = d + return nil +} + // Answers says what to do about a catalogue's catch-up request, and asks for them to be delivered. func (s *Server) Answers(r Replayer) error { if err := s.inbound.Also(KindCatchUp); err != nil { @@ -210,6 +227,8 @@ func (s *Server) act(ctx context.Context, m Control) { s.provisioner(ctx, m) case KindPullUpdated: s.pullUpdated(ctx, m) + case KindDecided: + s.decided(ctx, m) default: // Dropped: a message nothing understands will not be understood on the next attempt // either, and asking for it again would spin. @@ -656,6 +675,24 @@ func (s *Server) pullUpdated(ctx context.Context, m Control) { _ = m.Took() } +// decided hands the router's word on an ask to the decider; a failure to keep what it did is held for the +// store, like any word that must not be lost. +func (s *Server) decided(ctx context.Context, m Control) { + if s.decider == nil { + _ = m.Took() + return + } + err := s.decider.Decided(ctx, m.Body()) + switch s.decide(ctx, m, "the operator's word on an ask", "", "", err) { + case Hold: + return + } + if err != nil { + s.log.Printf("the operator's word on an ask could not be kept: %v", err) + } + _ = m.Took() +} + // saysWhatItDid states what a machine now runs, or what it would not take, as a fact on the bus // (novox/hq ADR 0134). // 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 new file mode 100644 index 00000000..52dfe165 --- /dev/null +++ b/vendor/git.novox.be/novox/mesh-sdk/go/asks/asks.go @@ -0,0 +1,377 @@ +// Package asks is the contract of asking a person and answering on a channel (novox/hq ADR 0259): the +// shapes an asker, the router and a channel exchange on the bus, and the subjects they travel on. No +// transport and no channel's service: an asker publishes an Ask under its own name and acts on the Warrant +// it hears; the router holds the ask, sends channels a Message, and judges the Choice a channel says; a +// channel shows a Message and says what was chosen and by whom, as its service authenticated it. +package asks + +import ( + "errors" + "fmt" + "regexp" + "strings" + "time" +) + +// The seats (novox/hq ADR 0259 §3). +const ( + // Seat is held by the router: an ask and its cancel are its accepts, a warrant its event, each named by + // the asker. + Seat = "operator-channel" + // ChannelSeat is the kinded bench a channel holds to show and say: its accepts carry the kind. + ChannelSeat = "channel" + // IntakeSeat is the kinded bench a channel holds to say what was chosen: its events and proofs carry + // the kind. + IntakeSeat = "intake" +) + +// AskSubject is where an asker publishes an ask, CancelSubject its cancel, and DecidedSubject where it +// hears the warrant, or the ask's end without one. +func AskSubject(asker string) string { return "mesh.seat." + Seat + ".accept.ask." + asker } +func CancelSubject(asker string) string { return "mesh.seat." + Seat + ".accept.cancel." + asker } +func DecidedSubject(asker string) string { return "mesh.seat." + Seat + ".event.decided." + asker } + +// The work a channel takes, on ChannelSubject. +const ( + Show = "show" // a message offering answers + Edit = "edit" // a message shown before, replaced + Send = "send" // a message offering nothing +) + +// ChannelSubject is where the router sends a channel of a kind its work. +func ChannelSubject(verb, kind string) string { + return "mesh.seat." + ChannelSeat + ".accept." + verb + "." + kind +} + +// What a channel says, on IntakeSubject. +const ( + Chosen = "choice" // a button tapped + Link = "link" // somebody asked to be linked as the operator +) + +// IntakeSubject is where a channel of a kind says what arrived. +func IntakeSubject(what, kind string) string { + return "mesh.seat." + IntakeSeat + ".event." + what + "." + kind +} + +// CodeProof is the one proof verb: a code the operator typed, carried by request and reply, never kept. +const CodeProof = "code" + +// ProofSubject is where a channel of a kind asks a proof. +func ProofSubject(verb, kind string) string { + return "mesh.seat." + IntakeSeat + ".proof." + verb + "." + kind +} + +// A Level is how much proof an option's answer needs (novox/hq ADR 0234 §8, the glossary's assurance level). +type Level string + +const ( + // Acknowledge performs only what any granted principal may already do: silencing, details. No proof. + Acknowledge Level = "acknowledge" + // Approve needs one proof: a verified sender (a linked account, linked an hour or more) or a code. + Approve Level = "approve" + // Destroy needs two proofs, one of them a code. + Destroy Level = "destroy" +) + +// Rank orders the levels; an unknown level ranks above every known one, so it is never taken as less. +func (l Level) Rank() int { + switch l { + case Acknowledge: + return 0 + case Approve: + return 1 + case Destroy: + return 2 + } + return 3 +} + +// Operator is the one role an ask may be answered by today. +const Operator = "operator" + +// Option is one answer an ask offers: a label for the button, what it does in plain words, its level. +type Option struct { + ID string `json:"id"` + Label string `json:"label"` + Does string `json:"does"` + Level Level `json:"level"` +} + +// Ask is a request for a person's word (novox/hq ADR 0259 §4). +type Ask struct { + // ID is the asker's own, unique to it. + ID string `json:"id"` + // Headline names the thing and what is wrong, in a few plain words; Explanation is what happened and + // what it means. Both are held to the plain rule and the content rule by the router. + Headline string `json:"headline"` + Explanation string `json:"explanation"` + Options []Option `json:"options"` + // Who may answer: Operator. + Who string `json:"who"` + // Expires is when the ask ends unanswered; OnExpiry is what the asker then does, in words the person + // is shown ("the delivery stays held"). An ask that authorises never defaults. + Expires time.Time `json:"expires"` + OnExpiry string `json:"on-expiry"` + // About is what the ask is about (a condition's key): a newer ask about it replaces the older. + About string `json:"about,omitempty"` + Urgent bool `json:"urgent,omitempty"` +} + +// The bounds of an ask (novox/hq ADR 0234 §8, ADR 0259 §4). +const ( + MostOptions = 4 + MostOpen = 3 + ApproveLasts = 24 * time.Hour + DestroyLasts = 10 * time.Minute + HeadlineLength = 60 + LabelLength = 24 +) + +var usableID = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9_-]{0,63}$`) + +// UsableID says whether a name can be an ask's or an option's id: a key and a subject token both. +func UsableID(id string) bool { return usableID.MatchString(id) } + +// Highest is the highest level among the ask's options. +func (a Ask) Highest() Level { + high := Acknowledge + for _, o := range a.Options { + if o.Level.Rank() > high.Rank() { + high = o.Level + } + } + return high +} + +// Option is the option of this id, or false. +func (a Ask) Option(id string) (Option, bool) { + for _, o := range a.Options { + if o.ID == id { + return o, true + } + } + return Option{}, false +} + +// Check is what an ask is held to before anything is shown: every refusal, in words its asker can act on. +func (a Ask) Check(now time.Time) error { + var problems []string + say := func(format string, args ...any) { problems = append(problems, fmt.Sprintf(format, args...)) } + if !UsableID(a.ID) { + say("its id %q is not letters, digits, - and _, at most 64", a.ID) + } + if strings.TrimSpace(a.Headline) == "" || len([]rune(a.Headline)) > HeadlineLength { + say("its headline is empty or longer than %d characters", HeadlineLength) + } + if strings.TrimSpace(a.Explanation) == "" { + say("it explains nothing") + } + if a.Who != Operator { + say("it is answered by %q, and only the operator answers today", a.Who) + } + if len(a.Options) == 0 || len(a.Options) > MostOptions { + say("it offers %d options, and an ask offers one to %d", len(a.Options), MostOptions) + } + seen := map[string]bool{} + for _, o := range a.Options { + switch { + case !UsableID(o.ID): + say("an option's id %q is not letters, digits, - and _", o.ID) + case seen[o.ID]: + say("the option %s is offered twice", o.ID) + } + seen[o.ID] = true + if strings.TrimSpace(o.Label) == "" || len([]rune(o.Label)) > LabelLength { + say("the option %s's label is empty or longer than %d characters", o.ID, LabelLength) + } + if strings.TrimSpace(o.Does) == "" { + say("the option %s does not say what it does", o.ID) + } + if o.Level.Rank() > Destroy.Rank() { + say("the option %s has the level %q, which is none of acknowledge, approve, destroy", o.ID, o.Level) + } + } + if !a.Expires.After(now) { + say("it expires before it is asked") + } + switch a.Highest() { + case Approve: + if a.Expires.After(now.Add(ApproveLasts)) { + say("an ask that approves lasts at most %s", ApproveLasts) + } + case Destroy: + if a.Expires.After(now.Add(DestroyLasts)) { + say("an ask that destroys lasts at most %s", DestroyLasts) + } + } + if a.Highest() != Acknowledge && strings.TrimSpace(a.OnExpiry) == "" { + say("it does not say what happens when nobody answers, and an ask that authorises never defaults") + } + if a.About != "" && strings.ContainsAny(a.About, " \n") { + say("what it is about is a key, without spaces") + } + if len(problems) > 0 { + return errors.New("the ask is refused: " + strings.Join(problems, "; ")) + } + return nil +} + +// Outcome is how an ask ended. +type Outcome string + +const ( + OutcomeChosen Outcome = "chosen" // a person chose an option: a warrant + OutcomeExpired Outcome = "expired" // nobody answered in time + OutcomeCancelled Outcome = "cancelled" // its asker took it back + OutcomeReplaced Outcome = "replaced" // a newer ask about the same thing replaced it + OutcomeRefused Outcome = "refused" // it was never shown: Words says why +) + +// Person is who chose, as the router verified them. +type Person struct { + // Who is the role: Operator. + Who string `json:"who"` + // Kind is the channel kind they answered on, Identity their account on that service, Display the name + // the service shows, and Verified how the router knew it was them. + Kind string `json:"kind"` + Identity string `json:"identity"` + Display string `json:"display,omitempty"` + Verified string `json:"verified"` +} + +// Warrant is the router's record that a person chose one option of one ask, or the ask's end without one +// (novox/hq ADR 0259 §6). It carries no secret. +type Warrant struct { + Ask string `json:"ask"` + Asker string `json:"asker"` + About string `json:"about,omitempty"` + Outcome Outcome `json:"outcome"` + // Option, Label and Level are the option chosen; By who chose it, Channel the module it came through, + // Proofs which proofs were present (P1, P2, P3). + Option string `json:"option,omitempty"` + Label string `json:"label,omitempty"` + Level Level `json:"level,omitempty"` + By *Person `json:"by,omitempty"` + Channel string `json:"channel,omitempty"` + Proofs []string `json:"proofs,omitempty"` + At time.Time `json:"at"` + // Words are why an ask ended without a choice, or what refused it. + Words string `json:"words,omitempty"` +} + +// Says is the warrant in the words an asker records with its act: "the operator, via telegram (user id +// verified), chose Release". +func (w Warrant) Says() string { + if w.Outcome != OutcomeChosen || w.By == nil { + return fmt.Sprintf("no person chose: the ask %s %s", w.Ask, w.Outcome) + } + via := w.By.Kind + if w.By.Verified != "" { + via += " (" + w.By.Verified + ")" + } + return fmt.Sprintf("the %s, via %s, chose %s", w.By.Who, via, w.Label) +} + +// For checks a warrant against the ask its asker made: the same asker and ask, a choice, an option the ask +// offered, at that option's level. An asker acts on nothing else. +func (w Warrant) For(asker string, a Ask) (Option, error) { + if w.Asker != asker || w.Ask != a.ID { + return Option{}, fmt.Errorf("the warrant is for %s's ask %s, not %s's %s", w.Asker, w.Ask, asker, a.ID) + } + if w.Outcome != OutcomeChosen || w.By == nil || w.By.Who != a.Who { + return Option{}, fmt.Errorf("the ask %s ended %s; no person chose", a.ID, w.Outcome) + } + o, offered := a.Option(w.Option) + if !offered { + return Option{}, fmt.Errorf("the ask %s offered no option %s", a.ID, w.Option) + } + 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) + } + return o, nil +} + +// Button is one answer a channel offers: its label and the router's one-time ticket for it. +type Button struct { + Label string `json:"label"` + Ticket string `json:"ticket"` +} + +// Message is the work a channel takes: shown with buttons (Show), shown again in place (Edit), or said +// (Send). Handle is the router's name for it, the same across a show and its edits; the channel keeps +// which of its own messages that is. The words are the router's, shown as given. +type Message struct { + Handle string `json:"handle"` + Title string `json:"title"` + Body string `json:"body"` + Buttons []Button `json:"buttons,omitempty"` + Urgent bool `json:"urgent,omitempty"` + Silent bool `json:"silent,omitempty"` + // Reply is the Choice or LinkAsked this answers, by its ID: the channel shows it where that was made. + Reply string `json:"reply,omitempty"` + // To is the account linked as the operator on this kind, for a channel that verifies its sender: where + // the channel sends what is not a reply. The router's word, from its list; empty when none is linked. + To string `json:"to,omitempty"` +} + +// Sender is who a channel's service says sent something: the account, the name it shows, and whether the +// service authenticated it. The router alone judges whether that is the operator. +type Sender struct { + Identity string `json:"identity"` + Display string `json:"display,omitempty"` + Authenticated bool `json:"authenticated"` +} + +// Failed is a message a channel could not deliver: its handle, why, and whether trying again could help. +type Failed struct { + Handle string `json:"handle"` + Why string `json:"why"` + Permanent bool `json:"permanent,omitempty"` + At time.Time `json:"at"` +} + +// Standing is what a channel says of itself, at least every five minutes and whenever it changes: +// whether it can send now, why not, and whether its edits notify nobody. +type Standing struct { + Ready bool `json:"ready"` + Why string `json:"why,omitempty"` + EditsSilently bool `json:"edits-silently,omitempty"` + At time.Time `json:"at"` +} + +// What a channel says of its own delivery, on IntakeSubject. +const ( + FailedWhat = "failed" + StandingWhat = "standing" +) + +// Choice is a button chosen on a channel. +type Choice struct { + // ID is the channel's own for this arrival, unique, so the router acts on it once. + ID string `json:"id"` + Ticket string `json:"ticket"` + Handle string `json:"handle,omitempty"` + Sender Sender `json:"sender"` + At time.Time `json:"at"` +} + +// LinkAsked is somebody on a channel asking to be linked as the operator. +type LinkAsked struct { + ID string `json:"id"` + Sender Sender `json:"sender"` + At time.Time `json:"at"` +} + +// Code is a code a person typed on a channel, asked as a proof: never in an event, never kept. +type Code struct { + Sender Sender `json:"sender"` + Code string `json:"code"` + Purpose string `json:"purpose"` +} + +// ProofAnswer is the router's answer to a proof: whether it was taken, and words to say to the person. +type ProofAnswer struct { + Accepted bool `json:"accepted"` + Words string `json:"words"` +} diff --git a/vendor/modules.txt b/vendor/modules.txt index 6eb9699c..cdf7769c 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -1,3 +1,6 @@ +# git.novox.be/novox/mesh-sdk/go v0.1.8-0.20261008145004-62367ce15ad6 +## explicit; go 1.22 +git.novox.be/novox/mesh-sdk/go/asks # github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op ## explicit; go 1.24.0 github.com/antithesishq/antithesis-sdk-go/assert