Files
mesh-controller/internal/link/queue.go
T
jochen 9873b3bf13 Keep a replay from moving the mesh backwards, and tighten the queue's edges (review of hq ADR 0219)
A registered replay asked now outranked newer asks of its module, and a plan
took any later outcome as its answer. A replay is now refused while the module
is asked anywhere, a plan module asked under an id is answered by that id
alone, and a rebuild of a commit asks what the module follows. An ask handed
back after a restart no longer reads as dead; a cancel that meets a start is
withdrawn; pause holds for an ask fetched as it lands; kill removes containers
before and after the build ends and says whether its outcome went out; a
holder may say only its own machine is paused.
2026-10-05 19:45:45 +02:00

466 lines
16 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()
reply, err := conn.RequestWithContext(asking, NodeSeatToolSubject(seat, verb, node), body)
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
}