Work slower than the window says so, and one address is the bus's

Three faults the mesh's own logs showed this morning. A handler that outlives the acknowledgement
window was handed its message again while it was still working: acting on a merge builds modules,
minutes against a thirty-second window, so one merge ran the whole catalogue five times over. The
transport now says the work is in progress while it runs, which is where the window belongs.

Everything the mesh hands out — a token, a membership, a person's credential — took its address
from the enrolment setting, which on a mesh that has moved still names the broker it moved from:
the first person issued after the move was handed the retired broker's port. There is one bus, and
its address is the one the control plane is connected to.

And `operator issue` documented an argument order its parser refused.
This commit is contained in:
2026-09-28 09:51:50 +02:00
parent 208388978a
commit 2f3bfda8c0
5 changed files with 123 additions and 3 deletions
+37
View File
@@ -154,10 +154,47 @@ func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Con
m.seq = meta.Sequence.Stream
m.delivered = meta.NumDelivered
}
// **Work that outlives the acknowledgement window says so while it runs.**
//
// The bus waits a fixed time to be told a message was taken, and then hands it to whoever
// consumes next — which is right for a consumer that died and wrong for one that is busy.
// Acting on a merge builds every module the merge changed: minutes of work against a
// thirty-second window. So the same merge was handed over again while the first build was
// still running, and again after that — on 2026-09-28 one merge ran the mesh's whole
// catalogue five times over and exhausted a public registry's pull limit.
//
// Here rather than in each handler, because the window belongs to the transport and every
// handler would otherwise have to remember it. It changes nothing about a handler that
// dies: a message is kept alive only while this goroutine is, so a controller that stops
// stops saying so, and the bus redelivers exactly as it should.
working := make(chan struct{})
defer close(working)
go stillWorking(msg, working)
}
act(ctx, m)
}
// heartbeatWhileWorking is how often a handler still running tells the bus so — comfortably inside
// the shortest acknowledgement window the mesh gives any of its consumers.
const heartbeatWhileWorking = 10 * time.Second
// stillWorking keeps one message alive until the work on it returns.
//
// An error is not worth reporting: what the bus does when it is not told is redeliver, which is
// exactly what happens if this fails, and the handler's own outcome is the thing worth logging.
func stillWorking(msg *nats.Msg, done <-chan struct{}) {
tick := time.NewTicker(heartbeatWhileWorking)
defer tick.Stop()
for {
select {
case <-done:
return
case <-tick.C:
_ = msg.InProgress()
}
}
}
// kindOfSubject is how this transport's addressing becomes what the mesh calls a message.
//
// By subject, which is the only thing the server enforces: a body claiming to be a report does not