Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -439,6 +439,126 @@ public void testResendChunkMessages() throws Exception {
Assert.assertEquals(((ConsumerImpl<String>) 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<String> 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<String> producer = pulsarClient.newProducer(Schema.STRING)
.topic(topicName)
.chunkMaxMessageSize(100)
.enableChunking(true)
.enableBatching(false)
.create();

ConsumerImpl<String> consumerImpl = (ConsumerImpl<String>) 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));

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.

[SUGGESTION] this test can assert before the duplicate chunks are processed, and the out-of-order path is not covered

producer.send() returning does not mean the consumer has processed the chunk, and chunkedMessagesMap.size() == 1 is already true after the first of the ten duplicates, so the assertions below can run before the rest have arrived (and between a replacement's remove and computeIfAbsent the map is briefly empty). It failed reliably for me with the fix reverted, so it is not vacuous in practice, but it is timing-dependent. A deterministic way is to send a complete second chunked message afterwards, receive it (messages are processed in order, so all duplicates are done), and then assert the three sizes.

It would also be good to cover the out-of-order discard (count-- plus queue.remove), and eviction order after a replacement.


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<String> consumer = pulsarClient.newConsumer(Schema.STRING)
.topic(topicName)
.subscriptionName("my-sub")
.maxPendingChunkedMessage(100)
.expireTimeOfIncompleteChunkedMessage(2, TimeUnit.SECONDS)
.subscribe();
@Cleanup
Producer<String> producer = pulsarClient.newProducer(Schema.STRING)
.topic(topicName)
.chunkMaxMessageSize(100)
.enableChunking(true)
.enableBatching(false)
.create();

ConsumerImpl<String> consumerImpl = (ConsumerImpl<String>) 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

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.

[NIT] tidy the new expiry test's comments

This comment has a fragment (WHILE still refreshing A: where the), and the Javadoc refers to "the reviewer's" scenario. Please reword both so the test reads standalone. The loop also sleeps up to ~12s; it passes quickly in the normal case, but it is worth checking it is not close to the 10s-style Awaitility limits used elsewhere in this class if CI is slow.

// 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";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,12 @@ public class ConsumerImpl<T> extends ConsumerBase<T> implements ConnectionHandle

protected Map<String, ChunkedMessageCtx> 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;
Expand All @@ -228,6 +234,11 @@ public class ConsumerImpl<T> extends ConsumerBase<T> implements ConnectionHandle
// it will be used to manage N outstanding chunked message buffers
private final BlockingQueue<String> pendingChunkedMessageUuidQueue;

@VisibleForTesting
int getPendingChunkedMessageUuidQueueSizeForTest() {
return pendingChunkedMessageUuidQueue.size();
}

private final boolean createTopicIfDoesNotExist;
private final boolean poolMessages;

Expand Down Expand Up @@ -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++;

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.

[SUGGESTION] consider deriving the pending count from chunkedMessagesMap instead of a separate counter

The bug here is that a hand-maintained int has to be adjusted on every path that touches the map (completion, replace, out-of-order discard, eviction, expiry), and this PR adds two more adjustments. Since chunkedMessagesMap is the source of truth, the eviction check could be chunkedMessagesMap.size() >= maxPendingChunkedMessage before inserting a new uuid, and pendingChunkedMessageCount (and its ForTest getter) could go away, which removes this whole class of drift. It would also avoid the plain int being modified from both the Netty thread and the expiry task on internalPinnedExecutor.

If you keep the counter, that is fine too; it just means every future change to these paths has to keep it in step by hand.

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
Expand Down Expand Up @@ -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());

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.

[NIT] queue.remove(uuid) is a linear scan on a locked array queue

GrowableArrayBlockingQueue.remove(Object) scans and shifts the backing array while holding its locks, and this runs on the Netty thread. With the default maxPendingChunkedMessage that is small, but with 0 (unbounded) a stream of out-of-order chunks against a large queue gets more expensive. The replace path now does the same remove-then-add. Not a blocker; an insertion-ordered map keyed by uuid would make these O(1).

}
chunkedMessagesMap.remove(msgMetadata.getUuid());
compressedPayload.release();
Expand Down