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