From b4c0204ac3b8cab01f351838e5a77e0b7795f148 Mon Sep 17 00:00:00 2001 From: Rahul Prasad Date: Sun, 20 Sep 2026 04:47:00 +0530 Subject: [PATCH 1/3] fix permit leak --- .../client/impl/MessageChunkingTest.java | 59 +++++++++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 14 +++++ 2 files changed, 73 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..45afc5c19a02e 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 @@ -365,6 +365,65 @@ 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. + * + * 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 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 and the consumer's available + * permits eventually drained to zero and dispatch stalled. + * + * Here we send N chunk-0's (all non-last, each credited at arrival) then expire them, then send + * N orphaned last chunks. Every chunk delivered must have its permit returned, so the client's + * availablePermits must equal the total number of chunks delivered (2 * N). + */ + @Test + public void testOrphanedLastChunkDoesNotLeakPermits() throws Exception { + final String topicName = "persistent://my-property/my-ns/orphanChunkPermitLeak"; + final String subName = "my-sub"; + @Cleanup + ConsumerImpl consumer = (ConsumerImpl) pulsarClient.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName(subName) + .maxPendingChunkedMessage(100) + .expireTimeOfIncompleteChunkedMessage(1, TimeUnit.SECONDS) + .autoAckOldestChunkedMessageOnQueueFull(true) + .subscribe(); + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName) + .chunkMaxMessageSize(100) + .enableChunking(true) + .enableBatching(false) + .create(); + + final int numMessages = 20; + + // Send only the first (non-last) chunk of each message; these are credited at arrival. + for (int i = 0; i < numMessages; i++) { + sendSingleChunk(producer, "orphan-" + i, 0, 2); + } + + // Let the scheduled expiry discard all the incomplete contexts. + Awaitility.await().atMost(15, TimeUnit.SECONDS) + .untilAsserted(() -> assertEquals(consumer.chunkedMessagesMap.size(), 0)); + + int permitsAfterExpiry = consumer.getAvailablePermits(); + + // Now send the orphaned LAST chunk of each message -> discard branch. Each must return its + // permit; without the fix none of these are credited. + for (int i = 0; i < numMessages; i++) { + sendSingleChunk(producer, "orphan-" + i, 1, 2); + } + + // Each orphaned last chunk consumed one permit that must be returned. + Awaitility.await().atMost(15, TimeUnit.SECONDS).untilAsserted(() -> + assertEquals(consumer.getAvailablePermits(), permitsAfterExpiry + numMessages, + "orphaned last chunks leaked flow-control permits")); + } + @Test public void testResendChunkMessagesWithoutAckHole() throws Exception { log.info().attr("method", methodName).log("Starting test"); 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..d56d14c7a8244 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 @@ -1669,6 +1669,12 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m if (isDuplicatedChunk) { doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); } + // A last chunk does not get its permit returned at the top of this method (only + // chunkId != last is credited there). Since this duplicated chunk is discarded here, + // return its permit to avoid leaking the broker's flow-control credit. + if (msgMetadata.getChunkId() == (msgMetadata.getNumChunksFromMsg() - 1)) { + increaseAvailablePermits(cnx); + } return null; } // means we lost the first chunk: should never happen @@ -1685,6 +1691,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); + } if (expireTimeOfIncompleteChunkedMessageMillis > 0 && System.currentTimeMillis() > (msgMetadata.getPublishTime() + expireTimeOfIncompleteChunkedMessageMillis)) { From 98628793f71c77b6dbfdd33afe01191d2b60ea0e Mon Sep 17 00:00:00 2001 From: Rahul Prasad Date: Wed, 23 Sep 2026 18:31:13 +0530 Subject: [PATCH 2/3] removed dead code --- .../java/org/apache/pulsar/client/impl/ConsumerImpl.java | 6 ------ 1 file changed, 6 deletions(-) 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 d86b82c61d8c0..6418e6e70a8f7 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 @@ -1679,12 +1679,6 @@ private ByteBuf processMessageChunk(ByteBuf compressedPayload, MessageMetadata m if (isDuplicatedChunk) { doAcknowledge(msgId, AckType.Individual, Collections.emptyMap(), null); } - // A last chunk does not get its permit returned at the top of this method (only - // chunkId != last is credited there). Since this duplicated chunk is discarded here, - // return its permit to avoid leaking the broker's flow-control credit. - if (msgMetadata.getChunkId() == (msgMetadata.getNumChunksFromMsg() - 1)) { - increaseAvailablePermits(cnx); - } return null; } // means we lost the first chunk: should never happen From 923f1f268ccd70294b47d7416c72c4c6bb40457e Mon Sep 17 00:00:00 2001 From: Rahul Prasad Date: Mon, 28 Sep 2026 01:03:33 +0530 Subject: [PATCH 3/3] updated testcase to simulate stall --- .../client/impl/MessageChunkingTest.java | 69 ++++++++++--------- 1 file changed, 38 insertions(+), 31 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 45afc5c19a02e..dcdf7e448ffb0 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 @@ -366,30 +366,34 @@ public void testMaxPendingChunkMessages() throws Exception { } /** - * Verifies that discarding an orphaned last chunk does not leak a flow-control permit. + * 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 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 and the consumer's available - * permits eventually drained to zero and dispatch stalled. + * 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. * - * Here we send N chunk-0's (all non-last, each credited at arrival) then expire them, then send - * N orphaned last chunks. Every chunk delivered must have its permit returned, so the client's - * availablePermits must equal the total number of chunks delivered (2 * N). + * 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 testOrphanedLastChunkDoesNotLeakPermits() throws Exception { + 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 - ConsumerImpl consumer = (ConsumerImpl) pulsarClient.newConsumer(Schema.STRING) + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) .topic(topicName) .subscriptionName(subName) - .maxPendingChunkedMessage(100) + .receiverQueueSize(receiverQueueSize) .expireTimeOfIncompleteChunkedMessage(1, TimeUnit.SECONDS) - .autoAckOldestChunkedMessageOnQueueFull(true) .subscribe(); @Cleanup Producer producer = pulsarClient.newProducer(Schema.STRING) @@ -399,29 +403,32 @@ public void testOrphanedLastChunkDoesNotLeakPermits() throws Exception { .enableBatching(false) .create(); - final int numMessages = 20; - - // Send only the first (non-last) chunk of each message; these are credited at arrival. - for (int i = 0; i < numMessages; i++) { + // 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) consumer).chunkedMessagesMap.size(), 0)); + // orphaned last chunk -> discard branch (must return its permit) + sendSingleChunk(producer, "orphan-" + idx, 1, 2); } - // Let the scheduled expiry discard all the incomplete contexts. - Awaitility.await().atMost(15, TimeUnit.SECONDS) - .untilAsserted(() -> assertEquals(consumer.chunkedMessagesMap.size(), 0)); - - int permitsAfterExpiry = consumer.getAvailablePermits(); - - // Now send the orphaned LAST chunk of each message -> discard branch. Each must return its - // permit; without the fix none of these are credited. - for (int i = 0; i < numMessages; i++) { - sendSingleChunk(producer, "orphan-" + i, 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); - // Each orphaned last chunk consumed one permit that must be returned. - Awaitility.await().atMost(15, TimeUnit.SECONDS).untilAsserted(() -> - assertEquals(consumer.getAvailablePermits(), permitsAfterExpiry + numMessages, - "orphaned last chunks leaked flow-control permits")); + Message 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