Files
mesh-controller/internal/link/queue.go
T
jochen 8ddc019cd2 Grant the self-check its ban-list question, say a refusal at once, judge the engine by its delivered version (hq to-be 45 Phase 1)
Live on 2026-10-06, two of the first self-check's findings were its own:

- D8 asked every machine's node-intrusion-prevention.banned, and the
  controller's grant did not name the subject: the bus refused it 24 times
  and D8 timed out after thirty seconds instead of saying so. The verbs the
  self-check asks are named in broker.VerbsTheSelfCheckAsks and granted
  (mesh.seat.<seat>.tool.<verb>.*); each probe declares the seat verbs it
  calls, askSeatTool refuses an undeclared one, and a test over the
  registry fails a probe whose question the controller is not granted.
  AskSeatTool now returns a refused publish at once ("the bus refused…")
  instead of waiting out its timeout; D8 asks the machines in parallel.
- D10 read every machine as behind right after a push: a node-engine says
  its version as the directory it is delivered into, the archive's digest
  (31045596c83a, catalogue versionOf), and D10 compared that with the
  build's commit (1545b00a). It now compares with the versions the
  registered build is delivered as, and a hand-placed engine's commit.
2026-10-06 10:44:38 +02:00

489 lines
17 KiB
Go

package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"regexp"
"sort"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
)
// The build queue, read and changed by hand (novox/hq ADR 0219).
//
// **A queue nobody can see is a queue nobody can trust.** Until this, the build seat's work queue
// was visible only as its consequences: a plan LATE with nothing to say why, a build that never
// started because five asks ahead of it were waiting on a laptop that was shut. ADR 0219 makes the
// queue the controller's to show and change, and the build running on one machine that machine's
// holder's to end.
//
// Three kinds of ask are in the seat's stream at any moment, told apart by the worker consumer
// every holder pulls from:
//
// - **waiting** — after the last one the worker handed out (stream seq > delivered.stream_seq);
// - **in flight** — handed out and not yet settled, which the holder that took it said with its
// `started` event, and which the worker counts among its acks pending;
// - **dead** — handed out as many times as the worker allows and never settled. A work queue
// drops what it settles, so one still in the stream behind the worker is either in flight or
// dead; the started events and the worker's pending count say which.
//
// Everything that drops an ask leaves a failed outcome for it, recorded the way a failed build's is
// (`cancelled by hand`, `killed by hand`), so a plan waiting on it fails visibly rather than waiting
// for ever on an ask nobody will answer.
// The words an outcome says when a person ended the ask.
const (
CancelledByHand = "cancelled by hand"
KilledByHand = "killed by hand"
)
// Ask states in a queue.
const (
AskWaiting = "waiting"
AskInFlight = "in flight"
AskDead = "dead"
)
// QueuedAsk is one ask in the seat's work queue.
type QueuedAsk struct {
Seq uint64 `json:"seq"`
Request BuildRequest `json:"-"`
ID string `json:"id"`
// Repository, Path and Ref are the ask's, as the request carried them; Held never travels out of
// here — it is every artifact the mesh has built, and a queue listing is not where that belongs.
Repository string `json:"repository"`
Path string `json:"path,omitempty"`
Ref string `json:"ref,omitempty"`
AskedAt time.Time `json:"asked-at,omitempty"`
State string `json:"state"`
// On and Started are the holder building an in-flight ask and when it started, from its started
// event. Was is the holder that started an in-flight ask and is no longer building it — handed
// back, to be handed out again — with Started its start there.
On string `json:"on,omitempty"`
Was string `json:"was-on,omitempty"`
Started time.Time `json:"started,omitempty"`
}
// Queue is the seat's work queue as it stands.
type Queue struct {
Seat string `json:"seat"`
// Delivered is the last stream sequence the worker handed out; AckPending how many it handed out
// and has not had settled; MaxDeliver how many times it hands one ask out.
Delivered uint64 `json:"delivered"`
AckPending int `json:"ack-pending"`
MaxDeliver int `json:"max-deliver"`
Asks []QueuedAsk `json:"asks"`
}
// Of is the asks in one state, in queue order.
func (q Queue) Of(state string) []QueuedAsk {
var out []QueuedAsk
for _, a := range q.Asks {
if a.State == state {
out = append(out, a)
}
}
return out
}
// Find is the ask with this id.
func (q Queue) Find(id string) (QueuedAsk, bool) {
for _, a := range q.Asks {
if a.ID == id {
return a, true
}
}
return QueuedAsk{}, false
}
// BuildHeard is what the seat's events say about builds: the latest start of each id, and every id
// whose outcome was heard.
type BuildHeard struct {
Started map[string]BuildStart
Outcomes map[string]bool
}
// ClassifyQueue says which state each ask is in. Pure: the stream's asks in sequence order, what the
// worker says, and what the seat's events said.
//
// **In flight is what the worker counts handed out and unsettled**, at most AckPending of the asks
// behind it, given out in this order:
//
// 1. what a machine is building now — its latest start, no outcome heard (a holder builds one ask
// at a time, ADR 0190), newest start first;
// 2. what a machine started and is no longer building — a holder restarted mid-build, the ask
// handed back and waiting to be handed out again while it has deliveries left — newest start
// first, naming the machine it was last on;
// 3. what was taken and not yet said started.
//
// Only what is left after that is dead: so nothing behind the worker is dead while there are no more
// of them than the worker counts pending. A machine that died mid-build keeps a latest start for
// ever; the worker's count is what says the ask was handed out as often as it may be.
//
// **The worker's count is the server's, and it lags in one case**: an ask handed out its last time
// and handed back is still counted pending until some holder pulls again, when the server finds it
// past its deliveries and lets it go. Holders pull whenever they are idle, so this is a moment —
// except while every holder is paused, when such an ask reads as in flight with no machine named.
func ClassifyQueue(asks []QueuedAsk, delivered uint64, ackPending int, heard BuildHeard) []QueuedAsk {
// Each machine's latest start.
latest := map[string]BuildStart{}
for _, s := range heard.Started {
at := startedAt(s)
if was, ok := latest[s.On]; !ok || at.After(startedAt(was)) {
latest[s.On] = s
}
}
current := map[string]BuildStart{}
for _, s := range latest {
if !heard.Outcomes[s.ID] {
current[s.ID] = s
}
}
out := make([]QueuedAsk, len(asks))
copy(out, asks)
var running, handedBack, unsaid []int
for i := range out {
a := &out[i]
if a.Seq > delivered {
a.State = AskWaiting
continue
}
a.State = AskDead
if s, ok := current[a.ID]; ok {
a.On, a.Started = s.On, startedAt(s)
running = append(running, i)
continue
}
if s, ever := heard.Started[a.ID]; ever {
a.Was, a.Started = s.On, startedAt(s)
handedBack = append(handedBack, i)
continue
}
unsaid = append(unsaid, i)
}
newestFirst := func(ix []int) {
sort.SliceStable(ix, func(x, y int) bool { return out[ix[x]].Started.After(out[ix[y]].Started) })
}
newestFirst(running)
newestFirst(handedBack)
slots := ackPending
for _, group := range [][]int{running, handedBack, unsaid} {
for _, i := range group {
if slots == 0 {
// Past what the worker counts: handed out as often as it may be.
out[i].On, out[i].Started = "", time.Time{}
continue
}
out[i].State = AskInFlight
slots--
}
}
return out
}
func startedAt(s BuildStart) time.Time {
t, _ := time.Parse(time.RFC3339Nano, s.At)
return t
}
// ReadQueue reads a build seat's work queue: the stream's asks, the worker's position, and what the
// seat's events said about the builds behind it. Through the JetStream API alone (STREAM.INFO,
// CONSUMER.INFO, STREAM.MSG.GET), which only the controller's account reaches.
func ReadQueue(ctx context.Context, js *broker.JetStream, seat string) (Queue, error) {
worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: seat, Accepts: []string{"build"}})
q := Queue{Seat: seat}
api, err := jetstream.New(js.Conn())
if err != nil {
return q, err
}
stream, err := api.Stream(ctx, worker.Stream)
if err != nil {
return q, fmt.Errorf("the %s work queue (%s) cannot be read: %w", seat, worker.Stream, err)
}
consumer, err := stream.Consumer(ctx, worker.Name)
switch {
case errors.Is(err, jetstream.ErrConsumerNotFound):
// No holder has ever been given a worker: everything in the stream is waiting.
case err != nil:
return q, fmt.Errorf("the worker of %s cannot be read: %w", seat, err)
default:
info, err := consumer.Info(ctx)
if err != nil {
return q, fmt.Errorf("the worker of %s cannot be read: %w", seat, err)
}
q.Delivered = info.Delivered.Stream
q.AckPending = info.NumAckPending
q.MaxDeliver = info.Config.MaxDeliver
}
// Every ask, by next-by-subject from the first: the stream holds only what is unsettled, so this
// is a handful of reads, never the history.
var asks []QueuedAsk
var seq uint64 = 1
for n := 0; n < 10000; n++ {
msg, err := stream.GetMsg(ctx, seq, jetstream.WithGetMsgSubject(BuildWorkOf(seat)))
if errors.Is(err, jetstream.ErrMsgNotFound) {
break
}
if err != nil {
return q, fmt.Errorf("reading the %s work queue: %w", seat, err)
}
ask := QueuedAsk{Seq: msg.Sequence}
if json.Unmarshal(msg.Data, &ask.Request) == nil {
r := ask.Request
ask.ID, ask.Repository, ask.Path, ask.Ref = r.ID, r.Repository, r.Path, r.Ref
if at, ok := BuildAskedAt(r.ID); ok {
ask.AskedAt = at
}
}
asks = append(asks, ask)
seq = msg.Sequence + 1
}
// What the events say, from shortly before the oldest ask behind the worker: its start, if any,
// came after it was asked.
var since time.Time
for _, a := range asks {
if a.Seq <= q.Delivered && (since.IsZero() || (!a.AskedAt.IsZero() && a.AskedAt.Before(since))) {
since = a.AskedAt
}
}
heard := BuildHeard{Started: map[string]BuildStart{}, Outcomes: map[string]bool{}}
if len(asks) > 0 && q.Delivered > 0 {
if since.IsZero() {
since = time.Now().Add(-7 * 24 * time.Hour)
}
if heard, err = ReadBuildEvents(ctx, js, seat, since.Add(-time.Minute)); err != nil {
return q, err
}
}
q.Asks = ClassifyQueue(asks, q.Delivered, q.AckPending, heard)
return q, nil
}
// ReadBuildEvents is every start and outcome the seat announced since a moment, from the events
// stream, with a consumer of its own that is gone when this returns.
func ReadBuildEvents(ctx context.Context, js *broker.JetStream, seat string, since time.Time) (BuildHeard, error) {
heard := BuildHeard{Started: map[string]BuildStart{}, Outcomes: map[string]bool{}}
// One token after `event`: started and built, never a build's log lines, which are two.
subject := "mesh.seat." + seat + ".event.*"
sub, err := js.Context().PullSubscribe(subject, "",
nats.BindStream(broker.EventsStream), nats.StartTime(since), nats.AckNone())
if err != nil {
return heard, fmt.Errorf("cannot read %s from the bus: %w", subject, err)
}
defer func() { _ = sub.Unsubscribe() }()
for {
batch, err := sub.Fetch(500, nats.MaxWait(2*time.Second))
if err != nil && !errors.Is(err, nats.ErrTimeout) && !errors.Is(err, context.DeadlineExceeded) {
return heard, fmt.Errorf("reading %s: %w", subject, err)
}
for _, msg := range batch {
switch msg.Subject {
case BuildStartedOf(seat):
var s BuildStart
if json.Unmarshal(msg.Data, &s) == nil && s.ID != "" {
if was, ok := heard.Started[s.ID]; !ok || startedAt(s).After(startedAt(was)) {
heard.Started[s.ID] = s
}
}
case BuildOutcomeOf(seat):
var r BuildResult
if json.Unmarshal(msg.Data, &r) == nil && r.ID != "" {
heard.Outcomes[r.ID] = true
}
}
}
if len(batch) < 500 {
break
}
if ctx.Err() != nil {
return heard, ctx.Err()
}
}
return heard, nil
}
// --- the cancelled set ----------------------------------------------------------------------
// safeKey is an id that can be a key, and a subject token, as it is: a build id, never a pattern.
var safeKey = regexp.MustCompile(`^[A-Za-z0-9_-]+$`)
// Cancelled is what the controller writes for an ask it cancelled.
type Cancelled struct {
At string `json:"at"`
Why string `json:"why"`
}
// MarkCancelled puts an ask's id into the seat's cancelled set (novox/hq ADR 0219). The controller's
// act, before it deletes the ask from the queue.
func MarkCancelled(ctx context.Context, js *broker.JetStream, seat, id string) error {
if !safeKey.MatchString(id) {
return fmt.Errorf("%q is not a build id", id)
}
api, err := jetstream.New(js.Conn())
if err != nil {
return err
}
kv, err := api.KeyValue(ctx, broker.CancelledSetName(seat))
if err != nil {
return fmt.Errorf("the cancelled set of %s is not on the bus — the controller asserts it at its "+
"start, so one older than this has not: %w", seat, err)
}
body, err := json.Marshal(Cancelled{At: time.Now().UTC().Format(time.RFC3339Nano), Why: CancelledByHand})
if err != nil {
return err
}
_, err = kv.Put(ctx, id, body)
return err
}
// UnmarkCancelled takes an ask's id out of the cancelled set: a cancel withdrawn, because the ask
// was taken as it was being cancelled.
func UnmarkCancelled(ctx context.Context, js *broker.JetStream, seat, id string) error {
if !safeKey.MatchString(id) {
return nil
}
api, err := jetstream.New(js.Conn())
if err != nil {
return err
}
kv, err := api.KeyValue(ctx, broker.CancelledSetName(seat))
if err != nil {
return err
}
return kv.Delete(ctx, id)
}
// IsCancelled asks the seat's cancelled set whether an ask was cancelled: one direct read of one key,
// which is all a holder is granted of it.
func IsCancelled(conn *nats.Conn, seat, id string) (bool, error) {
if !safeKey.MatchString(id) {
// Not a key the controller could have written.
return false, nil
}
set := broker.CancelledSetName(seat)
reply, err := conn.Request("$JS.API.DIRECT.GET.KV_"+set+".$KV."+set+"."+id, nil, 3*time.Second)
if err != nil {
return false, err
}
return cancelledFrom(reply)
}
// cancelledFrom reads a direct get's answer: a value is a cancel; not found, or a deletion marker,
// is not.
func cancelledFrom(reply *nats.Msg) (bool, error) {
if reply.Header != nil {
switch reply.Header.Get("Status") {
case "":
case "404":
return false, nil
default:
return false, fmt.Errorf("the cancelled set answered %s %s",
reply.Header.Get("Status"), reply.Header.Get("Description"))
}
if op := reply.Header.Get("KV-Operation"); op == "DEL" || op == "PURGE" {
return false, nil
}
}
return len(reply.Data) > 0, nil
}
// --- a holder's verbs -----------------------------------------------------------------------
// NodeSeatToolSubject is where a node-scoped seat's verb is asked of one machine's holder (design 33
// §4): the seat's verb subject with the machine as its last token.
func NodeSeatToolSubject(seat, verb, node string) string {
return SeatToolSubject(seat, verb) + "." + node
}
// AskSeatTool asks one machine's holder of a node seat one of its verbs and reads its answer.
func AskSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string, args any,
timeout time.Duration) (Answer, error) {
body, err := json.Marshal(args)
if err != nil {
return Answer{}, err
}
asking, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
subject := NodeSeatToolSubject(seat, verb, node)
// **A refused question is known at once** (novox/hq to-be 45, D8 live on 2026-10-06): the bus
// says it to the connection and the request would otherwise wait out its whole timeout, reading
// as a machine that did not answer.
refused, stop := refusalsOf(conn, subject)
defer stop()
type replied struct {
msg *nats.Msg
err error
}
done := make(chan replied, 1)
go func() {
msg, err := conn.RequestWithContext(asking, subject, body)
done <- replied{msg, err}
}()
var reply *nats.Msg
select {
case r := <-done:
reply, err = r.msg, r.err
case why := <-refused:
cancel()
return Answer{}, fmt.Errorf("the bus refused the controller asking %s.%s of %s — its grants do not "+
"name %s: %v", seat, verb, node, subject, why)
}
switch {
case errors.Is(err, nats.ErrNoResponders):
return Answer{}, fmt.Errorf("nothing on %s answers %s.%s: its holder is not running, or is "+
"older than the verb", node, seat, verb)
case errors.Is(err, context.DeadlineExceeded), errors.Is(err, nats.ErrTimeout):
return Answer{}, fmt.Errorf("%s did not answer %s.%s within %s", node, seat, verb, timeout)
case err != nil:
return Answer{}, err
}
var answer Answer
if err := json.Unmarshal(reply.Data, &answer); err != nil {
return Answer{}, fmt.Errorf("%s answered %s.%s with something unreadable: %w", node, seat, verb, err)
}
return answer, nil
}
// --- a holder saying whether it takes work ----------------------------------------------------
// HolderState is what a holder says about itself whenever it is paused or resumed, at its start,
// and while it stays paused: whether it takes work (novox/hq ADR 0219). Retained on the events
// stream under the machine's own subject, so the last one is the machine's state.
type HolderState struct {
On string `json:"on"`
Paused bool `json:"paused"`
At string `json:"at"`
}
// BuildPausedOf is where one machine's holder of a build seat says whether it takes work.
func BuildPausedOf(seat, node string) string { return "mesh.seat." + seat + ".event.paused." + node }
// PausedSaid reads what each machine last said about taking work. A machine that never said is not
// in the answer — a holder older than ADR 0219 takes work and never pauses.
func PausedSaid(js *broker.JetStream, seat string, nodes []string) (map[string]HolderState, error) {
out := map[string]HolderState{}
for _, node := range nodes {
msg, err := js.Context().GetLastMsg(broker.EventsStream, BuildPausedOf(seat, node))
if errors.Is(err, nats.ErrMsgNotFound) {
continue
}
if err != nil {
return out, err
}
var s HolderState
if json.Unmarshal(msg.Data, &s) == nil {
out[node] = s
}
}
return out, nil
}