diff --git a/internal/broker/derived.go b/internal/broker/derived.go index 2471c3e..466b88a 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -144,12 +144,18 @@ func ConsumerFor(p Principal) (Consumer, bool) { }, true } -// HolderConsumerFor is the worker a seat's holder gets on that seat's work queue. +// HolderConsumerFor is the worker a seat's holders share on that seat's work queue. // -// **A queue group even though the seat guarantees one holder.** The seat is *authority* — who may -// be the telegram sender — and the queue group is *delivery*. Tie delivery to the seat and the -// day somebody allows two holders for throughput, every message is processed twice with nothing -// reporting it. Kept separate, relaxing one changes nothing about the other. +// **One worker for every holder, and each holder pulls one ask when it is idle** (novox/hq ADR +// 0190). The seat is *authority* — who may be the telegram sender — and the worker is *delivery*, +// kept separate so that relaxing one changes nothing about the other: a node-scoped seat has a +// holder per machine, and all of them take from this one consumer, so the work is shared without +// any holder knowing about the others. Pulled rather than pushed because a push consumer hands the +// next ask to whichever subscriber the server picks, busy or not, and a pulled one is asked for by +// a holder that has just become free. Which is also what ends the race issue 186 describes — asks +// delivered behind the one being worked, expiring unacknowledged and dropped after the fifth +// redelivery: nothing is delivered that nobody asked for. A long build keeps its own ask alive +// (stillWorking); the ack wait is for a holder that died. func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) { if len(seat.Accepts) == 0 { return Consumer{}, false @@ -158,18 +164,13 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) Name: "SEAT_" + upperSnake(seat.Name) + "_worker", Stream: seatStreamName(seat.Name), Filters: []string{"mesh.seat." + seat.Name + ".accept.>"}, - Queue: "holders", AckWaitSeconds: 60, MaxDeliver: 5, - // **One in flight.** A holder works one ask at a time, so the server hands it one at a - // time: with the default of many, every ask behind the one being worked was delivered, - // left unacknowledged for the length of the work, redelivered after the ack wait, and - // after the fifth time dropped — on 2026-10-01 twenty-six of forty-three builds asked in - // two minutes were never built, and the queue read as empty (novox/hq issue 186). - MaxAckPending: 1, - Why: fmt.Sprintf("%s on %s holds %s; it acknowledges after the work is done, so a "+ - "crash mid-work redelivers rather than loses; one in flight, so a queue of asks is a "+ - "queue and not a race against the ack wait", module, node, seat.Name), + // As many in flight as there are holders working, which pulling bounds by itself: a holder + // fetches one and fetches again only after it acknowledged. The server's default stands. + Why: fmt.Sprintf("%s on %s holds %s; every holder pulls one ask at a time from this worker "+ + "and acknowledges after the work is done, so a crash mid-work redelivers rather than "+ + "loses and an idle holder is the one that takes the next ask", module, node, seat.Name), }, true } diff --git a/internal/broker/derived_test.go b/internal/broker/derived_test.go index 939d078..c7257d5 100644 --- a/internal/broker/derived_test.go +++ b/internal/broker/derived_test.go @@ -88,15 +88,20 @@ func TestAModuleThatConsumesNothingGetsNoConsumer(t *testing.T) { } } -// The seat is authority and the queue group is delivery. Tie them together and the day somebody -// allows two holders, every message is processed twice with nothing reporting it. -func TestAHoldersWorkerUsesAQueueGroupAnyway(t *testing.T) { +// The seat is authority and the worker is delivery (novox/hq ADR 0190): one worker per seat, shared +// by every holder and pulled from, so a second holder takes the next ask rather than a copy of the +// same one — which is what a queue group used to guard, and what pulling one durable gives outright. +func TestAHoldersWorkerIsOneSharedByItsHolders(t *testing.T) { c, ok := HolderConsumerFor("one", "telegram", telegramSeat()) if !ok { t.Fatal("the holder of a seat with inbound work got no worker") } - if c.Queue == "" { - t.Fatal("the worker is not in a queue group, so a second holder would double-process") + two, _ := HolderConsumerFor("two", "telegram", telegramSeat()) + if c.Name != two.Name || c.Stream != two.Stream { + t.Fatal("two holders got two workers, so each would process every ask") + } + if c.Push || c.Queue != "" { + t.Fatal("the worker is pushed, so the server would hand an ask to a busy holder") } if c.Stream != "SEAT_TELEGRAM_SENDER" { t.Fatalf("the worker reads %q, not the seat's own stream", c.Stream) @@ -154,15 +159,23 @@ func TestANodesDeclarationConsumerIsWhatItsOwnGrantAllows(t *testing.T) { has(t, perms.Subscribe, c.Filters[0]) } -// A holder works one ask at a time, so the server hands it one at a time (novox/hq issue 186): -// asks queued behind the one being worked wait in the stream rather than being delivered, -// left to expire and dropped after the fifth redelivery. -func TestAHoldersWorkerTakesOneAskAtATime(t *testing.T) { - c, found := HolderConsumerFor("anchor", "builder", DeclaredSeat{Name: "mesh-build-machine", Accepts: []string{"build"}}) +// Every holder of a seat shares one worker and pulls from it (novox/hq ADR 0190): no queue group +// and no delivery subject, because a push consumer hands the next ask to whichever subscriber the +// server picks, busy or not; and no cap of one in flight, because pulling bounds the asks in flight +// by the holders that are free — which is what ended the race of issue 186, where asks delivered +// behind the one being worked expired and were dropped. +func TestAHoldersWorkerIsPulledByEveryHolder(t *testing.T) { + c, found := HolderConsumerFor("anchor", "build-agent", DeclaredSeat{Name: "node-build-agent", Accepts: []string{"build"}}) if !found { t.Fatal("a seat that accepts work has no worker") } - if c.MaxAckPending != 1 { - t.Fatalf("the worker may have %d asks in flight; one, so a queue is a queue", c.MaxAckPending) + if c.Queue != "" || c.Push { + t.Fatalf("the worker is pushed (queue %q, push %v); a holder pulls when it is free", c.Queue, c.Push) + } + if c.MaxAckPending != 0 { + t.Fatalf("the worker caps asks in flight at %d; pulling bounds them by the holders working", c.MaxAckPending) + } + if c.Name != "SEAT_NODE_BUILD_AGENT_worker" || c.Stream != "SEAT_NODE_BUILD_AGENT" { + t.Fatalf("the worker is %s on %s; one per seat, shared by its holders", c.Name, c.Stream) } } diff --git a/internal/broker/nats.go b/internal/broker/nats.go index b0272bf..b52cb34 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -370,13 +370,17 @@ func PermissionsFor(p Principal) (Permissions, error) { // 3. Seats it holds: full participation. for _, s := range p.Holds { - // Taking work from the role's queue: the worker consumer it binds (asked about, - // delivered on, acknowledged), each on the seat's own stream. The first machine to - // take work over the new bus was refused the asking (2026-09-28). + // Taking work from the role's queue: the worker consumer every holder shares (asked + // about, pulled from, acknowledged), on the seat's own stream (novox/hq ADR 0190). A + // holder pulls — asks the consumer for its next message, answered on its own inbox — + // so what it needs is MSG.NEXT on that worker and nothing delivered to it. The first + // machine to take work over the new bus was refused the asking (2026-09-28). worker := "SEAT_" + upperSnake(s.Name) + "_worker" stream := seatStreamName(s.Name) - sub = append(sub, "_DELIVER."+worker, "_DELIVER."+worker+".>") - pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>") + pub = append(pub, + "$JS.API.CONSUMER.INFO."+stream+"."+worker, + "$JS.API.CONSUMER.MSG.NEXT."+stream+"."+worker, + "$JS.ACK."+stream+"."+worker+".>") for _, a := range s.Accepts { sub = append(sub, seatSubject(s, "accept", a)) } diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 79d73f3..3994266 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -37,8 +37,8 @@ accounts { subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] } } } { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { - publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } - subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } + publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.CONSUMER.MSG.NEXT.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } + subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: { diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index 85e32c3..0355857 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -126,29 +126,31 @@ func (m *natsMachine) Close() { } } -// Take binds to the role's worker and hands each request over, one at a time. +// Take binds to the role's worker and pulls one request at a time, handing each over. // // **Bound, never created.** The work queue and the worker on it are the controller's to define // (design 25 §3), and a build machine reaches no part of the JetStream API — so a missing one is said // as the mesh's to answer rather than quietly created with whatever this client defaults to. +// +// **Pulled, one at a time, by whichever holder is free** (novox/hq ADR 0190). Every machine holding +// the role binds this same worker; a machine asks for the next request only when it has finished +// the last, so a slow machine never holds an ask an idle one could take, and a machine that took +// five at once would run five container builds against one runtime and finish all of them slower +// than the first. func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) error { - worker, found := broker.HolderConsumerFor(m.on, "builder", + worker, found := broker.HolderConsumerFor(m.on, "build-agent", broker.DeclaredSeat{Name: m.seat, Accepts: []string{"build"}}) if !found { return fmt.Errorf("%s accepts no work, so there is nothing for this machine to take", m.seat) } - // One at a time, which the consumer's own ack-pending limit enforces rather than a prefetch - // setting: a machine that took five requests at once would run five container builds against one - // runtime and finish all of them slower than the first. - work := make(chan *nats.Msg, 1) // **The consumer's own filter, not the one subject this machine cares about.** The client checks // what is asked for against the consumer's filter and refuses anything that is not the same — // "subject does not match consumer" — so subscribing `…accept.build` against a consumer filtered // on `…accept.>` is rejected even though it is narrower. Learned twice now, on two different // consumers, which is why it is written down here. filter := worker.Filters[0] - sub, err := m.js.Context().ChanQueueSubscribe(filter, worker.Queue, work, + sub, err := m.js.Context().PullSubscribe(filter, worker.Name, nats.Bind(worker.Stream, worker.Name), nats.ManualAck()) if err != nil { return fmt.Errorf( @@ -159,13 +161,27 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) m.sub = sub for { - select { - case <-ctx.Done(): + if ctx.Err() != nil { return nil - case msg, ok := <-work: - if !ok { - return errors.New("the bus stopped delivering build work") + } + // 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)) + switch { + case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded): + return nil + case errors.Is(err, nats.ErrTimeout): + continue + case err != nil: + if sub.IsValid() { + // A transient fault in asking — a reconnect, a slow server — is asked past rather + // than ending the machine; one that outlasts the ack wait redelivers nothing lost. + time.Sleep(time.Second) + continue } + return fmt.Errorf("the bus stopped delivering build work: %w", err) + } + for _, msg := range fetched { var request BuildRequest if err := json.Unmarshal(msg.Data, &request); err != nil { // Unreadable: terminated rather than retried, because the next attempt reads the same