From c81ea5d134ec94704dd5b6aa3c833a45758e20d8 Mon Sep 17 00:00:00 2001 From: Rahul Prasad Date: Tue, 22 Sep 2026 19:05:14 +0530 Subject: [PATCH 1/3] counter drift --- .../client/impl/MessageChunkingTest.java | 53 +++++++++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 17 ++++++ 2 files changed, 70 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java index 9616614d421b6..cc0cb41c512c8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java @@ -439,6 +439,59 @@ public void testResendChunkMessages() throws Exception { Assert.assertEquals(((ConsumerImpl) consumer).getAvailablePermits(), 8); } + /** + * Proves that pendingChunkedMessageCount stays in sync with chunkedMessagesMap when duplicate + * (resent) first chunks arrive. + * + * When a first chunk (chunkId 0) arrives for a uuid that already has a ChunkedMessageCtx, + * processMessageChunk() removes the old context from chunkedMessagesMap without decrementing + * pendingChunkedMessageCount, then unconditionally increments it for the replacement. Before the + * fix, each duplicate first chunk inflated the count by 1 while the map size stayed 1 -- an + * upward drift that eventually pushes the count past maxPendingChunkedMessage and triggers + * spurious eviction of good in-flight messages. + * + * Invariant: pendingChunkedMessageCount == chunkedMessagesMap.size(). + * Buggy client: count == N, map == 1 -> FAIL. Fixed client: count == 1 -> PASS. + */ + @Test + public void testPendingChunkedMessageCountDriftOnDuplicateFirstChunk() throws Exception { + final String topicName = "persistent://my-property/my-ns/chunkCountDrift"; + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName("my-sub") + .maxPendingChunkedMessage(100) // high so eviction does not mask the drift + .autoAckOldestChunkedMessageOnQueueFull(true) + .subscribe(); + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName) + .chunkMaxMessageSize(100) + .enableChunking(true) + .enableBatching(false) + .create(); + + ConsumerImpl consumerImpl = (ConsumerImpl) consumer; + + // Send the first chunk (chunkId 0) of the SAME uuid many times (duplicate/resent first + // chunk), never completing it. Each duplicate replaces the context in the map (size stays 1) + // but bumps pendingChunkedMessageCount on the buggy client. + final int duplicates = 10; + for (int i = 0; i < duplicates; i++) { + sendSingleChunk(producer, "dup-uuid", 0, 2); + } + + Awaitility.await().atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertEquals(consumerImpl.chunkedMessagesMap.size(), 1)); + + int mapSize = consumerImpl.chunkedMessagesMap.size(); + int count = consumerImpl.getPendingChunkedMessageCountForTest(); + + assertEquals(count, mapSize, + "pendingChunkedMessageCount (" + count + ") drifted from chunkedMessagesMap.size (" + + mapSize + ") after " + duplicates + " duplicate first chunks"); + } + @Test public void testExpireIncompleteChunkMessage() throws Exception{ final String topicName = "persistent://my-property/my-ns/expireMsg"; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index 5862fce64de0e..cb7d45aa3f7ce 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -218,6 +218,12 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle protected Map chunkedMessagesMap = new ConcurrentHashMap<>(); private int pendingChunkedMessageCount = 0; + + @VisibleForTesting + int getPendingChunkedMessageCountForTest() { + return pendingChunkedMessageCount; + } + protected long expireTimeOfIncompleteChunkedMessageMillis = 0; private final AtomicBoolean expireChunkMessageTaskScheduled = new AtomicBoolean(false); private final int maxPendingChunkedMessage; @@ -1626,6 +1632,12 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m } chunkedMsgCtx.recycle(); chunkedMessagesMap.remove(msgMetadata.getUuid()); + // The replaced (old) context was counted when its first chunk arrived. It is being + // discarded here without completing, so decrement the count before the unconditional + // increment below re-counts the new replacement context. Otherwise each duplicate/ + // resent first chunk inflates pendingChunkedMessageCount by 1 while chunkedMessagesMap + // stays the same size -- drift that eventually triggers spurious eviction. + pendingChunkedMessageCount--; } pendingChunkedMessageCount++; if (maxPendingChunkedMessage > 0 && pendingChunkedMessageCount > maxPendingChunkedMessage) { @@ -1682,6 +1694,11 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); } chunkedMsgCtx.recycle(); + // A non-null context here means an out-of-order chunk is discarding a previously + // tracked (counted) assembly. It was counted when its first chunk arrived, so + // decrement the count to keep pendingChunkedMessageCount in sync with + // chunkedMessagesMap. (When chunkedMsgCtx is null nothing was counted for this uuid.) + pendingChunkedMessageCount--; } chunkedMessagesMap.remove(msgMetadata.getUuid()); compressedPayload.release(); From ed03dcb08f28d73ad330b80c3362ad106577080a Mon Sep 17 00:00:00 2001 From: Rahul Prasad Date: Tue, 22 Sep 2026 22:38:45 +0530 Subject: [PATCH 2/3] [fix][client] Keep chunk map, count and uuid-queue in sync on replace/discard On a duplicate/resent first chunk (or producer restart reusing the sequenceId), the consumer replaced the existing ChunkedMessageCtx but still incremented pendingChunkedMessageCount and re-added the uuid to pendingChunkedMessageUuidQueue, so both drifted above the real chunkedMessagesMap size on every duplicate. The count drift eventually crosses maxPendingChunkedMessage and evicts a good in-flight message; the queue growth is unbounded and a stale head can block removeExpireIncompleteChunkedMessages(). Handle replacement as reuse of the same tracking slot (no count++/queue.add), and only count+enqueue for a genuinely new uuid. Also decrement the count and remove the uuid from the queue in the out-of-order discard path when a tracked context is dropped. Adds a test asserting chunkedMessagesMap.size == pendingChunkedMessageCount == pendingChunkedMessageUuidQueue.size after repeated duplicate first chunks. Co-Authored-By: Claude Opus 4.8 --- .../client/impl/MessageChunkingTest.java | 6 +++ .../pulsar/client/impl/ConsumerImpl.java | 47 ++++++++++++------- 2 files changed, 36 insertions(+), 17 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java index cc0cb41c512c8..789c187c6490a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java @@ -486,7 +486,13 @@ public void testPendingChunkedMessageCountDriftOnDuplicateFirstChunk() throws Ex int mapSize = consumerImpl.chunkedMessagesMap.size(); int count = consumerImpl.getPendingChunkedMessageCountForTest(); + int queueSize = consumerImpl.getPendingChunkedMessageUuidQueueSizeForTest(); + // All three structures track the same single in-flight uuid; they must stay consistent + // regardless of how many duplicate first chunks arrive for it. + assertEquals(queueSize, mapSize, + "pendingChunkedMessageUuidQueue.size (" + queueSize + ") drifted from chunkedMessagesMap.size (" + + mapSize + ") after " + duplicates + " duplicate first chunks"); assertEquals(count, mapSize, "pendingChunkedMessageCount (" + count + ") drifted from chunkedMessagesMap.size (" + mapSize + ") after " + duplicates + " duplicate first chunks"); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index b09e01353d56e..1556bad580011 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -234,6 +234,11 @@ int getPendingChunkedMessageCountForTest() { // it will be used to manage N outstanding chunked message buffers private final BlockingQueue pendingChunkedMessageUuidQueue; + @VisibleForTesting + int getPendingChunkedMessageUuidQueueSizeForTest() { + return pendingChunkedMessageUuidQueue.size(); + } + private final boolean createTopicIfDoesNotExist; private final boolean poolMessages; @@ -1634,31 +1639,36 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m } }); } - // The first chunk of a new chunked-message received before receiving other chunks of previous - // chunked-message - // so, remove previous chunked-message from map and release buffer + // The first chunk of a new chunked-message received before receiving the other chunks + // of the previous chunked-message with the SAME uuid (a resent/duplicated first + // chunk, or a producer restart reusing the sequenceId). Discard the previous context. + // This is a REPLACEMENT of an existing tracking slot: the uuid is already represented + // exactly once in pendingChunkedMessageCount and once in pendingChunkedMessageUuidQueue, + // so this branch must NOT increment the count or re-add the uuid to the queue -- only + // the new-uuid branch below does that. Otherwise each duplicate first chunk would both + // drift the count above the real chunkedMessagesMap size (triggering spurious + // eviction) and leak a duplicate uuid into the queue (unbounded growth, and a stale + // head can block removeExpireIncompleteChunkedMessages()). The map entry itself is + // re-created by the shared computeIfAbsent below. if (chunkedMsgCtx.chunkedMsgBuffer != null) { ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); } chunkedMsgCtx.recycle(); chunkedMessagesMap.remove(msgMetadata.getUuid()); - // The replaced (old) context was counted when its first chunk arrived. It is being - // discarded here without completing, so decrement the count before the unconditional - // increment below re-counts the new replacement context. Otherwise each duplicate/ - // resent first chunk inflates pendingChunkedMessageCount by 1 while chunkedMessagesMap - // stays the same size -- drift that eventually triggers spurious eviction. - pendingChunkedMessageCount--; - } - pendingChunkedMessageCount++; - if (maxPendingChunkedMessage > 0 && pendingChunkedMessageCount > maxPendingChunkedMessage) { - removeOldestPendingChunkedMessage(); + } else { + // Genuinely new uuid: count it and enqueue it exactly once. Eviction is only checked + // here because only a new uuid grows the number of in-flight chunked messages. + pendingChunkedMessageCount++; + if (maxPendingChunkedMessage > 0 && pendingChunkedMessageCount > maxPendingChunkedMessage) { + removeOldestPendingChunkedMessage(); + } + pendingChunkedMessageUuidQueue.add(msgMetadata.getUuid()); } int totalChunks = msgMetadata.getNumChunksFromMsg(); ByteBuf chunkedMsgBuffer = PulsarByteBufAllocator.DEFAULT.buffer(msgMetadata.getTotalChunkMsgSize(), msgMetadata.getTotalChunkMsgSize()); chunkedMsgCtx = chunkedMessagesMap.computeIfAbsent(msgMetadata.getUuid(), (key) -> ChunkedMessageCtx.get(totalChunks, chunkedMsgBuffer)); - pendingChunkedMessageUuidQueue.add(msgMetadata.getUuid()); } // discard message if chunk is out-of-order @@ -1705,10 +1715,13 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m } chunkedMsgCtx.recycle(); // A non-null context here means an out-of-order chunk is discarding a previously - // tracked (counted) assembly. It was counted when its first chunk arrived, so - // decrement the count to keep pendingChunkedMessageCount in sync with - // chunkedMessagesMap. (When chunkedMsgCtx is null nothing was counted for this uuid.) + // tracked assembly. That assembly's uuid was counted AND added to + // pendingChunkedMessageUuidQueue when its first chunk arrived, so remove it from both + // to keep chunkedMessagesMap, pendingChunkedMessageCount and the uuid queue in sync. + // (When chunkedMsgCtx is null nothing was ever tracked for this uuid, so there is + // nothing to decrement or dequeue.) pendingChunkedMessageCount--; + pendingChunkedMessageUuidQueue.remove(msgMetadata.getUuid()); } chunkedMessagesMap.remove(msgMetadata.getUuid()); compressedPayload.release(); From 24d2088ae03bed9ed68b7ef3bf34f99801fb6951 Mon Sep 17 00:00:00 2001 From: Rahul Prasad Date: Thu, 1 Oct 2026 22:16:07 +0530 Subject: [PATCH 3/3] resend duplicate chunk stops expiry of other messages --- .../client/impl/MessageChunkingTest.java | 61 +++++++++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 23 ++++--- 2 files changed, 76 insertions(+), 8 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java index 789c187c6490a..7226c66be57cf 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java @@ -498,6 +498,67 @@ public void testPendingChunkedMessageCountDriftOnDuplicateFirstChunk() throws Ex + mapSize + ") after " + duplicates + " duplicate first chunks"); } + /** + * A duplicate first chunk resets its context's receivedTime. If the uuid is left at its original + * position in pendingChunkedMessageUuidQueue, queue order no longer matches expiry order, and + * because removeExpireIncompleteChunkedMessages() only inspects the head and returns at the first + * non-expired entry, a repeatedly-refreshed head uuid blocks expiry of genuinely-expired entries + * behind it (the reviewer's A/B head-of-line scenario). + * + * Here A's first chunk is refreshed continuously so A never expires, while B is left untouched + * past the expiry window. B must still be expired and removed. The fix re-positions A's uuid in + * the queue on each refresh so B reaches the head and is collected. + */ + @Test + public void testRefreshedFirstChunkDoesNotBlockExpiryOfLaterEntries() throws Exception { + final String topicName = "persistent://my-property/my-ns/refreshedHeadExpiry"; + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName("my-sub") + .maxPendingChunkedMessage(100) + .expireTimeOfIncompleteChunkedMessage(2, TimeUnit.SECONDS) + .subscribe(); + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName) + .chunkMaxMessageSize(100) + .enableChunking(true) + .enableBatching(false) + .create(); + + ConsumerImpl consumerImpl = (ConsumerImpl) consumer; + + // A arrives first, then B. Both incomplete. Queue: [A, B]. + sendSingleChunk(producer, "A", 0, 2); + sendSingleChunk(producer, "B", 0, 2); + Awaitility.await().atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertEquals(consumerImpl.chunkedMessagesMap.size(), 2)); + + // Keep refreshing A's first chunk continuously so A's receivedTime is never older than the + // 2s expiry window -- A must never be eligible for expiry while we observe. We check B's + // removal WHILE still refreshing A: where the + // perpetually-refreshed head uuid actively shields the entry behind it. B is never refreshed, + // so once >2s pass it is expired -- unless A (left at the stale queue head without the fix) + // blocks the head-only expiry scan. Refresh every 400ms for up to ~12s, asserting B is gone. + long bArrivalNanos = System.nanoTime(); + boolean bCollected = false; + for (int i = 0; i < 30 && !bCollected; i++) { + sendSingleChunk(producer, "A", 0, 2); // keep A's receivedTime fresh (head stays non-expired) + Thread.sleep(400); + // only meaningful once B is well past its 2s expiry window + if (System.nanoTime() - bArrivalNanos > TimeUnit.SECONDS.toNanos(4)) { + bCollected = consumerImpl.chunkedMessagesMap.get("B") == null; + } + } + + // B must have been expired and removed while A was still being refreshed at the (old) queue + // head. On the buggy client B stays stuck behind the refreshed A indefinitely. + assertNull(consumerImpl.chunkedMessagesMap.get("B"), + "expired message B was not collected while A was continually refreshed: a refreshed " + + "head uuid blocked head-only expiry of later entries"); + } + @Test public void testExpireIncompleteChunkMessage() throws Exception{ final String topicName = "persistent://my-property/my-ns/expireMsg"; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index 1556bad580011..9fe57cf0b07fc 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -1642,19 +1642,26 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m // The first chunk of a new chunked-message received before receiving the other chunks // of the previous chunked-message with the SAME uuid (a resent/duplicated first // chunk, or a producer restart reusing the sequenceId). Discard the previous context. - // This is a REPLACEMENT of an existing tracking slot: the uuid is already represented - // exactly once in pendingChunkedMessageCount and once in pendingChunkedMessageUuidQueue, - // so this branch must NOT increment the count or re-add the uuid to the queue -- only - // the new-uuid branch below does that. Otherwise each duplicate first chunk would both - // drift the count above the real chunkedMessagesMap size (triggering spurious - // eviction) and leak a duplicate uuid into the queue (unbounded growth, and a stale - // head can block removeExpireIncompleteChunkedMessages()). The map entry itself is - // re-created by the shared computeIfAbsent below. + // This is a REPLACEMENT of an existing tracking slot: the uuid is already counted once + // in pendingChunkedMessageCount, so this branch must NOT increment the count -- only + // the new-uuid branch below does. Otherwise each duplicate first chunk would drift the + // count above the real chunkedMessagesMap size (triggering spurious eviction). + // + // The uuid must, however, be re-positioned in pendingChunkedMessageUuidQueue. The + // replacement context below gets a fresh receivedTime, so leaving the uuid at its + // original (older) position would make queue order no longer match expiry order: + // removeExpireIncompleteChunkedMessages() only inspects the head and returns at the + // first non-expired entry, so a repeatedly-refreshed head uuid would indefinitely + // block expiry of genuinely-expired entries behind it. Remove the stale entry and + // re-add it so its position reflects the refreshed receivedTime. The count is + // unchanged (one entry out, one back in). if (chunkedMsgCtx.chunkedMsgBuffer != null) { ReferenceCountUtil.safeRelease(chunkedMsgCtx.chunkedMsgBuffer); } chunkedMsgCtx.recycle(); chunkedMessagesMap.remove(msgMetadata.getUuid()); + pendingChunkedMessageUuidQueue.remove(msgMetadata.getUuid()); + pendingChunkedMessageUuidQueue.add(msgMetadata.getUuid()); } else { // Genuinely new uuid: count it and enqueue it exactly once. Eviction is only checked // here because only a new uuid grows the number of in-flight chunked messages.