Skip to content

fix: TimedOutboxSweeper starvation and publish-confirm floods (#4560) - #4562

Merged
iancooper merged 42 commits into
masterfrom
bugfix/4560-timed-outbox-sweeper-starvation
Oct 10, 2026
Merged

iancooper merged 42 commits into
masterfrom
bugfix/4560-timed-outbox-sweeper-starvation

Conversation

@iancooper

@iancooper iancooper commented Oct 8, 2026 •

Copy link
Copy Markdown
Member

Symptom

Some teams use CommandProcessor.Post with an InMemoryOutbox and a producer that confirms out of band, usually Kafka. Under load on Kubernetes (EKS), TimedOutboxSweeper sweeps 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.

  • What starved the sweeper: its schedule depended on the shared thread pool. TimedOutboxSweeper ran async void Sweep from a System.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).
  • What made it worse: each confirmation was raised with its own 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.
  • Defects found along the way:
    • An exception in a sweep escaped async void and crashed the process.
    • CreateScope() sat outside the try, so a failure there leaked the distributed lock.
    • StopAsync did not wait for a sweep already running.
    • A TimerInterval of 0 swept once and then never again.
  • Out of scope: the outstanding-message check, which runs on every Post under default DI and blocks the pool, belongs to Each Post sorts the whole in-memory outbox to count outstanding messages, even with no limit set #4554. Our evidence is posted there.

Fix

  • The sweeper runs on a dedicated thread. The schedule keeps to its due times. It re-anchors only when a sweep starts a whole interval or more late, so it never sweeps back to back to catch up.
    • StopAsync waits for a sweep in flight.
    • A failed sweep is caught, logged and recorded.
    • The lock is released on every path.
    • TimerInterval < 1 throws a ConfigurationException.
  • Publish confirmations drain through a new SerialCallbackQueue in Paramore.Brighter.Tasks. It has one dedicated thread per producer, raises callbacks one at a time in order, and runs them in BrighterAsyncContext. Kafka and RMQ.Sync use it.
    • The queue is unbounded on purpose. A bounded queue would block Confluent's delivery thread and the RMQ connection loop.
    • A confirmation that arrives after Complete is dropped rather than thrown.
  • Sweeper health metrics, per ADR 0081-timed-outbox-sweeper-health-metrics: paramore.brighter.outbox_sweeper.tick.lag, .sweep.duration and .sweeps, the last tagged by outcome (completed, failed, lock_unavailable). AddBrighterInstrumentation registers them; without it the sweeper uses a no-op meter.
  • A release note under ## Master covers the behaviour changes and the cost of one thread per confirming producer on 1–2 CPU pods.

Known limits (documented, not fixed)

  • A producer or outbox that awaits with 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".
  • A 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.

  1. A sweep now stays on the sweeper thread across awaited sends (review Support for multiple Application Layer Protocols in Task Queues #1/No API Documentation #2).
    • Why it didn't before: Confluent.Kafka 2.15.1 completes ProduceAsync with RunContinuationsAsynchronously; 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".
    • Running the sweep in a BrighterAsyncContext was not enough on its own. A ConfigureAwait(false) continuation is not inlined on a thread that has a SynchronizationContext; it goes to the pool. Two such hops had to go:
      • background dispatch, whose only caller is the sweeper, now passes continueOnCapturedContext: true;
      • ExecuteWithResiliencePipelineAsync now passes the caller's continueOnCapturedContext to Polly. Before, it always used Polly's default of false.
    • Effect on a Proactor handler that calls PostAsync or ClearOutboxAsync: Polly's continuations, such as OnRetry, 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 on PostAsync from the pump thread already deadlocked.
  2. Publish confirmations drain in batches of up to 32 on the producer's thread (review Support async versions of Send and Publish #3). One at a time capped a database outbox at one mark-dispatched round trip at a time.
    • SerialCallbackQueue is renamed BatchedCallbackQueue.
    • The Kafka and RMQ.Sync "one at a time, in order" tests are replaced. The new tests assert that callbacks overlap, never exceed the batch size, and never run on the pool.
    • New metric: paramore.brighter.publish_confirmation.queue.depth, by messaging.system and messaging.destination.name.
    • ADR 0081 is amended for both changes.
  3. The sweeper and confirmation threads start with an empty ExecutionContext (review Support fluent configuration as well as attibutes #4), so they no longer keep the starting caller's AsyncLocal state.
  4. The schedule uses elapsed time, not the wall clock (review Automate the build #5), so a clock stepping backwards no longer delays the next sweep.
  5. A race found while testing: the sweeper re-checks its due time after creating its timer. Under load, a test clock could otherwise miss the due time.
  6. Not acted on:
    • 6a, StopAsync's continuation needs a pool thread. The host's shutdown timeout bounds it.
    • 6b, a sweep in flight isn't cancelled on stop. That's by design, and a test covers it.
    • 6c, logging the backlog at dispose: this is already done. WaitForConfirmationCallbacks logs how many callbacks are still in flight.

RMQ.Sync coalesced publisher confirms (found by this PR's CI)

Symptom. rabbitmq-sync-ci failed intermittently.

  • The failing test was RmqConfirmationBatchingTests...overlapping_batches_off_the_thread_pool, with "Not every message was confirmed".
  • In the same run it timed out after 30 s on one framework and passed in 1 s on the other.

Root cause. This predates the branch; the same code is on master. This PR's new 50-message test exposed it.

  • RabbitMQ may coalesce publisher confirms: one basic.ack or basic.nack with multiple=true covers every delivery tag up to and including its own.
  • RabbitMQ.Client 6.8.1 passes that flag through unchanged as e.Multiple. I checked this in its IL.
  • The Sync producer settled only e.DeliveryTag. The other messages the ack covered never raised a confirmation, so an outbox would send them again.
  • RMQ.Async already handled this.

Fixes (RmqMessageProducer, each with a red-first regression test):

  1. Coalesced acks and nacks. Settle every pending tag up to and including DeliveryTag, in ascending order.
  2. A throwing subscriber. Each confirmation is claimed with TryRemove before it is raised, and a sync OnMessagePublished subscriber that throws is caught and logged.
    • Before, a throw skipped the rest of the range.
    • Because the handler is subscribed once per Send, the next call also raised the same confirmation again.
  3. Channel reset. Clear pending confirmations on the IOException reset, as RMQ.Async does.
    • The new channel numbers its tags from 1, so its acks were confirming the old channel's messages, which the broker had never confirmed.
    • Range settlement would have made that a whole range of such messages per ack.

Tests. The regression tests use a DispatchProxy channel, CoalescingConfirmsRmqChannel. It holds back the real per-tag acks and raises one frame on demand, so the tests don't depend on broker timing.

  • With the CI filter, RMQ.Sync gives 143 passed, 1 skipped on net9.0 and net10.0.
  • The original CI test and the 4 new tests were green in 20 of 20 repeated runs.
  • The full record is in bugfixes/0053-rmq-sync-multiple-confirm/bugfix.md.

Follow-ups (not fixed here):

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:

  • Core.Tests: 1699 passed, 7 skipped.
  • Kafka.Tests: 265/265 on each framework, before the merge.
  • RMQ.Sync.Tests: everything passes except the 9 mTLS tests.
  • Extensions, Testing and InMemory: all pass.
  • RMQ.Async: everything passes except mTLS and the shared mary queue.

Failures that remain are environmental:

  • the mTLS tests need generated certificates;
  • the RMQ.Sync and RMQ.Async dispatcher tests share a mary queue 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's Initializer warm-up list. This predates this branch, and a rerun passes 65/65.

Merge notes

Docs

The user guide's sweeper metrics section (alerts, histogram buckets, System.Runtime thread-pool metrics) is in BrighterCommand/Docs#199.

Fixes #4560

🤖 Generated with Claude Code

https://claude.ai/code/session_01SpjnqGWbARr3Z3ofGc6C2x

iancooper and others added 9 commits October 8, 2026 17:36
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
@iancooper

Copy link
Copy Markdown
Member Author

@claude Let's review the changes here with a view to trying to remove issues with respect to thread starvation

@claude

claude Bot commented Oct 8, 2026 •

Copy link
Copy Markdown

Claude encountered an error after 0s —— View job


I'll analyze this and get back to you.

1 similar comment
@claude

claude Bot commented Oct 8, 2026 •

Copy link
Copy Markdown

Claude encountered an error after 0s —— View job


I'll analyze this and get back to you.

@iancooper

Copy link
Copy Markdown
Member Author

@claude Let's review the changes here with a view to trying to remove issues with respect to thread starvation

@claude

claude Bot commented Oct 8, 2026 •

Copy link
Copy Markdown

Claude finished @iancooper's task in 3m 7s —— View job


Review: what can still starve the sweeper (PR #4562)

  • Gather context (diff, bugfix record, mediator/producer call paths)
  • Review TimedOutboxSweeper
  • Review SerialCallbackQueue
  • Review Kafka / RMQ.Sync producer changes
  • Review starvation tests
  • Post findings

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

TimedOutboxSweeper.cs SweepOnce:

return SweepAsync().GetAwaiter().GetResult();

The sweeper thread has no SynchronizationContext, so every await continuation inside the sweep goes back to the pool. Here is the Kafka path from the sweeper:

  • OutboxSweeper.SweepAsync
  • → ClearOutstandingFromOutboxAsync
  • → DispatchAsync (with continueOnCapturedContext: false)
  • → KafkaMessageProducer.SendWithDelayAsync
  • → KafkaMessagePublisher.PublishMessageAsync
  • → await producer.ProduceAsync(...) (KafkaMessagePublisher.cs:53)

As far as I know, Confluent completes ProduceAsync with a TaskCompletionSource created with RunContinuationsAsynchronously, from librdkafka's delivery thread. I couldn't open the package source here to check. If that holds, the continuation at line 53 is queued to the pool with no context. With every pool worker blocked, a 100-message batch moves forward one message each time the pool injects a thread (about 1–2 per second once it is at its limit). The sweep starts on time and then crawls. That looks the same as the symptom in #4560.

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:

  • Run the sweep inside the existing BrighterAsyncContext on the sweeper thread:
    return BrighterAsyncContext.Run(SweepAsync);
  • KafkaMessagePublisher and KafkaMessageProducer use plain awaits. Their continuations would post back to the sweeper thread.
  • The ConfigureAwait(false) calls in DispatchAsync resume on whichever thread completed the inner task. That would be the sweeper thread, so they stay off the pool too.
  • What still needs the pool: database outboxes and locks (RelationDatabaseOutbox has 20 ConfigureAwait(false)), plus Linux socket completions. Those are the honest known limits.
  • SerialCallbackQueue already uses this pattern, so the sweeper would match it.

Fix this →

2. 🔴 The starvation test can't detect #1

When_thread_pool_workers_are_all_blocked_should_still_sweep_the_outbox uses InMemoryMessageProducer, whose SendAsync completes inline. No await in the sweep ever yields, so the test only shows that the tick fires.

Add a test double whose SendAsync returns a task completed from a dedicated non-pool thread through a TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously), the way Confluent does. Send two or more messages, or wait on the outbox being marked dispatched. That test should fail on the current branch and pass once #1 is fixed. Per CLAUDE.md, write it with /bugfix:test / /test-first.

3. 🟠 One confirmation queue per producer means one slow callback stalls all confirmations

SerialCallbackQueue fixes the flood of Task.Run calls, but it now runs one confirmation at a time per producer.

  • In-memory outbox: MarkDispatchedAsync is effectively synchronous, so this is fine.

  • Database outbox: each callback is a database round trip. Its ConfigureAwait(false) continuations land on the pool, so under starvation the single drain thread waits on the pool once per message. A database stall or Polly retry back-offs also hold up every confirmation behind it.

  • The queue is unbounded on purpose, so a backlog shows up as:

    1. growing memory;
    2. messages that stay unmarked past MinimumMessageAge;
    3. the sweeper resending them, which produces more confirmations.

    That feedback loop resembles the original symptom: messages pile up as undispatched.

Suggestions:

  • Minimum: add a queue-depth gauge alongside the ADR 0081 metrics, for example paramore.brighter.publish_confirmation.queue.depth, so the backlog is visible.
  • Consider: marking a message dispatched doesn't need ordering. A small fixed number of drain threads (2–4), or draining everything queued and awaiting it together, would cap thread use without serialising database I/O. If strict ordering is a requirement, the ADR or bugfix record should say why. Today the record only notes that out-of-order completion was observed.

4. 🟡 Both new threads capture the starting caller's ExecutionContext

new Thread(Drain).Start() (SerialCallbackQueue.cs:751) and new Thread(SweepUntilStopped).Start() both capture the ExecutionContext of whoever starts them.

  • For the queue, that is whoever made the first Enqueue call. On the Kafka async path this is a ProduceAsync continuation inside a Post, with that request's AsyncLocals attached: logging scopes, baggage, RequestContext-related state.
  • The thread keeps those values for the producer's lifetime, so they never get collected, and every later confirmation runs with them.
  • CreateConfirmationSpan already avoids nesting under Activity.Current, so tracing is protected. Other AsyncLocal users are not.

Wrap both starts in using (ExecutionContext.SuppressFlow()) { ... }. That works on netstandard2.0, where UnsafeStart doesn't exist. The old Task.Run captured a fresh context for each item, so this is a small behaviour change too.

5. 🟡 The schedule uses the wall clock

WaitUntil and NextDue work from _timeProvider.GetUtcNow(). If the node's clock steps backwards (an NTP correction, or a VM resuming), due - now grows by the size of the step. The real-time WaitAny timeout then wakes the thread and goes back to sleep until the wall clock catches up. In other words, the sweeper skips sweeps again. Scheduling from GetTimestamp() / GetElapsedTime() (monotonic) avoids that. Tick lag could still be reported from the same monotonic values.

6. ⚪ Minor

  • StopAsync needs a pool thread to finish. _stopped uses RunContinuationsAsynchronously, so the StopAsync continuation is queued to the pool. On a starved pool, shutdown waits for a worker. The host's shutdown timeout bounds this, so it's low impact, but it fits this PR's theme.
  • StopAsync doesn't cancel a sweep in progress. It cancels _stopping, but SweepAsync passes CancellationToken.None to the lock and the sweep. Shutdown therefore waits for a whole batch. On a starved pool combined with Support for multiple Application Layer Protocols in Task Queues #1, that can outlast the host timeout.
  • Callbacks can outlive Dispose. WaitForConfirmationCallbacks gives up after its timeout, and then Complete() lets the drain thread keep running callbacks after the producer, and possibly the DI container, has been disposed. The callbacks catch their own errors, so it is safe. Because the queue is now serial, though, the backlog at dispose can be longer than it used to be. A log line with the remaining depth would help.

What looks good

  • The re-anchoring rule (only after a whole interval of lag) avoids catch-up bursts.
  • The lock is released in finally on every path, and CreateScope now sits inside the try.
  • async void and process-crash paths are gone. The meter is guarded, so metrics can't fail a sweep.
  • Dropping a confirmation after Complete instead of throwing it back onto Confluent's or RMQ's thread is the right choice.

Suggested order: #2 then #1 (red, then green), then #3's gauge. #4 and #5 are small and could go in as separate commits.

iancooper and others added 3 commits October 9, 2026 01:36
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
iancooper and others added 2 commits October 9, 2026 11:24
…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
iancooper added a commit to BrighterCommand/Docs that referenced this pull request Oct 9, 2026
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
iancooper and others added 7 commits October 9, 2026 14:37
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
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
@iancooper

Copy link
Copy Markdown
Member Author

@claude This update addresses the issues that you raised in #4562 (comment)

Review item 1 was fixed in two commits:

  • eece13a (the test): When_..._complete_off_the_pool_should_still_sweep_a_batch, using a new OffPoolCompletingProducer test double. It sends complete on another thread while the thread pool is starved, the way Confluent.Kafka completes ProduceAsync. Before the fix, it failed with "dispatched 1 of 5".
  • 80512a4 (the fix): each sweep now runs in BrighterAsyncContext.Run, so it stays on the sweeper's own thread. That alone wasn't enough, because a ConfigureAwait(false) continuation still went back to the pool, so two more changes were needed:
    • in TimedOutboxSweeper, background dispatch now passes continueOnCapturedContext: true;
    • in OutboxProducerMediator, ExecuteWithResiliencePipelineAsync now passes the caller's continueOnCapturedContext to Polly. Before, it always used Polly's default of false.

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

claude Bot commented Oct 9, 2026 •

Copy link
Copy Markdown

Claude finished @iancooper's task in 1m 46s —— View job


Re-review of the follow-up changes (#4562)

  • Gather context (diff since the last review)
  • Check items 1/2: sweep stays on the sweeper thread, plus the new test
  • Check item 3: BatchedCallbackQueue and the queue-depth gauge
  • Check items 4/5: ExecutionContext suppression, monotonic schedule
  • Review the RMQ.Sync coalesced-confirm fix (0053)
  • Post findings

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 (IsConfirmedBy / RemovePendingConfirmations). I found one doc gap that matters for operators, and one place where stale comments are still in the code. Neither blocks the starvation fix. As before, I only read the code: I didn't build it or run the tests.


What I checked and agree with

Item Where Verdict
1. Sweep stays off the pool TimedOutboxSweeper.cs:233 BrighterAsyncContext.Run(SweepAsync); OutboxProducerMediator.cs background dispatch passes true; Polly gets ResilienceContextPool.Shared.Get(continueOnCapturedContext, ct) and the context is returned in finally ✅ ClearOutstandingFromOutboxAsync's only caller is OutboxSweeper, so hard-coding true there is safe
2. Test can detect it OffPoolCompletingProducer + ..._complete_off_the_pool_should_still_sweep_a_batch ✅
3. Batched drain + depth gauge BatchedCallbackQueue (32 per batch, in BrighterAsyncContext), PublishConfirmationMeter, a registry holding weak references ✅ Complete/Drain both unregister, so the race between an Enqueue and Complete is harmless
4. Empty ExecutionContext ExecutionContext.SuppressFlow() around both Thread.Start calls ✅
5. Monotonic schedule GetTimestamp/GetElapsedTime, plus a re-check after the timer is created ✅
0053 RMQ.Sync TryRemove claims a tag before raising it; tags are settled in ascending order; the sync subscriber is wrapped in try; Clear() runs on the IOException reset ✅

1. 🟠 A hung sweep now emits no sweeper metrics at all, and ADR 0081 says it will show as lock_unavailable

ADR 0081 (line 57 in "Silent stalls", and lines 222–223) says lock_unavailable "means a sweep is overrunning or hung". That held while each tick was a separate timer callback. It no longer holds within one process:

  • The schedule now runs on one thread, so while SweepOnce is blocked, SweepUntilStopped cannot reach WaitUntil again. No tick runs, so nothing records tick.lag, sweep.duration or sweeps.
  • lock_unavailable only appears if another pod shares the distributed lock. With one replica or the in-memory lock, a hung sweep is completely silent. Tick lag is recorded only when a sweep starts, so it never shows the hang either.
  • The new context makes a hang slightly more likely. Code in the sweep path that blocks on a task (sync-over-async) can deadlock on the sweeper thread. This includes a custom IDistributedLock, producer or outbox that does .GetAwaiter().GetResult() on a plain await. Brighter's own async paths look clean; I checked the GetAwaiter().GetResult() sites in producers and outboxes, and they are sync-only methods. Third-party implementations are the risk. Before this change the same code only stalled a pool thread.

Suggested fix (docs only): in ADR 0081 and Docs#199, say that the main alert for a stall is no increase in paramore.brighter.outbox_sweeper.sweeps (any outcome) over N intervals. Say lock_unavailable points to an overrun on another instance. The meter needs no change.

Fix this →

2. 🟡 Kafka still says confirmations run "one at a time, in delivery order"

The update says these comments were corrected. The RMQ.Sync ones were, but KafkaMessageProducer.cs still has two:

  • :80: "Confirmation raises run one at a time, in delivery order, on a single worker…"
  • :461: "…confirmations reach subscribers one at a time, in delivery order."

Both should say "in batches of up to 32, started in delivery order", matching the RMQ.Sync wording.

Fix this →

3. ⚪ Minor (no action needed)

  • The default flag is applied unevenly. Some calls to ExecuteWithResiliencePipelineAsync don't pass a flag: mark-dispatched in DispatchAsync/BulkDispatchAsync at lines 1162 and 1230, EndBatchAddToOutboxAsync, and HandleAsyncPublishConfirmation. They pick up the default true, which now reaches Polly. This is fine for the sweeper and for the confirmation queue, where it keeps Polly on the drain thread, which is what you want.
    • It's also not a new deadlock: each of those outer awaits already captured the context.
    • It does mean "passes the caller's choice to Polly" is only true where the caller passes one. That may be worth a word in the release note.
  • SettleConfirmations with multiple takes a sorted copy of _pendingConfirmations.Keys on every coalesced ack. That's fine at RMQ confirm rates; RMQ.Sync producer subscribes confirm handlers and calls ConfirmSelect on every Send #4567 is the bigger cost on that path.
  • The RMQ.Sync producer: auto-recovered channel can confirm messages pending on the old channel #4568 write-up is honest about the risk. I agree it needs its own /bugfix diagnosis before this ships in a release that has automatic recovery turned on.

Bottom line: For the starvation problem (#4560) this is in good shape. Brighter's own code no longer sends the sweeper or confirmation thread back to the pool. The remaining pool dependencies are in third-party ConfigureAwait(false) code, and the release note covers them. I'd make the two doc fixes above before merging.

iancooper and others added 3 commits October 9, 2026 17:03
…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
@iancooper

Copy link
Copy Markdown
Member Author

Following up on item 1 of the re-review, the hung sweep that goes silent.

In this PR (docs only):

  • ADR 0081 (e373af9) now says that in-process, a hung sweep records nothing at all, and that lock_unavailable only shows an overrun on another instance.
  • Docs#199 (BrighterCommand/Docs@6468dd9) says the same. Both also record the sync-over-async deadlock risk on the sweeper's thread.
  • The stall alert stays on completed sweeps summed across instances, not "any outcome". While the lock holder is hung, the other instances keep recording lock_unavailable, so an alert on any outcome would never fire.
  • The Kafka comments (item 2) are corrected in 26cb7b9.

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.

@iancooper
iancooper merged commit d1251f8 into master Oct 10, 2026
56 of 58 checks passed
@iancooper
iancooper deleted the bugfix/4560-timed-outbox-sweeper-starvation branch October 10, 2026 11:09
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

3 - Done Bug .NET Pull requests that update .net code V10.X

Projects

None yet

Development

Successfully merging this pull request may close these issues.

TimedOutboxSweeper misses ticks when pool threads are blocked, and one Task.Run per publish confirm delays MarkDispatched under load

1 participant