Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -365,6 +365,72 @@ public void testMaxPendingChunkMessages() throws Exception {
assertNull(consumer.receive(5, TimeUnit.SECONDS));
}

/**
* Verifies that discarding an orphaned last chunk does not leak a flow-control permit, by
* asserting the user-visible symptom: dispatch must not stall.
*
* Each chunk the broker dispatches consumes one permit. Non-last chunks are credited back at
* arrival; the last chunk is normally credited when the assembled message is consumed. When a
* chunked message is torn apart (its first chunk expires, then the orphaned last chunk arrives
* with no assembly context), the last chunk hits the discard branch. Without the fix that branch
* never returned the permit, so every torn message leaked one permit.
*
* The consumer only sends fresh permits to the broker once its returned-permit accumulator
* reaches receiverQueueSize/2. With a small receiver queue, leaking more than receiverQueueSize/2
* permits means that threshold is never reached again, the broker's credit drains to zero, and
* dispatch stalls permanently. We reproduce exactly that: tear more than receiverQueueSize/2
* messages, then send a normal complete chunked message and assert it is still delivered.
*
* Without the fix the final message never arrives (stall). With the fix it is received.
*/
@Test
public void testOrphanedLastChunkDoesNotStallDispatch() throws Exception {
final String topicName = "persistent://my-property/my-ns/orphanChunkPermitLeak";
final String subName = "my-sub";
final int receiverQueueSize = 4; // flush threshold = receiverQueueSize/2 = 2
@Cleanup
Consumer<String> consumer = pulsarClient.newConsumer(Schema.STRING)
.topic(topicName)
.subscriptionName(subName)
.receiverQueueSize(receiverQueueSize)
.expireTimeOfIncompleteChunkedMessage(1, TimeUnit.SECONDS)
.subscribe();
@Cleanup
Producer<String> producer = pulsarClient.newProducer(Schema.STRING)
.topic(topicName)
.chunkMaxMessageSize(100)
.enableChunking(true)
.enableBatching(false)
.create();

// Tear apart well more than receiverQueueSize/2 chunked messages. Each torn message leaks
// one permit on the buggy client; once cumulative leaks exceed receiverQueueSize the broker
// credit is exhausted and never replenished.
final int tornMessages = receiverQueueSize * 3; // 12
for (int i = 0; i < tornMessages; i++) {
// first (non-last) chunk -> creates an incomplete context
sendSingleChunk(producer, "orphan-" + i, 0, 2);
// wait for the scheduled expiry to discard it
final int idx = i;
Awaitility.await().atMost(10, TimeUnit.SECONDS)
.untilAsserted(() -> assertEquals(
((ConsumerImpl<String>) consumer).chunkedMessagesMap.size(), 0));
Comment thread
programmerahul marked this conversation as resolved.
// orphaned last chunk -> discard branch (must return its permit)
sendSingleChunk(producer, "orphan-" + idx, 1, 2);
}

// Now send a NORMAL complete chunked message. On a healthy consumer it is dispatched and
// received; on the buggy client the leaked permits have stalled dispatch and it never arrives.
sendSingleChunk(producer, "live", 0, 2);
sendSingleChunk(producer, "live", 1, 2);

Message<String> msg = consumer.receive(15, TimeUnit.SECONDS);
assertNotNull(msg, "dispatch stalled: orphaned last chunks leaked flow-control permits until "
+ "the broker stopped dispatching (receiverQueueSize=" + receiverQueueSize + ")");
assertEquals(msg.getValue(), "chunk-live-0|chunk-live-1|");
consumer.acknowledge(msg);
}

@Test
public void testResendChunkMessagesWithoutAckHole() throws Exception {
log.info().attr("method", methodName).log("Starting test");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1696,6 +1696,14 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m
}
chunkedMessagesMap.remove(msgMetadata.getUuid());
compressedPayload.release();
// This discarded chunk consumed a broker flow-control permit. Non-last chunks already
// had their permit returned at the top of this method (increaseAvailablePermits when
// chunkId != last); the last chunk did not. Return it here so that tearing a chunked
// message apart (expiry/eviction/orphaned last chunk) does not leak permits, which would
// otherwise drain the consumer's available permits to zero and stall dispatch.
if (msgMetadata.getChunkId() == (msgMetadata.getNumChunksFromMsg() - 1)) {
increaseAvailablePermits(cnx);
}
Comment thread
programmerahul marked this conversation as resolved.
if (expireTimeOfIncompleteChunkedMessageMillis > 0
&& System.currentTimeMillis() > (msgMetadata.getPublishTime()
+ expireTimeOfIncompleteChunkedMessageMillis)) {
Expand Down
Loading