938 lines
35 KiB
Go
938 lines
35 KiB
Go
package main
|
|
|
|
// The controller asks, and acts on the operator's warrant (novox/hq ADR 0259 §6). It holds no channel, no
|
|
// identity and no factor: it asks the router like any other module, and performs the answer chosen with its
|
|
// own grant.
|
|
//
|
|
// - **For every open, unsilenced condition that needs the operator and names its answers**, one ask is
|
|
// published on the `operator-channel` seat under the controller's own name: the condition's words, its
|
|
// actions as options at their levels (Silence acknowledges; Release, Stop, Start and Restart approve),
|
|
// answered by the operator, expiring after a day (a week when every option only acknowledges). A
|
|
// condition that clears, is silenced, or changes its answers has its ask cancelled; an ask that expired
|
|
// unanswered is asked again while the condition lasts; one the router refused is asked again after a wait
|
|
// that grows with each refusal in a row, at most half an hour, so a refusal is never believed longer. Each ask is kept in the controller's bucket
|
|
// `asked`, so a restart neither asks twice nor forgets.
|
|
// - **On a warrant**, heard on the seat's event under the controller's own name (which only the router may
|
|
// say), the controller acts once per ask: only for an ask it holds, only for the option it offered at
|
|
// that option's level, and only while the condition is still open. It performs the action as itself —
|
|
// a silence through its own conditions, any other through the verb the action names — with the warrant's
|
|
// words as its why, and records it in the hand-act log as the operator's decision, naming the channel,
|
|
// the ask and the proofs. An ask that ended without a choice is recorded and nothing is done.
|
|
// - **A warrant it missed** while away is read from the router's record of its asks, under its own name.
|
|
|
|
import (
|
|
"context"
|
|
"crypto/rand"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"git.novox.be/novox/mesh-sdk/go/asks"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"github.com/novox/mesh-controller/internal/conditions"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// The asker's name on the seat: the controller's module.
|
|
const askerName = broker.ControllerSeat
|
|
|
|
// How long an ask lasts: a day when an answer approves, a week when every answer only acknowledges.
|
|
const (
|
|
// askApproveFor 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
|
|
// askCatchUpAfter is how old an open ask is before the router's record of it is read: a warrant heard
|
|
// on the event needs no reading.
|
|
askCatchUpAfter = 2 * time.Minute
|
|
// askAgainAfterAnswer is how long a condition the operator answered is not asked about again with the
|
|
// same answers: what was chosen takes a while to clear it, and asking again at once would ask twice.
|
|
askAgainAfterAnswer = time.Hour
|
|
// 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
|
|
// askRefusedRetryMost is the longest a refusal is believed without asking again (novox/hq issue 369): the
|
|
// router refused while its channels had not yet said they could send, the channels could send twenty
|
|
// minutes later, and the controller repeated that refusal for eleven hours. A refused ask is asked again
|
|
// after askEvery, then twice as long after each refusal in a row, never longer than this — and at once when
|
|
// the channels change.
|
|
askRefusedRetryMost = 30 * time.Minute
|
|
// askRefusedCounted is how far back refusals in a row are counted for that wait.
|
|
askRefusedCounted = 6 * time.Hour
|
|
// askVerdictWait is how long an ask made again after a refusal waits for the router's word before it is
|
|
// taken as taken: the router refuses an ask as it reads it, so a refusal comes within seconds. Meanwhile
|
|
// the condition saying it was not delivered stands as it was, neither cleared nor raised again.
|
|
askVerdictWait = 2 * time.Minute
|
|
)
|
|
|
|
// refusedRetryAfter is how long a part refused n times in a row waits before it is asked again.
|
|
func refusedRetryAfter(n int) time.Duration {
|
|
wait := askEvery
|
|
for i := 1; i < n && wait < askRefusedRetryMost; i++ {
|
|
wait *= 2
|
|
}
|
|
return min(wait, askRefusedRetryMost)
|
|
}
|
|
|
|
// What became of an ask, as the controller keeps it.
|
|
const (
|
|
askOpen = "open"
|
|
askCancelled = "cancelled"
|
|
)
|
|
|
|
// asked is one ask the controller made, as it keeps it.
|
|
type asked struct {
|
|
ID string `json:"id"`
|
|
Condition string `json:"condition"`
|
|
// Channels is what the channels were when it was asked (asker.channels): an ask the router refused is asked
|
|
// again at once when the condition's answers or the channels change, and otherwise after a wait that grows
|
|
// with each refusal in a row (refusedRetryAfter, novox/hq issue 369).
|
|
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"`
|
|
Opened time.Time `json:"opened"`
|
|
Ended time.Time `json:"ended,omitempty"`
|
|
Warrant *asks.Warrant `json:"warrant,omitempty"`
|
|
// Acted is what the controller did on the warrant: empty before it did anything, "acting" while it acts,
|
|
// then "done", "failed: …" or "nothing: …". Anything but empty is never acted on again.
|
|
Acted string `json:"acted,omitempty"`
|
|
// 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.
|
|
Rehearsal bool `json:"rehearsal,omitempty"`
|
|
// Proposal is a settings layer proposed through a verb (proposals.go, novox/hq ADR 0277): about no
|
|
// condition, set on Approve by the serving controller itself, and left alone by the reconciling of conditions.
|
|
Proposal *settingsProposal `json:"proposal,omitempty"`
|
|
}
|
|
|
|
// ofACondition says an ask is one of a condition's: not a rehearsal, not a proposal.
|
|
func (r asked) ofACondition() bool { return !r.Rehearsal && r.Proposal == nil }
|
|
|
|
// 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)
|
|
// 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)
|
|
}
|
|
|
|
// 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)
|
|
silence func(ctx context.Context, key string, d time.Duration, by, why string) error
|
|
store askedStore
|
|
// publish puts a message on a subject's stream, de-duplicated by id.
|
|
publish func(ctx context.Context, subject string, body []byte, id string) error
|
|
// call performs an action's verb with its arguments, as the controller.
|
|
call func(ctx context.Context, a conditions.Action, args map[string]string) error
|
|
// setLayer sets a proposed settings layer on the operator's warrant (novox/hq ADR 0277); nil cannot.
|
|
setLayer setOnWarrant
|
|
// record writes the hand-act log.
|
|
record func(ctx context.Context, act link.HandAct) error
|
|
// routerRecord reads the router's record of an ask for a warrant missed; nil reads nothing.
|
|
routerRecord func(ctx context.Context, id string) (*asks.Warrant, error)
|
|
// 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)
|
|
// grantHeld says whether the bus holds the controller's grant to ask: the user list the bus's machine was
|
|
// last sent is the one the mesh composes now (novox/hq issue 353). The grant is composed from the router's
|
|
// assignment and reaches the bus only when that machine is next pushed, so between `assign` and `push`
|
|
// the record says a router is here and the bus refuses every ask. why says what to do; nil is yes.
|
|
grantHeld func(ctx context.Context) (held bool, why string, err error)
|
|
// channels is what the channels are now, as a fingerprint: who holds which kind, promising what.
|
|
channels func(ctx context.Context) string
|
|
// raise keeps the asker's own condition (sourceAsker): which conditions needing the operator could not be
|
|
// asked, and why. Nil raises nothing (a test that does not look).
|
|
raise func(ctx context.Context, obs []conditions.Observation) error
|
|
now func() time.Time
|
|
logf func(string, ...any)
|
|
|
|
saidNoRouter bool
|
|
saidNoGrant bool
|
|
|
|
mu sync.Mutex
|
|
nudged chan struct{}
|
|
}
|
|
|
|
func (a *asker) nudge() {
|
|
if a == nil {
|
|
return
|
|
}
|
|
a.mu.Lock()
|
|
if a.nudged == nil {
|
|
a.nudged = make(chan struct{}, 1)
|
|
}
|
|
ch := a.nudged
|
|
a.mu.Unlock()
|
|
select {
|
|
case ch <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// keep asks until ctx ends: now, on every change of a condition, and every askEvery.
|
|
func (a *asker) keep(ctx context.Context) {
|
|
a.nudge()
|
|
tick := time.NewTicker(askEvery)
|
|
defer tick.Stop()
|
|
a.mu.Lock()
|
|
nudged := a.nudged
|
|
a.mu.Unlock()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-tick.C:
|
|
case <-nudged:
|
|
}
|
|
if err := a.reconcile(ctx); err != nil {
|
|
a.logf("what the operator is asked could not be brought up to date: %v", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// wants says whether a condition is one to ask about now.
|
|
func wants(c conditions.Condition, now time.Time) bool {
|
|
return len(c.Actions) > 0 && c.Needs != "" && !c.SilencedAt(now)
|
|
}
|
|
|
|
func sameAsked(a []conditions.Action, b []conditions.Action) bool {
|
|
x, _ := json.Marshal(a)
|
|
y, _ := json.Marshal(b)
|
|
return string(x) == string(y)
|
|
}
|
|
|
|
// reconcile brings what is asked in line with what is open.
|
|
func (a *asker) reconcile(ctx context.Context) error {
|
|
now := a.now()
|
|
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
|
|
}
|
|
open, err := a.open(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
var unasked []conditions.Condition
|
|
for _, c := range open {
|
|
if wants(c, now) {
|
|
unasked = append(unasked, c)
|
|
}
|
|
}
|
|
return a.sayUnasked(ctx, unasked, "no router takes the controller's asks: no module holding "+
|
|
broker.AsksSeat+" that takes an ask under its asker's name is assigned")
|
|
}
|
|
a.saidNoRouter = false
|
|
}
|
|
if a.grantHeld != nil {
|
|
held, why, err := a.grantHeld(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !held {
|
|
if !a.saidNoGrant {
|
|
a.logf("the bus does not hold the controller's grant to ask yet: %s; the operator is asked nothing until it does", why)
|
|
a.saidNoGrant = true
|
|
}
|
|
open, err := a.open(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
var unasked []conditions.Condition
|
|
for _, c := range open {
|
|
if wants(c, now) {
|
|
unasked = append(unasked, c)
|
|
}
|
|
}
|
|
return a.sayUnasked(ctx, unasked, "the bus does not hold the controller's grant to ask yet: "+why)
|
|
}
|
|
a.saidNoGrant = false
|
|
}
|
|
channels := ""
|
|
if a.channels != nil {
|
|
channels = a.channels(ctx)
|
|
}
|
|
open, err := a.open(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
all, err := a.store.All(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
byCondition := map[string]asked{} // by partKey
|
|
// What is open and not a condition's — a rehearsal, a proposal — still counts toward what the router holds
|
|
// open for the controller (askMostOpen); expired unanswered, it is kept so, as the router says it too.
|
|
otherOpen := 0
|
|
for _, r := range all {
|
|
if r.State == askOpen && !r.ofACondition() && !now.Before(r.Ask.Expires) {
|
|
if _, err := a.store.Change(ctx, r.ID, func(x *asked) bool {
|
|
if x.State != askOpen || x.Acted != "" {
|
|
return false
|
|
}
|
|
x.State, x.Ended, x.Acted = string(asks.OutcomeExpired), now, "nothing: the ask expired unanswered"
|
|
return true
|
|
}); err != nil {
|
|
return err
|
|
}
|
|
continue
|
|
}
|
|
if r.State == askOpen && !r.ofACondition() {
|
|
otherOpen++
|
|
}
|
|
if r.State == askOpen && r.ofACondition() {
|
|
k := partKey(r.Condition, r.Part)
|
|
if prior, held := byCondition[k]; !held || r.Opened.After(prior.Opened) {
|
|
byCondition[k] = r
|
|
}
|
|
}
|
|
}
|
|
// A warrant missed while away, read from the router's record: for a condition's ask, and for a proposal's or a
|
|
// rehearsal's alike.
|
|
if a.routerRecord != nil {
|
|
var lookedUp []asked
|
|
for _, r := range byCondition {
|
|
lookedUp = append(lookedUp, r)
|
|
}
|
|
for _, r := range all {
|
|
if r.State == askOpen && !r.ofACondition() {
|
|
lookedUp = append(lookedUp, r)
|
|
}
|
|
}
|
|
for _, r := range lookedUp {
|
|
if now.Sub(r.Opened) < askCatchUpAfter {
|
|
continue
|
|
}
|
|
if w, err := a.routerRecord(ctx, r.ID); err == nil && w != nil {
|
|
body, _ := json.Marshal(w)
|
|
if err := a.Decided(ctx, body); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
if all, err = a.store.All(ctx); err != nil {
|
|
return err
|
|
}
|
|
byCondition = map[string]asked{}
|
|
for _, r := range all {
|
|
if r.State == askOpen && r.ofACondition() {
|
|
byCondition[partKey(r.Condition, r.Part)] = r
|
|
}
|
|
}
|
|
}
|
|
// 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 {
|
|
k := partKey(r.Condition, r.Part)
|
|
if r.State == string(asks.OutcomeChosen) && now.Sub(r.Ended) < askAgainAfterAnswer {
|
|
answered[k] = r
|
|
}
|
|
if r.State == string(asks.OutcomeRefused) {
|
|
if prior, has := refused[k]; !has || r.Opened.After(prior.Opened) {
|
|
refused[k] = r
|
|
}
|
|
}
|
|
}
|
|
// Refusals in a row, by part: those since the last ask that was not refused, within askRefusedCounted.
|
|
// superseded: a part asked since its last refusal, by an ask the router did not refuse — open, answered,
|
|
// expired or cancelled. Its refusal is history then, never said again (the review of PR 198: an answered
|
|
// retry brought the refusal back).
|
|
inRow, superseded := map[string]int{}, map[string]asked{}
|
|
for k, last := range refused {
|
|
for _, r := range all {
|
|
if partKey(r.Condition, r.Part) == k && r.State != string(asks.OutcomeRefused) && r.State != askUnsent &&
|
|
r.Opened.After(last.Opened) && r.Opened.After(superseded[k].Opened) {
|
|
superseded[k] = r
|
|
}
|
|
}
|
|
var since time.Time
|
|
for _, r := range all {
|
|
if partKey(r.Condition, r.Part) == k && r.State != string(asks.OutcomeRefused) && r.State != askUnsent &&
|
|
!r.Opened.After(last.Opened) && r.Opened.After(since) {
|
|
since = r.Opened
|
|
}
|
|
}
|
|
for _, r := range all {
|
|
if partKey(r.Condition, r.Part) == k && r.State == string(asks.OutcomeRefused) && r.Opened.After(since) &&
|
|
now.Sub(r.Opened) < askRefusedCounted {
|
|
inRow[k]++
|
|
}
|
|
}
|
|
}
|
|
wanted := map[string]bool{}
|
|
var unasked []conditions.Condition // refused by the router, and nothing it was refused for changed
|
|
var refusedWords []string
|
|
// 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 := otherOpen
|
|
for _, c := range open {
|
|
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
|
|
}
|
|
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) {
|
|
cur, held := byCondition[key]
|
|
ended := r.Ended
|
|
if ended.IsZero() {
|
|
ended = r.Opened
|
|
}
|
|
later, asked := superseded[key]
|
|
waiting := !held && !asked && r.Channels == channels && now.Sub(ended) < refusedRetryAfter(inRow[key])
|
|
// Asked again now (the wait over), or lately and the router's word not in yet: said as it was until
|
|
// that word, so the condition neither clears nor is raised again at each try.
|
|
retrying := !held && !asked && !waiting
|
|
verdictDue := held && asked && later.ID == cur.ID && now.Sub(cur.Opened) < askVerdictWait
|
|
if waiting || retrying || verdictDue {
|
|
if !saidUnasked {
|
|
unasked, saidUnasked = append(unasked, c), true
|
|
}
|
|
if r.Warrant != nil && r.Warrant.Words != "" {
|
|
refusedWords = append(refusedWords, r.Warrant.Words)
|
|
}
|
|
}
|
|
if waiting {
|
|
continue // refused lately, and nothing it was refused for has changed: asked again after the wait
|
|
}
|
|
}
|
|
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++
|
|
}
|
|
}
|
|
stillOpen := map[string]conditions.Condition{}
|
|
for _, c := range open {
|
|
stillOpen[c.Key] = c
|
|
}
|
|
for key, r := range byCondition {
|
|
if wanted[key] {
|
|
continue
|
|
}
|
|
// **A silence never takes an approval back** (the confirmation review of 2026-10-09, M1). Silence is an
|
|
// acknowledgement — anyone at the desk may give it — so a condition silenced while its approval is asked
|
|
// keeps that ask open, unchanged, until it is answered on a channel that proves who answered, or expires.
|
|
// It is not asked again once it ends, while the silence lasts.
|
|
if c, open := stillOpen[r.Condition]; open && c.SilencedAt(now) && r.Ask.Highest() != asks.Acknowledge &&
|
|
now.Before(r.Ask.Expires) && keepsItsAnswers(c, r) {
|
|
continue
|
|
}
|
|
if err := a.cancel(ctx, r, "the condition ended, was silenced or needs nothing now"); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
why := "the router refused the ask"
|
|
if len(refusedWords) > 0 {
|
|
why += ": " + refusedWords[0]
|
|
}
|
|
return a.sayUnasked(ctx, unasked, why)
|
|
}
|
|
|
|
// keepsItsAnswers says a condition still offers the answers an ask kept was asked with.
|
|
func keepsItsAnswers(c conditions.Condition, r asked) bool {
|
|
for _, p := range partsOf(c) {
|
|
if partKey(c.Key, p.name) == partKey(r.Condition, r.Part) {
|
|
return sameAsked(r.Actions, p.actions)
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// sourceAsker raises the asker's own condition.
|
|
const sourceAsker = "asker"
|
|
|
|
// sayUnasked keeps the asker's one condition: while a condition that needs the operator could not be asked
|
|
// on any channel, said loudly (failure must be loud), cleared when every one could be.
|
|
func (a *asker) sayUnasked(ctx context.Context, unasked []conditions.Condition, why string) error {
|
|
if a.raise == nil {
|
|
return nil
|
|
}
|
|
var obs []conditions.Observation
|
|
if len(unasked) > 0 {
|
|
keys := make([]string, 0, len(unasked))
|
|
severity := conditions.Warning
|
|
for _, c := range unasked {
|
|
keys = append(keys, c.Key)
|
|
if c.Severity == conditions.Urgent {
|
|
severity = conditions.Urgent
|
|
}
|
|
}
|
|
sort.Strings(keys)
|
|
obs = append(obs, conditions.Observation{Scope: conditions.ScopeSeat, ID: broker.AsksSeat, Token: "unasked",
|
|
Kind: "asks-undelivered", Severity: severity, Source: sourceAsker,
|
|
Summary: fmt.Sprintf("%d condition(s) that need the operator could not be asked on any channel: %s; %s",
|
|
len(keys), strings.Join(keys, ", "), why),
|
|
Headline: "Questions for you not delivered",
|
|
Explanation: "The mesh could not send you its questions on any channel.",
|
|
Needs: "answer them from the mesh MCP server, and check why no channel carries them.",
|
|
Resolved: "The mesh can ask you again"})
|
|
}
|
|
if err := a.raise(ctx, obs); err != nil {
|
|
a.logf("whether the operator could be asked could not be kept as a condition: %v", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// optionID is an action's label as an option's id: "Silence for a week" is silence-for-a-week.
|
|
func optionID(label string) string {
|
|
var b strings.Builder
|
|
dash := false
|
|
for _, r := range strings.ToLower(label) {
|
|
switch {
|
|
case r >= 'a' && r <= 'z', r >= '0' && r <= '9':
|
|
b.WriteRune(r)
|
|
dash = false
|
|
case !dash && b.Len() > 0:
|
|
b.WriteByte('-')
|
|
dash = true
|
|
}
|
|
}
|
|
return strings.TrimSuffix(b.String(), "-")
|
|
}
|
|
|
|
// doesWords is what an action does, in the words an option says it with.
|
|
func doesWords(act conditions.Action) string {
|
|
switch {
|
|
case act.Arguments["silence"] != "":
|
|
return "nothing more is said of it for a week"
|
|
case act.Verb == "mesh-delivery.release":
|
|
return "the delivery goes on"
|
|
case act.Verb == "mesh-delivery.stop":
|
|
return "the delivery ends"
|
|
case act.Verb == broker.ControllerSeat+".plans" && act.Arguments["go"] != "":
|
|
return "the delivery starts"
|
|
case act.Verb == broker.ControllerSeat+".plans" && act.Arguments["stop"] != "":
|
|
return "the delivery is stopped"
|
|
case strings.HasSuffix(act.Verb, ".restart"):
|
|
return "its service is restarted on " + act.Machine
|
|
}
|
|
return strings.ToLower(act.Label)
|
|
}
|
|
|
|
// 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 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: p.about,
|
|
Urgent: c.Severity == conditions.Urgent}
|
|
options := map[string]int{}
|
|
approves := false
|
|
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
|
|
// Every option binds the exact act it stands for (novox/hq ADR 0259 §6): the verb, the machine and
|
|
// every argument. The warrant then authorises that act and no other.
|
|
binds, _ := asks.ActDigest(boundAct(act))
|
|
q.Options = append(q.Options, asks.Option{ID: oid, Label: act.Label, Does: doesWords(act), Level: level,
|
|
Binds: binds})
|
|
}
|
|
q.Expires = now.Add(askAcknowledgeFor)
|
|
if approves {
|
|
q.Expires = now.Add(askApproveFor)
|
|
}
|
|
return q, options
|
|
}
|
|
|
|
// 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 {
|
|
var b [8]byte
|
|
_, _ = rand.Read(b[:])
|
|
return "c" + hex.EncodeToString(b[:])
|
|
}
|
|
|
|
// 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, p, now)
|
|
if err := q.Check(now); err != nil {
|
|
return err
|
|
}
|
|
body, err := json.Marshal(q)
|
|
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 nil
|
|
}
|
|
|
|
// 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 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)
|
|
return nil
|
|
}
|
|
|
|
// Decided takes the router's word on one of the controller's asks (link.Decider). An error is returned only
|
|
// when what was decided could not be kept, so the word is held and heard again.
|
|
func (a *asker) Decided(ctx context.Context, body []byte) error {
|
|
var w asks.Warrant
|
|
if err := json.Unmarshal(body, &w); err != nil {
|
|
a.logf("the router's word on an ask could not be read; ignored: %v", err)
|
|
return nil
|
|
}
|
|
if w.Asker != askerName {
|
|
a.logf("REFUSED a warrant for %s's ask %s: the controller acts only on its own", w.Asker, w.Ask)
|
|
return nil
|
|
}
|
|
r, err := a.store.Get(ctx, w.Ask)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if r == nil {
|
|
a.logf("REFUSED a warrant for the ask %s, which the controller does not hold", w.Ask)
|
|
return nil
|
|
}
|
|
if r.Acted != "" {
|
|
return nil // heard again: acted on once
|
|
}
|
|
now := a.now()
|
|
if w.Outcome != asks.OutcomeChosen {
|
|
acted := "nothing: the ask " + string(w.Outcome)
|
|
if 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 nil
|
|
}
|
|
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)
|
|
return nil
|
|
}
|
|
index, offered := r.Options[option.ID]
|
|
if !offered || index >= len(r.Actions) {
|
|
a.logf("REFUSED a warrant for the ask %s: it chose %s, which no action stands for", r.ID, option.ID)
|
|
return nil
|
|
}
|
|
act := r.Actions[index]
|
|
// The act about to be performed is the one the option bound when the controller asked: a record changed
|
|
// since is refused, never performed.
|
|
if err := option.Performs(boundAct(act)); err != nil {
|
|
a.logf("REFUSED a warrant for the ask %s: %v", r.ID, err)
|
|
return nil
|
|
}
|
|
open, err := a.open(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
stillOpen := r.Rehearsal || r.Proposal != nil // a rehearsal and a proposal are 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).
|
|
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 nil
|
|
}
|
|
// Claimed before acting, by compare-and-set: only the delivery whose write stands acts (security review
|
|
// 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
|
|
}
|
|
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 {
|
|
args[k] = v
|
|
}
|
|
if v, takes := args["why"]; takes && v == "" {
|
|
args["why"] = why
|
|
}
|
|
var acted error
|
|
outcome := "done"
|
|
switch {
|
|
case r.Rehearsal && act.Verb == rehearsalVerb:
|
|
// A rehearsal's answer performs nothing: it is recorded below as the operator's decision.
|
|
case r.Proposal != nil && act.Verb == proposalVerb:
|
|
// A proposed settings layer, set by this controller itself on Approve (novox/hq ADR 0277): the act's
|
|
// digest of the values is held to the record's own values before anything is set.
|
|
outcome, acted = a.decideProposal(ctx, *r, act, w)
|
|
if acted == nil && outcome != "" && !strings.HasPrefix(outcome, "nothing") {
|
|
outcome = "done: " + outcome
|
|
}
|
|
case act.Arguments["silence"] != "":
|
|
acted = a.silence(ctx, act.Arguments["silence"], conditions.MaxSilence, byWords(w), why)
|
|
default:
|
|
acted = a.call(ctx, act, args)
|
|
}
|
|
ended := a.now()
|
|
if acted != nil {
|
|
outcome = "failed: " + acted.Error()
|
|
}
|
|
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 != "" {
|
|
verbArgs = append(verbArgs, "on "+act.Machine)
|
|
}
|
|
keys := make([]string, 0, len(args))
|
|
for k := range args {
|
|
keys = append(keys, k)
|
|
}
|
|
sort.Strings(keys)
|
|
for _, k := range keys {
|
|
if k != "why" {
|
|
verbArgs = append(verbArgs, k+"="+args[k])
|
|
}
|
|
}
|
|
if err := a.record(ctx, link.HandAct{Verb: handActWarrant, Args: verbArgs, Why: why, By: byWords(w),
|
|
Cause: conditions.CauseOperatorAnswer, Condition: r.Condition, Via: viaWords(w), Ask: r.ID,
|
|
Proofs: w.Proofs, RequestedBy: r.Condition, Outcome: r.Acted}); err != nil {
|
|
a.logf("%s was done, and could NOT be recorded in the hand-act log: %v", why, err)
|
|
}
|
|
a.logf("%s: %s", why, r.Acted)
|
|
return nil
|
|
}
|
|
|
|
// handActWarrant is the verb an act the operator chose on a warrant is recorded under: a person's decision,
|
|
// never a repair (handActVerbs).
|
|
const handActWarrant = "warrant"
|
|
|
|
// byWords is who chose, as the hand-act log says it: "the operator, as telegram identity 42".
|
|
func byWords(w asks.Warrant) string {
|
|
if w.By == nil {
|
|
return "the operator"
|
|
}
|
|
return fmt.Sprintf("the %s, as %s identity %s", w.By.Who, w.By.Kind, w.By.Identity)
|
|
}
|
|
|
|
// viaWords is the channel an answer came through: its module and kind, and how the sender was known.
|
|
func viaWords(w asks.Warrant) string {
|
|
if w.By == nil {
|
|
return w.Channel
|
|
}
|
|
via := w.Channel + " (" + w.By.Kind + ")"
|
|
if w.By.Verified != "" {
|
|
via += ", " + w.By.Verified
|
|
}
|
|
return via
|
|
}
|
|
|
|
// errNotGranted is an action whose verb the controller's grant does not name.
|
|
var errNotGranted = errors.New("the controller's grant does not name this verb")
|