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 }