Area: mq — external NATS history stream (SSE replay / live events)
Two edge cases in internal/mq/external.go's two ways of reading the history stream back: Subscribe (live events, an ordered consumer) and ReplaySince (catch-up replay, a one-off consumer).
1. The live ordered consumer can skip rows across its own reconnect/reset.
Subscribe (internal/mq/external.go:707) opens a fresh jetstream.OrderedConsumer with DeliverPolicy: jetstream.DeliverNewPolicy each time it is called, and the function's own doc comment says this "skips nothing across reconnects." nats.go's ordered consumer transparently recreates itself under the hood on certain errors (a leader change, a consumer deleted out from under it), and when it does, it reissues the same DeliverNewPolicy — "deliver only messages published from now" — rather than resuming from the last row it actually delivered. If that internal reset happens before the consumer has delivered its first row (cursor still at its starting point), any rows published between the original Subscribe call and the reset are never delivered: they were published before "now" as the recreated consumer understands it.
- Evidence: inferred, from reading nats.go's ordered-consumer reconnect behavior; not reproduced against a live reset.
- Scenario: a client calls
Subscribe for live SSE events; before any row has arrived, the underlying ordered consumer hits a reset (e.g., the history stream's consumer leader changes); rows published in that narrow window between Subscribe and the reset are silently missing from the live stream, even though the doc comment promises otherwise.
2. ReplaySince can end a replay early, and ignores context cancellation between pulls.
ReplaySince (internal/mq/external.go:945) counts NumPending once at the start, then loops fetching in batches of up to replayBatch (256) with jetstream.FetchMaxWait(replayPullWait) (replayPullWait = 2 * time.Second, external.go:90). When a fetch returns zero rows (got == 0), the code treats that as "caught up" and returns success as long as the connection is still up:
if got == 0 {
if e.nc.IsConnected() {
return nil
}
return fmt.Errorf("replay fetch: %w", nats.ErrConnectionClosed)
}
A pull can legitimately return zero rows within that 2-second window for a reason other than "nothing left to deliver" — for example while the history stream's consumer leader is changing. That is indistinguishable here from a genuine end-of-replay, so the replay can truncate silently while remaining (the counted NumPending) is still greater than zero, and the caller sees a normal, successful return.
Separately, the Fetch call itself takes no context.Context — only FetchMaxWait. ctx.Err() is checked before each fetch and while iterating messages, but not while a fetch call is actually in flight, so a replay whose caller cancels its context can keep running for up to replayPullWait (~2s) past the cancellation before the next check notices.
- Evidence: inferred, from reading the loop and the nats.go
Fetch signature; not reproduced against an actual leader election.
Reproduce/confirm: for (1), force an ordered-consumer reset (e.g. kill the history stream's consumer leader) immediately after opening a live Subscribe, before any row has arrived, and check whether rows published in that window reach the subscriber. For (2), trigger a leader election on the history stream mid-replay and check whether ReplaySince returns early with remaining > 0; separately, cancel a replay's context mid-fetch and measure how long the call takes to return.
Found in review of #624.
Related: #624, #613.
Area: mq — external NATS history stream (SSE replay / live events)
Two edge cases in
internal/mq/external.go's two ways of reading the history stream back:Subscribe(live events, an ordered consumer) andReplaySince(catch-up replay, a one-off consumer).1. The live ordered consumer can skip rows across its own reconnect/reset.
Subscribe(internal/mq/external.go:707) opens a freshjetstream.OrderedConsumerwithDeliverPolicy: jetstream.DeliverNewPolicyeach time it is called, and the function's own doc comment says this "skips nothing across reconnects." nats.go's ordered consumer transparently recreates itself under the hood on certain errors (a leader change, a consumer deleted out from under it), and when it does, it reissues the sameDeliverNewPolicy— "deliver only messages published from now" — rather than resuming from the last row it actually delivered. If that internal reset happens before the consumer has delivered its first row (cursor still at its starting point), any rows published between the originalSubscribecall and the reset are never delivered: they were published before "now" as the recreated consumer understands it.Subscribefor live SSE events; before any row has arrived, the underlying ordered consumer hits a reset (e.g., the history stream's consumer leader changes); rows published in that narrow window betweenSubscribeand the reset are silently missing from the live stream, even though the doc comment promises otherwise.2.
ReplaySincecan end a replay early, and ignores context cancellation between pulls.ReplaySince(internal/mq/external.go:945) countsNumPendingonce at the start, then loops fetching in batches of up toreplayBatch(256) withjetstream.FetchMaxWait(replayPullWait)(replayPullWait = 2 * time.Second,external.go:90). When a fetch returns zero rows (got == 0), the code treats that as "caught up" and returns success as long as the connection is still up:A pull can legitimately return zero rows within that 2-second window for a reason other than "nothing left to deliver" — for example while the history stream's consumer leader is changing. That is indistinguishable here from a genuine end-of-replay, so the replay can truncate silently while
remaining(the countedNumPending) is still greater than zero, and the caller sees a normal, successful return.Separately, the
Fetchcall itself takes nocontext.Context— onlyFetchMaxWait.ctx.Err()is checked before each fetch and while iterating messages, but not while a fetch call is actually in flight, so a replay whose caller cancels its context can keep running for up toreplayPullWait(~2s) past the cancellation before the next check notices.Fetchsignature; not reproduced against an actual leader election.Reproduce/confirm: for (1), force an ordered-consumer reset (e.g. kill the history stream's consumer leader) immediately after opening a live
Subscribe, before any row has arrived, and check whether rows published in that window reach the subscriber. For (2), trigger a leader election on the history stream mid-replay and check whetherReplaySincereturns early withremaining > 0; separately, cancel a replay's context mid-fetch and measure how long the call takes to return.Found in review of #624.
Related: #624, #613.