A seat's work is shared by its holders: node-build-agent, pulled one ask at a time (hq ADR 0190) #228
+16
-15
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
+2
-2
@@ -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: {
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user