Repository navigation
fix: TimedOutboxSweeper starvation and publish-confirm floods (#4560) - #4562
Conversation
TimedOutboxSweeper records tick lag, sweep duration and sweeps by outcome directly through a new IAmABrighterSweeperMeter role, so a stalled or late sweeper is visible when traces are sampled out. Refs #4560 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
…irm floods Pins the defects confirmed for #4560. These tests are red until the fix lands: - a sweep still runs while every thread-pool worker is blocked - StopAsync waits for an in-flight sweep - a throwing sweep or CreateScope does not crash the process, and the lock is released - TimerInterval <= 0 is rejected - Kafka and RMQ.Sync raise publish confirmations one at a time, in order - the sweeper records tick lag, sweep duration and sweeps by outcome (ADR 0081) Refs #4560 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
…tboxSweeper Structural only, per ADR 0081: IAmABrighterSweeperMeter, SweepOutcome, NullSweeperMeter and the instrument names. TimedOutboxSweeper gains optional timeProvider and meter constructor parameters, not yet used. Refs #4560 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
The sweeper ran each sweep as an async void callback on a thread-pool Timer, so it stopped sweeping while every pool worker was blocked. Under load that left unconfirmed messages in the InMemoryOutbox. - Sweeps run on a dedicated thread, scheduled by an optional TimeProvider - StopAsync waits for a sweep in flight - A failed sweep is logged and retried on the next tick instead of crashing the process; the lock is released on every path, including CreateScope - TimerInterval below one second is rejected with a ConfigurationException - SweeperMeter records tick lag, sweep duration and sweeps by outcome (ADR 0081), registered by AddBrighterInstrumentation Fixes #4560 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
… order Both producers queued one thread-pool work item per delivery report or ack, so under load confirmations ran concurrently (50 of 50 at once on Kafka) and the queue grew with the post rate, delaying MarkDispatched and causing the sweeper to re-send delivered messages. A new SerialCallbackQueue in Paramore.Brighter.Tasks runs the callbacks on a single worker, in delivery order, and keeps the in-flight tracking that Dispose waits on. Refs #4560 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
SerialCallbackQueue started its drain worker with Task.Run, so confirmations still waited for a thread-pool thread and stalled when the pool was starved. It now drains on one background thread per queue, running each callback in a BrighterAsyncContext so its continuations stay on that thread. A confirmation that arrives after Complete is dropped instead of throwing on the broker's thread. Refs #4560 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
) Covers the dedicated sweeper and confirmation threads, the new sweeper metrics, the TimerInterval and subscriber-ordering behaviour changes, and the cost of one thread per confirming producer on small pods. Refs #4560 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Triage, confirm evidence, regression tests, fix and verify for #4560. diagnostic/ keeps the starve and pileup experiments as an end-to-end tool; the samples/ harness was dropped, since the starvation and confirm-ordering regression tests already prove the fix. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Conflicts: - KafkaMessageProducer.Dispose: keep master's idempotent KafkaResourceCleanup (#4556) and complete the SerialCallbackQueue as one more guarded step. - docs/adr/index.md: regenerated; master added a second 0081 ADR, which the index allows (identity is the filename stem). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
|
@claude Let's review the changes here with a view to trying to remove issues with respect to thread starvation |
|
Claude encountered an error after 0s —— View job I'll analyze this and get back to you. |
1 similar comment
|
Claude encountered an error after 0s —— View job I'll analyze this and get back to you. |
|
@claude Let's review the changes here with a view to trying to remove issues with respect to thread starvation |
|
Claude finished @iancooper's task in 3m 7s —— View job Review: what can still starve the sweeper (PR #4562)
Summary: The PR fixes the schedule. It does not yet fix the sweep itself. Ticks now fire on time, but for Kafka each message a sweep sends still needs a free pool thread before the sweep can move on. That case can be fixed here cheaply, and the current starvation test cannot detect it. I did not build the code or run the tests; everything below comes from reading it. 1. 🔴 The sweep runs off the pool, but each Kafka send still waits for a pool thread
return SweepAsync().GetAwaiter().GetResult();The sweeper thread has no
As far as I know, Confluent completes The PR lists this as a known limit, but for Kafka with the in-memory outbox (the case in the issue) it can be fixed:
2. 🔴 The starvation test can't detect #1
Add a test double whose 3. 🟠 One confirmation queue per producer means one slow callback stalls all confirmations
Suggestions:
4. 🟡 Both new threads capture the starting caller's
|
Review of #4562 (findings #1/#2): the existing starvation test uses InMemoryMessageProducer, which completes inline, so it only proves the tick fires. OffPoolCompletingProducer completes each send from a dedicated delivery thread through a RunContinuationsAsynchronously TCS and awaits it plainly, as KafkaMessagePublisher awaits Confluent's ProduceAsync (Confluent.Kafka 2.15.1 uses RunContinuationsAsynchronously). RED (3/3): "dispatched 1 of 5 outstanding messages within 00:00:03". Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Root cause: SweepOnce blocked on SweepAsync with no context, so the continuation after each awaited Kafka send needed a pool thread; under starvation a sweep dispatched one message and stalled. Running the sweep in a BrighterAsyncContext was not enough on its own: a ConfigureAwait(false) continuation is not inlined on a thread with a SynchronizationContext, so two hops on the sweep path also sent it back to the pool. - TimedOutboxSweeper runs the sweep with BrighterAsyncContext.Run. - BackgroundDispatchUsingAsync (only caller: OutboxSweeper) passes continueOnCapturedContext: true to DispatchAsync/BulkDispatchAsync. - ExecuteWithResiliencePipelineAsync passes the caller's continueOnCapturedContext to Polly through a pooled ResilienceContext; that overload previously always used Polly's default of false. Still a known limit: a producer or outbox that awaits with ConfigureAwait(false) internally resumes on the pool. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Review of #4562 found that raising confirmations one at a time caps a producer with a database outbox at one mark-dispatched round trip at a time. The queue becomes BatchedCallbackQueue: one drain thread per producer, batches of up to 32 started in order inside one BrighterAsyncContext. A new ObservableUpDownCounter, paramore.brighter.publish_confirmation.queue.depth, reads live queues from an internal registry. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
…ix record - Release note: the sweep stays on its thread across awaited sends, elapsed-time scheduling, batched confirmations, empty ExecutionContext, the queue-depth metric, and continueOnCapturedContext reaching Polly (including the Proactor case). - ADR 0081 amendment: the batch size is a private constant of 32. - Bugfix record: review tests 9-13, fixes, and re-verification. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
…-starvation' Picks up master (#4552, #4563) merged on GitHub. release_notes.md: kept the updated #4560 cost paragraph and master's #4541 section; ADR index regenerated. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Follows the review changes in BrighterCommand/Brighter#4562: a Kafka sweep now resumes on the sweeper thread, the schedule uses elapsed time, and a new paramore.brighter.publish_confirmation.queue.depth instrument reports confirmation backlogs. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
RabbitMQ may coalesce publisher confirms: one basic.ack with multiple=true settles every delivery tag up to and including its own. The Sync producer settled only the exact tag, so the other messages the ack covered never raised a publish confirmation, and an outbox would send them again. RMQ.Async already handles this. Exposed by the intermittent rabbitmq-sync-ci failure of RmqConfirmationBatchingTests on #4562. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
A basic.nack with multiple=true rejects every delivery tag up to and including its own. Settle that whole range, as for acks, and have both handlers share one SettleConfirmations method. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
…peats a confirmation The producer raised a confirmation before removing it from the pending set. If an OnMessagePublished subscriber threw, the rest of a coalesced range was skipped, and because the ack handler is subscribed once per Send, the next invocation raised the same confirmation again. Claim each pending confirmation with TryRemove before raising it, and log a faulting sync subscriber rather than letting it escape the settle loop. The test channel now calls each ack/nack subscriber in turn and swallows its exception, as RabbitMQ.Client does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
…e old one After an IOException the producer resets its connection, and the next channel numbers its delivery tags from 1 again. Pending confirmations from the old channel were never cleared, so the new messages' entries collided with them, and an ack on the new channel (above all a coalesced one) confirmed messages the broker had never confirmed. Clear them before the reset, as RMQ.Async does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
The comments predate BatchedCallbackQueue, which starts callbacks in ack order but runs up to 32 at once. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
|
@claude This update addresses the issues that you raised in #4562 (comment) Review item 1 was fixed in two commits:
It's written up in three places:
The release note also gives the limit: a producer or outbox that awaits with ConfigureAwait(false) itself, as database clients and the AWS SDK do, still resumes on the pool. The fix removes Brighter's own trips back to the pool, not ones made inside other libraries. In addition, testing of this revealed an additional issue The RMQ.Sync producer now raises a confirmation for every message one RabbitMQ ack or nack covers, not only the last one. It also no longer raises a confirmation twice when a subscriber throws, and it no longer confirms old-channel messages after a connection reset. With the CI filter, the RMQ.Sync tests pass on net9.0 and net10.0 (143 passed, 1 skipped). The test that failed in CI was green in 20 of 20 repeated runs. Issues raised and not covered in this fix:
Other changes:
|
|
Claude finished @iancooper's task in 1m 46s —— View job Re-review of the follow-up changes (#4562)
Summary: Every item from the last review is addressed, and the fixes match the descriptions. The RMQ.Sync fix is correct and now matches RMQ.Async ( What I checked and agree with
1. 🟠 A hung sweep now emits no sweeper metrics at all, and ADR 0081 says it will show as
|
…pleted sweeps Sweeps now run one after another on the sweeper's thread, so a hung sweep records no tick lag, duration or outcome until it ends. lock_unavailable shows only an overrun on another instance sharing the lock. The stall alert is no increase in completed sweeps across instances; an alert on any outcome would never fire while the lock holder hangs. Also records the sync-over-async deadlock risk on the sweeper's thread. Raised in the re-review of #4562. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
Raised in the re-review of #4562; the RMQ.Sync comments were corrected in fa78671. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x
…-starvation' into bugfix/4560-timed-outbox-sweeper-starvation
|
Following up on item 1 of the re-review, the hung sweep that goes silent. In this PR (docs only):
Detection is #4573: a running-time gauge for the current sweep, and a Degraded health check that builds on #4555. Both are read on the exporter's or health-check host's thread, so they keep reporting while the sweeper's thread is stuck. The issue also asks whether our health checks should tell "restart" apart from "investigate and decide" for an operator. That makes it a feature, with its own spec and ADR, so it's not in this PR. |
Symptom
Some teams use
CommandProcessor.Postwith anInMemoryOutboxand a producer that confirms out of band, usually Kafka. Under load on Kubernetes (EKS),TimedOutboxSweepersweeps rarely or not at all. Unconfirmed messages build up in the in-memory outbox, and a pod restart loses them.Confirmed root cause
The bugfix record holds the full triage, the confirm evidence and the experiment tables. It is in
bugfixes/0052-timed-outbox-sweeper-starvation/bugfix.md.TimedOutboxSweeperranasync void Sweepfrom aSystem.Threading.Timer. When pool workers were blocked or kept busy, ticks were 3–4.5× late. A dedicated thread under the same load ticked on time (Experiments A and C).Task.Run, so Kafka and RMQ.Sync could run 50 of 50 and 15–23 handlers at once. A long queue of short work items did not delay ticks on its own (Experiment A, flood mode). It does add pool load, though, and the confirmations completed out of order.async voidand crashed the process.CreateScope()sat outside thetry, so a failure there leaked the distributed lock.StopAsyncdid not wait for a sweep already running.TimerIntervalof 0 swept once and then never again.Fix
StopAsyncwaits for a sweep in flight.TimerInterval < 1throws aConfigurationException.SerialCallbackQueueinParamore.Brighter.Tasks. It has one dedicated thread per producer, raises callbacks one at a time in order, and runs them inBrighterAsyncContext. Kafka and RMQ.Sync use it.Completeis dropped rather than thrown.0081-timed-outbox-sweeper-health-metrics:paramore.brighter.outbox_sweeper.tick.lag,.sweep.durationand.sweeps, the last tagged by outcome (completed,failed,lock_unavailable).AddBrighterInstrumentationregisters them; without it the sweeper uses a no-op meter.## Mastercovers the behaviour changes and the cost of one thread per confirming producer on 1–2 CPU pods.Known limits (documented, not fixed)
ConfigureAwait(false)internally, as database clients and the AWS SDK do, still resumes on the pool during a sweep. Kafka's send no longer does. See "Changes after review".ConfigureAwait(false)continuation inside a confirmation callback can still land on the pool.Changes after review
The review comment is addressed below. Each item has a test committed red first, or a characterisation test checked with a named mutation.
ProduceAsyncwithRunContinuationsAsynchronously; I checked this in its IL. So under starvation a Kafka sweep dispatched one message and then stalled. The new test fails with "dispatched 1 of 5".BrighterAsyncContextwas not enough on its own. AConfigureAwait(false)continuation is not inlined on a thread that has aSynchronizationContext; it goes to the pool. Two such hops had to go:continueOnCapturedContext: true;ExecuteWithResiliencePipelineAsyncnow passes the caller'scontinueOnCapturedContextto Polly. Before, it always used Polly's default offalse.PostAsyncorClearOutboxAsync: Polly's continuations, such asOnRetry, now stay on the pump thread. The mediator's and Kafka's own awaits already did. A characterisation test pins this, and the release note covers it. This adds no new deadlock: blocking onPostAsyncfrom the pump thread already deadlocked.SerialCallbackQueueis renamedBatchedCallbackQueue.paramore.brighter.publish_confirmation.queue.depth, bymessaging.systemandmessaging.destination.name.ExecutionContext(review Support fluent configuration as well as attibutes #4), so they no longer keep the starting caller'sAsyncLocalstate.StopAsync's continuation needs a pool thread. The host's shutdown timeout bounds it.WaitForConfirmationCallbackslogs how many callbacks are still in flight.RMQ.Sync coalesced publisher confirms (found by this PR's CI)
Symptom.
rabbitmq-sync-cifailed intermittently.RmqConfirmationBatchingTests...overlapping_batches_off_the_thread_pool, with "Not every message was confirmed".Root cause. This predates the branch; the same code is on master. This PR's new 50-message test exposed it.
basic.ackorbasic.nackwithmultiple=truecovers every delivery tag up to and including its own.e.Multiple. I checked this in its IL.e.DeliveryTag. The other messages the ack covered never raised a confirmation, so an outbox would send them again.Fixes (
RmqMessageProducer, each with a red-first regression test):DeliveryTag, in ascending order.TryRemovebefore it is raised, and a syncOnMessagePublishedsubscriber that throws is caught and logged.Send, the next call also raised the same confirmation again.IOExceptionreset, as RMQ.Async does.Tests. The regression tests use a
DispatchProxychannel,CoalescingConfirmsRmqChannel. It holds back the real per-tag acks and raises one frame on demand, so the tests don't depend on broker timing.bugfixes/0053-rmq-sync-multiple-confirm/bugfix.md.Follow-ups (not fixed here):
Sendsubscribes the ack/nack handlers and callsConfirmSelect()on every call. This is a performance problem, not a correctness one.IOExceptionpath running. Acks could then confirm messages pending on the old channel, and range settlement makes one such ack confirm a whole range. This is unverified and needs its own/bugfixdiagnosis.Tests
Regression tests: see
bugfixes/0052-timed-outbox-sweeper-starvation/bugfix.md(Regression Tests 1–13).Full suites, net9.0 and net10.0, at the merge with master that includes #4552 and #4563:
maryqueue.Failures that remain are environmental:
maryqueue with different settings.Extensions.AspNetCore.Tests is intermittently flaky. A static logger in
UseInboxHandler<T>is created from a logger factory that a disposed test host left behind.UseInboxHandler<T>is missing from that project'sInitializerwarm-up list. This predates this branch, and a rerun passes 65/65.Merge notes
Picked up master (fix(outbox): enforce maxOutStandingMessages for an async-only outbox #4552, fix(service-activator)!: coordinate consumer startup and shutdown #4563) as merged on GitHub. In
release_notes.md, both sections are kept. fix(outbox): enforce maxOutStandingMessages for an async-only outbox #4552's outstanding-count check still runs inTask.Run, with noSynchronizationContext, so it never runs on the sweeper or pump thread.Merged
master.KafkaMessageProducer.Disposekeeps fix(kafka): make producer and consumer disposal idempotent #4556's idempotentKafkaResourceCleanupand completes the confirm queue as one more guarded step.Master also added an ADR numbered 0081. Both are kept, because the index's identity is the filename stem.
Docs
The user guide's sweeper metrics section (alerts, histogram buckets,
System.Runtimethread-pool metrics) is in BrighterCommand/Docs#199.Fixes #4560
🤖 Generated with Claude Code
https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x