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..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 @@ -439,6 +439,126 @@ 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(); + 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"); + } + + /** + * 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 2a81cf336944d..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 @@ -219,6 +219,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; @@ -228,6 +234,11 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle // 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; @@ -1628,25 +1639,43 @@ 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 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()); - } - pendingChunkedMessageCount++; - if (maxPendingChunkedMessage > 0 && pendingChunkedMessageCount > maxPendingChunkedMessage) { - removeOldestPendingChunkedMessage(); + 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. + 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 @@ -1692,6 +1721,14 @@ 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 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();