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.
This commit is contained in:
jochen
2026-10-05 19:45:45 +02:00
parent 541603c15c
commit 9873b3bf13
13 changed files with 600 additions and 106 deletions
+3 -1
View File
@@ -16,7 +16,9 @@ func TestTheBuildQueueIsControlledWithTheGrantsItNeedsAndNoMore(t *testing.T) {
}
has(t, holder.Publish, "$JS.API.DIRECT.GET.KV_SEAT_NODE_BUILD_AGENT_cancelled.$KV.SEAT_NODE_BUILD_AGENT_cancelled.>")
hasNot(t, holder.Publish, "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>")
has(t, holder.Publish, "mesh.seat.node-build-agent.event.paused.*")
has(t, holder.Publish, "mesh.seat.node-build-agent.event.paused.ace")
hasNot(t, holder.Publish, "mesh.seat.node-build-agent.event.paused.*")
has(t, holder.Publish, "mesh.seat.node-build-agent.event.log.*")
has(t, holder.Subscribe, "mesh.seat.node-build-agent.tool.kill.ace")
hasNot(t, holder.Subscribe, "mesh.seat.node-build-agent.tool.kill.g14")
+11
View File
@@ -124,6 +124,10 @@ type Principal struct {
// goes with the retired seat row.
var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"}
// perMachineEvents are a node-scoped seat's events about the holder itself, whose last token is the
// holder's machine (novox/hq ADR 0219): `paused.<node>`, the build agent saying whether it takes work.
var perMachineEvents = map[string]bool{"paused.*": true}
// enrolmentPrefix is the space every enrolling node's user and inbox live under, so the one place the
// controller may answer an enrolment is derived from the same constant the user is named from.
const enrolmentPrefix = "enrol"
@@ -421,6 +425,13 @@ func PermissionsFor(p Principal) (Permissions, error) {
sub = append(sub, seatSubject(s, "accept", a))
}
for _, e := range s.Emits {
// **A machine says its own state and no other's** (novox/hq ADR 0219): on a node-scoped
// seat, an event about the holder itself carries the machine as its last token, and
// each holder is granted its own machine's alone.
if s.Scope == "node" && p.Node != "" && perMachineEvents[e] {
pub = append(pub, seatSubject(s, "event", strings.TrimSuffix(e, "*")+p.Node))
continue
}
pub = append(pub, seatSubject(s, "event", e))
}
for _, t := range s.Serves {
+21 -3
View File
@@ -175,8 +175,12 @@ func MachineOverNATSWith(js *broker.JetStream, on, seat string, opts MachineOpti
return &natsMachine{js: js, on: on, seat: seat, opts: opts}
}
// pausedPoll is how often a paused holder looks again whether it was resumed.
const pausedPoll = 2 * time.Second
// pausedPoll is how often a paused holder looks again whether it was resumed, and pausedFetch how
// long one pull of a holder that can be paused waits.
const (
pausedPoll = 2 * time.Second
pausedFetch = 5 * time.Second
)
func (m *natsMachine) Close() {
if m.sub != nil {
@@ -235,7 +239,14 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build))
}
// One, and wait a while for it; an empty queue is a timeout, which is the normal state of a
// machine with nothing to build, and is asked again.
fetched, err := sub.Fetch(1, nats.Context(ctx))
// Asked for a few seconds at a time when the holder can be paused, so a pause reaches a pull
// already waiting within that, rather than when the client's own wait runs out.
asking, endAsking := ctx, func() {}
if m.opts.Paused != nil {
asking, endAsking = context.WithTimeout(ctx, pausedFetch)
}
fetched, err := sub.Fetch(1, nats.Context(asking))
endAsking()
switch {
case ctx.Err() != nil:
// Ours ended: the machine is being stopped.
@@ -257,6 +268,13 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build))
return fmt.Errorf("the bus stopped delivering build work: %w", err)
}
for _, msg := range fetched {
// **Paused while the pull was answered** (ADR 0219): "takes no new build" holds even for
// the one that arrived in that moment. Handed back at once for another holder — counted as
// a delivery, which a pause landing exactly then costs and nothing else does.
if m.opts.Paused != nil && m.opts.Paused() {
_ = msg.Nak()
continue
}
var request BuildRequest
if err := json.Unmarshal(msg.Data, &request); err != nil {
// Unreadable: terminated rather than retried, because the next attempt reads the same
+4 -3
View File
@@ -303,8 +303,9 @@ func TestNatsTwoMachinesShareTheWorkAndNeitherIsHandedMoreThanItCanTake(t *testi
case <-time.After(2 * time.Second):
}
// One finishes, and only then is the third taken — by that machine, the one that is free.
close(release["anchor"])
release["anchor"] = make(chan struct{})
// Released by a send, never by replacing the channel: the map is read by both machines' builds
// while this runs, and writing it raced them.
release["anchor"] <- struct{}{}
select {
case got := <-took:
if got.machine != "anchor" {
@@ -318,7 +319,7 @@ func TestNatsTwoMachinesShareTheWorkAndNeitherIsHandedMoreThanItCanTake(t *testi
// an explicit hand-back, because the real wait is a minute. Here: the laptop goes, anchor
// finishes, and with nothing queued nothing more is taken by the machine that is left.
machines["laptop"].Close()
close(release["anchor"])
release["anchor"] <- struct{}{}
select {
case got := <-took:
t.Fatalf("%s took %s; the queue should be empty", got.machine, got.id)
+54 -28
View File
@@ -62,8 +62,11 @@ type QueuedAsk struct {
Ref string `json:"ref,omitempty"`
AskedAt time.Time `json:"asked-at,omitempty"`
State string `json:"state"`
// On and Started are the holder that took an in-flight ask and when, from its started event.
// 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"`
}
@@ -109,19 +112,24 @@ type BuildHeard struct {
// 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 a machine is building now.** A holder builds one ask at a time (ADR 0190), so
// what a machine is building is its latest start, if no outcome has been heard for it. An ask
// behind the worker that is some machine's current build is in flight; at most AckPending of them
// are, newest start first — a machine that died mid-build has a latest start forever, and the
// worker's own count is what bounds that. An ask behind the worker that no machine has said it
// started is in flight only while the worker counts more pending than starts explain (taken, and
// not yet said); otherwise it is dead.
// **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.
// Cancel takes an ask in flight that no machine said it started, for exactly that reason.
func ClassifyQueue(asks []QueuedAsk, delivered uint64, ackPending int, heard BuildHeard) []QueuedAsk {
// Each machine's latest start.
latest := map[string]BuildStart{}
@@ -140,7 +148,7 @@ func ClassifyQueue(asks []QueuedAsk, delivered uint64, ackPending int, heard Bui
out := make([]QueuedAsk, len(asks))
copy(out, asks)
var running, unsaid []int
var running, handedBack, unsaid []int
for i := range out {
a := &out[i]
if a.Seq > delivered {
@@ -153,28 +161,29 @@ func ClassifyQueue(asks []QueuedAsk, delivered uint64, ackPending int, heard Bui
running = append(running, i)
continue
}
if _, ever := heard.Started[a.ID]; !ever {
unsaid = append(unsaid, i)
}
}
sort.SliceStable(running, func(x, y int) bool { return out[running[x]].Started.After(out[running[y]].Started) })
slots := ackPending
for _, i := range running {
if slots == 0 {
// A machine's latest start the worker no longer counts: it died with the ask, and the
// ask was handed out as often as it may be.
out[i].On, out[i].Started = "", time.Time{}
if s, ever := heard.Started[a.ID]; ever {
a.Was, a.Started = s.On, startedAt(s)
handedBack = append(handedBack, i)
continue
}
out[i].State = AskInFlight
slots--
unsaid = append(unsaid, i)
}
for _, i := range unsaid {
if slots == 0 {
break
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--
}
out[i].State = AskInFlight
slots--
}
return out
}
@@ -336,6 +345,23 @@ func MarkCancelled(ctx context.Context, js *broker.JetStream, seat, id string) e
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) {
+77 -8
View File
@@ -3,6 +3,7 @@ package link
import (
"context"
"encoding/json"
"sync"
"testing"
"time"
@@ -35,9 +36,9 @@ func TestTheQueueTellsWaitingInFlightAndDeadApart(t *testing.T) {
},
Outcomes: map[string]bool{"done": true},
}
got := ClassifyQueue(asks, 4, 3, heard)
got := ClassifyQueue(asks, 4, 2, heard)
want := map[string]string{"dead-1": AskDead, "running-g": AskInFlight, "running-a": AskInFlight,
"unsaid": AskInFlight, "waiting-1": AskWaiting, "waiting-2": AskWaiting}
"unsaid": AskDead, "waiting-1": AskWaiting, "waiting-2": AskWaiting}
for _, a := range got {
if a.State != want[a.ID] {
t.Errorf("%s is %s, want %s", a.ID, a.State, want[a.ID])
@@ -48,16 +49,45 @@ func TestTheQueueTellsWaitingInFlightAndDeadApart(t *testing.T) {
t.Errorf("an ask in flight does not say where: %+v", a)
}
// **The worker's count bounds it**: with two pending, the machine whose start is oldest died with
// its ask, and the ask nobody said is not in flight either.
got = ClassifyQueue(asks, 4, 2, heard)
// **The worker's count bounds it**: with one pending, the machine whose start is oldest died
// with its ask.
got = ClassifyQueue(asks, 4, 1, heard)
q = Queue{Asks: got}
if a, _ := q.Find("running-g"); a.State != AskInFlight {
t.Errorf("running-g is %s", a.State)
if a, _ := q.Find("running-a"); a.State != AskInFlight {
t.Errorf("running-a is %s", a.State)
}
if a, _ := q.Find("running-g"); a.State != AskDead || a.On != "" {
t.Errorf("past the worker's count running-g is %s on %q", a.State, a.On)
}
// **Handed back after a holder restarted mid-build**: started on g14, g14 then started something
// else, and the worker still counts it — it waits to be handed out again, and is not dead. With
// no more asks behind the worker than it counts pending, none is dead.
restarted := []QueuedAsk{{Seq: 1, ID: "dead-1"}, {Seq: 2, ID: "running-g"}}
for _, a := range ClassifyQueue(restarted, 2, 2, heard) {
if a.State != AskInFlight {
t.Errorf("%s is %s with two behind the worker and two pending", a.ID, a.State)
}
if a.ID == "dead-1" && (a.On != "" || a.Was != "g14") {
t.Errorf("an ask handed back reads on %q, was on %q", a.On, a.Was)
}
}
// The handed-back one takes a leftover slot before an ask nobody said started.
got = ClassifyQueue(asks, 4, 4, heard)
q = Queue{Asks: got}
for _, id := range []string{"dead-1", "running-g", "running-a", "unsaid"} {
if a, _ := q.Find(id); a.State != AskInFlight {
t.Errorf("with four pending %s is %s", id, a.State)
}
}
got = ClassifyQueue(asks, 4, 3, heard)
q = Queue{Asks: got}
if a, _ := q.Find("dead-1"); a.State != AskInFlight {
t.Errorf("the handed-back ask lost its slot to one nobody said: %s", a.State)
}
if a, _ := q.Find("unsaid"); a.State != AskDead {
t.Errorf("an ask no start explains, past the worker's count, is %s", a.State)
t.Errorf("unsaid is %s", a.State)
}
// None pending: everything behind the worker is dead, and names no machine.
for _, a := range ClassifyQueue(asks, 4, 0, heard) {
if a.Seq <= 4 && (a.State != AskDead || a.On != "") {
@@ -265,3 +295,42 @@ func TestNatsAWaitingAskerHearsItsAskCancelledOrKilled(t *testing.T) {
t.Fatalf("the waiter heard %+v (%v)", result, err)
}
}
// Paused in the moment a pull was answered: the ask is handed back, not built (ADR 0219).
func TestNatsAHolderPausedAsItsPullWasAnsweredHandsTheAskBack(t *testing.T) {
js := aBusWithACancelledSet(t)
ctx, stop := context.WithCancel(context.Background())
defer stop()
// Not paused when it asks; paused by the time the answer is read.
var looks int
var mu sync.Mutex
paused := func() bool {
mu.Lock()
defer mu.Unlock()
looks++
return looks > 1
}
took := make(chan string, 1)
machine := MachineOverNATSWith(js, "anchor", TheBuildMachine, MachineOptions{Paused: paused})
defer machine.Close()
go func() {
_ = machine.Take(ctx, func(ctx context.Context, work Build) { took <- work.Request().ID })
}()
time.Sleep(200 * time.Millisecond)
body, _ := json.Marshal(BuildRequest{ID: "build-handed-back", Repository: "/r"})
if _, err := js.Context().Publish(BuildWork(), body); err != nil {
t.Fatal(err)
}
select {
case id := <-took:
t.Fatalf("a holder paused as its pull was answered built %s", id)
case <-time.After(2 * time.Second):
}
q, err := ReadQueue(ctx, js, TheBuildMachine)
if err != nil {
t.Fatal(err)
}
if len(q.Asks) != 1 || q.Asks[0].ID != "build-handed-back" {
t.Fatalf("the ask is not back in the queue: %+v", q.Asks)
}
}