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..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 @@ -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 consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName(subName) + .receiverQueueSize(receiverQueueSize) + .expireTimeOfIncompleteChunkedMessage(1, TimeUnit.SECONDS) + .subscribe(); + @Cleanup + Producer 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) consumer).chunkedMessagesMap.size(), 0)); + // 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 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"); 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 fc3ce28b41216..2a8077f713a9e 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 @@ -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); + } if (expireTimeOfIncompleteChunkedMessageMillis > 0 && System.currentTimeMillis() > (msgMetadata.getPublishTime() + expireTimeOfIncompleteChunkedMessageMillis)) {