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