Skip to content

[fix][client] Serialize chunked-message bookkeeping to fix use-after-free and count/queue drift - #26084

Open
SongOf wants to merge 2 commits into
apache:masterfrom
SongOf:fix/client-chunked-message-bookkeeping-race
Open

[fix][client] Serialize chunked-message bookkeeping to fix use-after-free and count/queue drift#26084
SongOf wants to merge 2 commits into
apache:masterfrom
SongOf:fix/client-chunked-message-bookkeeping-race

Conversation

@SongOf

@SongOf SongOf commented Jun 24, 2026

Copy link
Copy Markdown
Contributor

Motivation

ConsumerImpl's chunked-message reassembly state — the per-uuid ChunkedMessageCtx,
its chunkedMsgBuffer, pendingChunkedMessageCount and pendingChunkedMessageUuidQueue
— is mutated from two different threads with no synchronization:

  • the receive/assembly path (processMessageChunk, and the last-chunk finalize in
    messageReceived) runs on the Netty IO event-loop thread
    (ClientCnx.handleMessage calls consumer.messageReceived(...) directly);
  • the incomplete-chunk expiry path (removeExpireIncompleteChunkedMessages) runs on
    the client's internalPinnedExecutor — a separate single-thread pool
    (client.getInternalExecutorService()), not the eventLoopGroup.

When a late chunk for a uuid arrives while the expiry task is removing that same ctx,
the expiry thread can release() / recycle() a ChunkedMessageCtx and its buffer while
the receive thread is still writing into it. This races into:

  • use-after-free / double-free of chunkedMsgBuffer (IllegalReferenceCountException,
    or worse — writing into memory the allocator already handed to someone else);
  • double ChunkedMessageCtx.recycle(), which corrupts the Netty Recycler pool and can
    hand the same instance to two different chunked messages;
  • pendingChunkedMessageCount drift (non-atomic int mutated from both threads).

Incomplete-chunk expiry is enabled by default (expireTimeOfIncompleteChunkedMessageMillis
= 1 minute), so any consumer of chunked messages is exposed.

Separately, a redelivered first chunk (chunkId == 0 for a uuid that already has an
in-progress ctx) replaced the old ctx without decrementing pendingChunkedMessageCount
and re-enqueued the uuid, so the counter over-counted and pendingChunkedMessageUuidQueue
ended up with duplicate / mis-ordered entries (the queue is meant to be ordered oldest-first
by receivedTime, with one entry per in-progress uuid, since removeExpire/removeOldest
clean from the head).

Modifications

  • Add a dedicated lock chunkedMessageLock. processMessageChunk,
    removeOldestPendingChunkedMessage and removeExpireIncompleteChunkedMessages now acquire
    it (thin wrapper over an extracted body), and messageReceived holds it across the
    last-chunk assembly + finalize so the assemble→finalize window is closed. All chunk
    bookkeeping is serialized, and pendingChunkedMessageCount is mutated only under the lock
    (so it needs no volatile/Atomic). Heavy per-message work (newMessage, decryption of
    non-chunk payloads, executeNotifyCallback) stays outside the lock, and non-chunked
    messages never touch it.
    • The lock only ever guards chunk bookkeeping; the calls made under it
      (ConcurrentHashMap ops, doAcknowledge which is async/non-blocking) never acquire it
      in reverse, so there is no new lock-ordering / deadlock and no blocking call held across
      the lock.
  • Fix the redelivered-first-chunk path: decrement pendingChunkedMessageCount for the
    replaced ctx, remove(uuid) its stale entry from pendingChunkedMessageUuidQueue and
    re-add(uuid) at the tail — keeping the queue ordered oldest-first with exactly one entry
    per in-progress uuid (net 0 count change on replace, +1 on a genuinely new uuid).

This is an internal client-side locking change only; no public API, wire protocol, schema,
config defaults or metrics are changed.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • ConsumerImplTest.testChunkedMessageCountRaceBetweenReceiveAndExpiry — drives
    processMessageChunk on a "receiver" thread concurrently with the real
    removeExpireIncompleteChunkedMessages on an "expirer" thread, for the "a late chunk
    arrives for a uuid that is concurrently being expired" scenario. Verified to fail before
    the fix
    (IllegalReferenceCountException / corrupted pendingChunkedMessageCount) and
    pass after.
  • ConsumerImplTest.testDuplicateFirstChunkOvercountsPendingChunkedMessageCount
    deterministic; delivers a redelivered first chunk and asserts
    pendingChunkedMessageCount == chunkedMessagesMap.size() and that the uuid appears exactly
    once in pendingChunkedMessageUuidQueue. Verified to fail before and pass after.

Full pulsar-client unit suite (712 tests) passes with no regressions.

Does this pull request potentially affect one of the following parts:

  • Dependencies (add or upgrade a dependency)
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

Documentation

  • doc-required
  • doc-not-needed
    (internal bug fix; no user-facing behavior or config change)
  • doc
  • doc-complete

Comment thread pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java Outdated
Comment thread pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java Outdated
@SongOf
SongOf force-pushed the fix/client-chunked-message-bookkeeping-race branch 2 times, most recently from 69fb6f6 to 910db42 Compare June 25, 2026 03:17
@SongOf
SongOf requested a review from lhotari June 25, 2026 04:00

@congbobo184 congbobo184 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.

LGTM

@congbobo184 congbobo184 added the type/bug The PR fixed a bug or issue reported a bug label Jun 28, 2026
@congbobo184 congbobo184 added this to the 5.0.0-M2 milestone Jun 28, 2026
@SongOf
SongOf force-pushed the fix/client-chunked-message-bookkeeping-race branch 3 times, most recently from d9b2d45 to 33da7ad Compare June 29, 2026 17:14

@lhotari lhotari left a comment

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.

Review performed with AI assistance (Claude Code / Claude Fable 5 combined with a Codex gpt-5.6-sol review pass); findings below were verified against the code before posting.

Overall the change looks sound and worthwhile. The receive/assembly path (Netty IO thread) and the incomplete-chunk expiry path (internalPinnedExecutor) genuinely race on ChunkedMessageCtx, its buffer and pendingChunkedMessageCount, and the new chunkedMessageLock serializes them correctly. The lock scope is well chosen — decompression, newMessage and callback dispatch stay outside it, non-chunked messages never touch it — and I verified there are no lock-ordering cycles (nothing called under the lock — doAcknowledge, increaseAvailablePermits, trackMessage — re-enters chunk code) and that buffer refcounts stay balanced (uncompressPayloadIfNeeded does not consume its input). The PR also fixes several real pre-existing bugs beyond the headline race: the assemble→finalize NPE window, the duplicate-first-chunk overcount, the forward-gap discard desync, and — most impactful — the ghost-head bug where a completed uuid at the queue head permanently stalled expiry. All four new tests pass locally (:pulsar-client-original:test, race test ~4s).

One regression should be addressed before merge:

1. Decompression failure now permanently drops the earlier chunks' message IDs (final-chunk path in messageReceived)

The ctx is removed from chunkedMessagesMap and recycled inside the lock before decompression. If uncompressPayloadIfNeeded then fails, the code returns without acking or tracking the captured chunkedMessageIds; discardCorruptedMessage inside it only handles the final chunk's ID. Since the ctx is gone, the expiry sweep can never reach those IDs either — the first n−1 chunks stay unacked until the consumer reconnects (permanent ack hole / stuck backlog on that subscription). Previously the ctx stayed in the map and expiry eventually acked all recorded chunk IDs (before tripping the latent double-release this PR fixes — broken differently, but disposal did happen). Suggested fix: in the uncompressedPayload == null branch, individually doAcknowledge each non-null captured ID, mirroring the corrupted-chunk handling in doProcessMessageChunk / removeChunkMessage(..., autoAck=true) semantics.

Test / hygiene items:

2. Reflection into private state in the new tests. Project convention is no reflection into private state — use @VisibleForTesting package-private accessors instead. ConsumerImplTest is already in org.apache.pulsar.client.impl, so making processMessageChunk, pendingChunkedMessageCount, pendingChunkedMessageUuidQueue and expireChunkMessageTaskScheduled package-private removes all setAccessible/FieldUtils use. Some reflection is unnecessary even now: expireTimeOfIncompleteChunkedMessageMillis is protected and can be assigned directly, as the test already does with chunkedMessagesMap.

3. Stale javadoc on testDuplicateFirstChunkOvercountsPendingChunkedMessageCount: it still says "Currently FAILS deterministically … Enable once the duplicate-first-chunk path decrements the counter", which is pre-fix wording — the fix is included in this PR and the test passes. It also hardcodes source line references ("ConsumerImpl.java:1628-1634") that will rot. Please rewrite it to describe the guarded invariant.

4. The race test has no timeout. testChunkedMessageCountRaceBetweenReceiveAndExpiry uses unbounded receiver.join()/expirer.join() and no @Test(timeOut = …). The scenario it guards is exactly the kind whose future regression could be a deadlock between the two paths — which would hang the suite instead of failing. Please add a generous timeOut.

5. The PR description understates the change (in a good way): the ghost-head expiry-stall fix in doRemoveExpireIncompleteChunkedMessages (previously a completed uuid at the queue head blocked all expiry behind it indefinitely) and the forward-gap discard count/queue sync fix aren't mentioned, and "Verifying this change" lists 2 tests while 4 were added. Please update Motivation/Modifications so the full behavior change is captured for reviewers and release notes.

6. The PR title is truncated — it literally ends with a Unicode ellipsis ("…to fix use-after-…", GitHub's auto-fill from a long commit subject). Please complete it, e.g. [fix][client] Serialize chunked-message bookkeeping to fix use-after-free and count drift.

@SongOf SongOf changed the title [fix][client] Serialize chunked-message bookkeeping to fix use-after-… [fix][client] Serialize chunked-message bookkeeping to fix use-after-free and count/queue drift Jul 27, 2026
maxlisongsong added 2 commits July 28, 2026 02:34
# Conflicts:
#	pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java
@SongOf
SongOf force-pushed the fix/client-chunked-message-bookkeeping-race branch from d46ad7b to 3e4b95b Compare July 28, 2026 16:27
@SongOf

SongOf commented Jul 28, 2026

Copy link
Copy Markdown
Contributor Author

@lhotari
Fixed 1–4; Edited 5–6.


Thanks for the thorough review — fixed 1–4:

  1. Fixed: the ack loop in messageReceived now skips the last entry (chunkedMessageIds[0..length-2]), since it's already acked by discardCorruptedMessage inside uncompressPayloadIfNeeded's failure path. The earlier n−1 chunks
    are now acked individually instead of being dropped.
  2. Fixed: pendingChunkedMessageCount, pendingChunkedMessageUuidQueue, expireChunkMessageTaskScheduled, and processMessageChunk are now package-private + @VisibleForTesting; all setAccessible/FieldUtils reflection removed from
    the test file.
  3. Fixed: rewrote the javadoc on testDuplicateFirstChunkOvercountsPendingChunkedMessageCount to describe the guarded invariant instead of the pre-fix failure state; dropped the hardcoded line references.
  4. Fixed: added @test(timeOut = 60000) to testChunkedMessageCountRaceBetweenReceiveAndExpiry.

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

Labels

type/bug The PR fixed a bug or issue reported a bug

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants