Skip to content

Fix consumer hangs after mailbox failure - #374

Open
bartoszs321 wants to merge 6 commits into
fsprojects:developfrom
bartoszs321:fix/consumer-mailbox-failure-recovery
Open

bartoszs321 wants to merge 6 commits into
fsprojects:developfrom
bartoszs321:fix/consumer-mailbox-failure-recovery

Conversation

@bartoszs321

Copy link
Copy Markdown
Contributor

Summary

  • keep failed single-topic consumer mailboxes draining requests so calls made after shutdown complete instead of hanging
  • propagate multi-topic poller failures to the parent consumer and fail pending receives
  • drain failed multi-topic mailbox requests and close child consumers when DisposeAsync is called
  • add integration coverage for single-topic and multi-topic mailbox failures

Problem

When a consumer mailbox faulted, its reader stopped. Requests posted afterwards were accepted by the unbounded channel but never received a response. This could leave operations such as GetStats, ReceiveAsync, and DisposeAsync waiting indefinitely. A multi-topic poller failure was only logged, leaving the parent consumer apparently ready even though it could no longer deliver messages.

Implementation

After a mailbox exits, a lightweight drain continues reading its channel. Requests that can no longer be performed fail with AlreadyClosedException; close requests complete after the necessary local or child-consumer cleanup. Multi-topic poller failures are posted to the parent mailbox, which moves the consumer to a failed state and completes pending receives.

Verification

  • Unit tests: 216 passed, 2 skipped
  • Mailbox integration tests against Apache Pulsar 3.3.0: 3 passed

bartoszs321 and others added 2 commits October 5, 2026 14:07
Co-authored-by: Cursor <cursoragent@cursor.com>
…lbox-failure-recovery

Co-authored-by: Cursor <cursoragent@cursor.com>

# Conflicts:
#	tests/IntegrationTests/Basic.fs

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Faulted single-topic consumers remain in the wrong lifecycle state after disposal, and pending-receive coverage is incomplete.

Review effort: Balanced
Findings: 2 Medium severity

Open (2)
What changed in this PR

Improves consumer failure handling so mailbox failures no longer leave requests hanging.

Changes:

  • Drains stopped consumer mailboxes and completes requests.
  • Propagates multi-topic poller failures.
  • Adds mailbox-failure integration coverage.
File Description
ConsumerImpl.fs Drains requests after mailbox termination.
MultiTopicsConsumerImpl.fs Handles poller failures and child cleanup.
Basic.fs Tests mailbox-failure behavior.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread src/Pulsar.Client/Internal/ConsumerImpl.fs
Comment thread tests/IntegrationTests/Basic.fs
let mb = Channel.CreateUnbounded<ConsumerMessage<'T>>(UnboundedChannelOptions(SingleReader = true, AllowSynchronousContinuations = true))

let replyFromStoppedMailbox (msg: ConsumerMessage<'T>) =
let ex = AlreadyClosedException "Consumer is already closed"

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This should be NotConnectedException

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Updated in 29edfc2

@Lanayx

Lanayx commented Oct 5, 2026

Copy link
Copy Markdown
Member

Thank you for the PR, would you be able to implement similar behavior for Producer, so it will be single PR for both

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Stopped consumers can retain pooled payloads, and failed partition cleanup is incorrectly reported as successful.

Review effort: Balanced
Findings: 1 Medium severity

Open (1)
Resolved since last review (2)
Previously missed (1)

In code that hasn't changed since last review

Medium severity Mailbox failure drops messages without disposing payloads

src/​Pulsar.Client/​Internal/​ConsumerImpl.fs:891

The catch-all also drops queued MessageReceived values without disposing their RawMessage.Payload. These payloads come from MemoryStreamManager, and every normal discard path explicitly disposes them (for example, ConsumerImpl.fs:1002-1052), so messages racing with a mailbox failure retain pooled buffers. Add a stopped-mailbox case that disposes the payload before discarding the message.

Comment thread src/Pulsar.Client/Internal/PartitionedProducerImpl.fs

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

A partitioned producer request can still hang when that request triggers the mailbox failure, and the multi-topic test does not exercise waiter cleanup.

Review effort: Balanced
Findings: 1 High severity

Open (1)
Resolved since last review (1)

Comment on lines +339 to +340
finally
drainStoppedMailbox())

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants