From 11b20b10ffb76b17ab31dc349e190434b45355cf Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 1 Oct 2026 16:16:13 +0200 Subject: [PATCH] The build machine takes one ask at a time, and says so while it builds With the worker consumer's default of many deliveries in flight, every ask behind the one being built was delivered at once, left unacknowledged for the length of the build, redelivered after the ack wait and dropped after the fifth time: on 2026-10-01 twenty-six of forty-three builds asked in two minutes were never built and the queue read as empty (hq issue 186). The holder's worker now has one in flight, and a running build tells the bus it is still working, as the controller's long handlers do, so a build longer than the ack wait is neither redelivered nor counted out. --- internal/broker/derived.go | 9 ++++++++- internal/broker/derived_test.go | 13 +++++++++++++ internal/link/builds_nats.go | 6 ++++++ 3 files changed, 27 insertions(+), 1 deletion(-) diff --git a/internal/broker/derived.go b/internal/broker/derived.go index b6bbb4d..2471c3e 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -161,8 +161,15 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) 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", module, node, seat.Name), + "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), }, true } diff --git a/internal/broker/derived_test.go b/internal/broker/derived_test.go index b4529d4..939d078 100644 --- a/internal/broker/derived_test.go +++ b/internal/broker/derived_test.go @@ -153,3 +153,16 @@ func TestANodesDeclarationConsumerIsWhatItsOwnGrantAllows(t *testing.T) { has(t, perms.Publish, "$JS.ACK.NODES."+c.Name+".>") 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"}}) + 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) + } +} diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index 786f783..85e32c3 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -173,7 +173,13 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) _ = msg.Term() continue } + // A build outlives the acknowledgement window many times over; said while it runs, + // as the controller says it for its own long handlers, so the server neither hands + // the ask to a second machine nor counts the wait against its deliveries. + working := make(chan struct{}) + go stillWorking(msg, working) do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js}) + close(working) } } }