Skip to content
Open
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 @@ -1185,4 +1185,85 @@ public void testExceptionBySeekFunction() throws Exception {
assertTrue(e.getCause().getMessage().contains("Only support seek by messageId or timestamp"));
}
}

/**
* Seeking exclusively to a messageId re-dispatches that boundary message, which the consumer
* filters out instead of delivering. The dropped message's flow-control permit must still be
* returned; otherwise, with receiverQueueSize=1, the single leaked permit exhausts the budget and
* the consumer stalls right after the seek, never delivering the message after the seek target.
*/
@Test
public void testSeekBoundaryDropDoesNotLeakPermit() throws Exception {
final String topicName = "persistent://prop/ns-abc/seekBoundaryPermitLeak";

@Cleanup
Producer<byte[]> producer = pulsarClient.newProducer().topic(topicName)
.enableBatching(false).create();

List<MessageId> ids = new ArrayList<>();
for (int i = 0; i < 5; i++) {
ids.add(producer.send(("seek-msg-" + i).getBytes()));
}

@Cleanup
org.apache.pulsar.client.api.Consumer<byte[]> consumer = pulsarClient.newConsumer()
.topic(topicName)
.subscriptionName("my-sub")
.receiverQueueSize(1) // tiny budget: a single leaked permit stalls the consumer
.subscribe();

// Seek exclusively to index 1. The boundary message (index 1) is filtered; the messages
// after it (index 2, 3, 4) must still be deliverable.
consumer.seek(ids.get(1));

Message<byte[]> msg = consumer.receive(10, TimeUnit.SECONDS);
assertNotNull(msg, "consumer stalled after seek: the boundary-message drop leaked its permit "
+ "and the receiverQueueSize=1 budget was exhausted");
assertEquals(msg.getValue(), "seek-msg-2".getBytes());
consumer.acknowledge(msg);
}

/**
* The chunked variant of {@link #testSeekBoundaryDropDoesNotLeakPermit}: the boundary message
* dropped on seek is a chunked message. The ConsumerImpl drop block credits the non-last chunks
* at arrival and must repay the last chunk's permit when the assembled message is filtered;
* otherwise, with receiverQueueSize=1, the consumer stalls and the next chunked message is never
* delivered.
*/
@Test
public void testSeekBoundaryDropDoesNotLeakPermitForChunkedMessage() throws Exception {
final String topicName = "persistent://prop/ns-abc/seekBoundaryPermitLeakChunked";

@Cleanup
Producer<byte[]> producer = pulsarClient.newProducer().topic(topicName)
.enableBatching(false)
.enableChunking(true)
.chunkMaxMessageSize(100) // force multi-chunk from a modest payload
.create();

// Each message is ~3 chunks at chunkMaxMessageSize=100.
byte[] payload = new byte[250];
Arrays.fill(payload, (byte) 'x');
List<MessageId> ids = new ArrayList<>();
for (int i = 0; i < 3; i++) {
ids.add(producer.send(payload));
}

@Cleanup
org.apache.pulsar.client.api.Consumer<byte[]> consumer = pulsarClient.newConsumer()
.topic(topicName)
.subscriptionName("my-sub")
.receiverQueueSize(1)
.subscribe();

// Seek exclusively to the first chunked message; it is filtered as the boundary. The chunked
// message after it must still be delivered.
consumer.seek(ids.get(0));

Message<byte[]> msg = consumer.receive(15, TimeUnit.SECONDS);
assertNotNull(msg, "consumer stalled after seek: dropping the chunked boundary message leaked "
+ "a flow-control permit and the receiverQueueSize=1 budget was exhausted");
assertEquals(msg.getMessageId(), ids.get(1));
consumer.acknowledge(msg);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -1545,6 +1545,13 @@ void messageReceived(CommandMessage cmdMessage, ByteBuf headersAndPayload, Clien
.log("Ignoring message from before the startMessageId");

uncompressedPayload.release();
// This message is dropped instead of being delivered to the application, so its
// outstanding flow-control permit is never returned via messageProcessed(). Return
// it here to avoid leaking a permit for the boundary message that a seek/startMessageId
// caused to be re-dispatched. Refund numMessages : the broker charges
// one permit per message in the entry, so an undecryptable batch (which can also reach
// this block) is repaid exactly what it consumed. For the plain and chunked cases numMessages is 1.
increaseAvailablePermits(cnx, numMessages);
return;
}

Expand Down