From dc7dcf367733a71a2a1f721237648775b9945ea1 Mon Sep 17 00:00:00 2001 From: Denovo1998 Date: Sat, 19 Sep 2026 18:33:46 +0800 Subject: [PATCH 1/2] [fix][broker] Prevent stale consumer removal from removing replacements --- .../AbstractDispatcherMultipleConsumers.java | 15 +++++ .../service/ConsumerPrioritySelector.java | 9 +++ ...PersistentDispatcherMultipleConsumers.java | 10 +-- ...entDispatcherMultipleConsumersClassic.java | 6 +- ...tStickyKeyDispatcherMultipleConsumers.java | 5 ++ ...KeyDispatcherMultipleConsumersClassic.java | 9 ++- .../service/ConsumerPrioritySelectorTest.java | 16 +++++ ...criptionUnackedMessagesAccountingTest.java | 65 +++++++++++++++++++ 8 files changed, 125 insertions(+), 10 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java index c1358d75aeeb3..ab1854d531de5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java @@ -60,6 +60,10 @@ protected final synchronized void removeConsumerFromList(Consumer consumer) { consumerPrioritySelector.remove(consumer); } + protected final synchronized void removeConsumerInstanceFromList(Consumer consumer) { + consumerPrioritySelector.removeInstance(consumer); + } + protected final synchronized void removeConsumersFromList(Predicate predicate) { consumerPrioritySelector.removeIf(predicate); } @@ -88,6 +92,17 @@ protected final boolean containsConsumerInstance(Consumer consumer) { return consumerSetImpl.indexExists(index) && consumerSetImpl.indexGet(index) == consumer; } + /** + * Removes only the exact registered instance. The caller must hold the dispatcher monitor. + */ + protected final boolean removeConsumerInstance(Consumer consumer) { + if (!containsConsumerInstance(consumer)) { + return false; + } + consumerSet.removeAll(consumer); + return true; + } + public boolean isClosed() { return isClosed == TRUE; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ConsumerPrioritySelector.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ConsumerPrioritySelector.java index a69ddd5fd11fa..1a11527360d09 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ConsumerPrioritySelector.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ConsumerPrioritySelector.java @@ -66,6 +66,15 @@ void remove(T consumer) { } } + void removeInstance(T consumer) { + for (int index = 0; index < consumerList.size(); index++) { + if (consumerList.get(index) == consumer) { + decrementPriorityCount(consumerList.remove(index)); + return; + } + } + } + void removeIf(Predicate predicate) { consumerList.removeIf(predicate); // Bulk removal repairs inconsistent dispatcher membership. Rebuild from survivors so even an diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index dffc09c3e18a6..079e612ed1e98 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -241,12 +241,12 @@ protected boolean isConsumersExceededOnSubscription() { @Override public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { - if (consumerSet.removeAll(consumer) == 1) { + if (removeConsumerInstance(consumer)) { // decrement unack-message count for removed consumer. Only the removal that actually // unregisters the consumer may debit it, otherwise removing an already-removed consumer // debits the same messages again and drives the subscription counter negative. addUnAckedMessages(-consumer.getUnackedMessages()); - removeConsumerFromList(consumer); + removeConsumerInstanceFromList(consumer); log.info() .attr("consumer", consumer) .attr("pendingAcks", consumer.getPendingAcks().size()) @@ -284,7 +284,7 @@ public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceE // The debit belongs to the removal that unregisters the consumer; do not repeat it here. // The add-consumer failure path can also unregister via internalRemoveConsumer, but that // consumer has not received any messages and therefore has nothing to debit. - removeConsumersFromList(c -> consumer.equals(c)); + removeConsumersFromList(c -> c == consumer); if (consumerList.isEmpty()) { clearComponentsAfterRemovedAllConsumers(); } @@ -292,8 +292,8 @@ public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceE } protected synchronized void internalRemoveConsumer(Consumer consumer) { - consumerSet.removeAll(consumer); - removeConsumerFromList(consumer); + removeConsumerInstance(consumer); + removeConsumerInstanceFromList(consumer); } protected synchronized void clearComponentsAfterRemovedAllConsumers() { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java index 1f5e70a741bc9..84a94efc81117 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java @@ -231,12 +231,12 @@ protected boolean isConsumersExceededOnSubscription() { @Override public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { - if (consumerSet.removeAll(consumer) == 1) { + if (removeConsumerInstance(consumer)) { // decrement unack-message count for removed consumer. Only the removal that actually // unregisters the consumer may debit it, otherwise removing an already-removed consumer // debits the same messages again and drives the subscription counter negative. addUnAckedMessages(-consumer.getUnackedMessages()); - removeConsumerFromList(consumer); + removeConsumerInstanceFromList(consumer); log.info() .attr("consumer", consumer) .attr("size", consumer.getPendingAcks().size()) @@ -267,7 +267,7 @@ public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceE */ log.error().attr("consumer", consumer).log("Trying to remove a non-connected consumer"); // The debit belongs to the removal that unregisters the consumer; do not repeat it here. - removeConsumersFromList(c -> consumer.equals(c)); + removeConsumersFromList(c -> c == consumer); if (consumerList.isEmpty()) { clearComponentsAfterRemovedAllConsumers(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java index b8acced90c9ab..4cd4fc6c07a33 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java @@ -226,6 +226,11 @@ private synchronized void registerDrainingHashes(Consumer skipConsumer, @Override public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { + if (!containsConsumerInstance(consumer)) { + // Let the superclass repair stale list membership without touching a replacement's selector state. + super.removeConsumer(consumer); + return; + } // The consumer must be removed from the selector before calling the superclass removeConsumer method. Optional impactedConsumers = selector.removeConsumer(consumer); super.removeConsumer(consumer); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersClassic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersClassic.java index 2235611778647..41d4f8e54b057 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersClassic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersClassic.java @@ -159,8 +159,8 @@ public synchronized CompletableFuture addConsumer(Consumer consumer) { selector.addConsumer(consumer).handle((result, ex) -> { if (ex != null) { synchronized (PersistentStickyKeyDispatcherMultipleConsumersClassic.this) { - consumerSet.removeAll(consumer); - removeConsumerFromList(consumer); + removeConsumerInstance(consumer); + removeConsumerInstanceFromList(consumer); } throw FutureUtil.wrapToCompletionException(ex); } @@ -234,6 +234,11 @@ private void sortRecentlyJoinedConsumersIfNeeded() { @Override public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { + if (!containsConsumerInstance(consumer)) { + // Let the superclass repair stale list membership without touching a replacement's selector state. + super.removeConsumer(consumer); + return; + } // The consumer must be removed from the selector before calling the superclass removeConsumer method. // In the superclass removeConsumer method, the pending acks that the consumer has are added to // redeliveryMessages. If the consumer has not been removed from the selector at this point, diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsumerPrioritySelectorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsumerPrioritySelectorTest.java index bbe3f91c56fe4..92abcc24b5fe7 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsumerPrioritySelectorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsumerPrioritySelectorTest.java @@ -176,6 +176,22 @@ public void priorityCountsFollowActualRemovals() { assertThat(selector.priorityLevelCount()).isZero(); } + @Test + public void instanceRemovalPreservesEqualReplacement() { + List consumers = new CopyOnWriteArrayList<>(); + ConsumerPrioritySelector selector = new ConsumerPrioritySelector<>(consumers, c -> c.priority, + c -> c.available); + Slot original = new Slot(0, 0); + Slot replacement = new Slot(0, 3); + selector.add(replacement); + selector.removeInstance(original); + assertThat(consumers).singleElement().isSameAs(replacement); + assertThat(selector.priorityLevelCount()).isEqualTo(1); + selector.removeInstance(replacement); + assertThat(consumers).isEmpty(); + assertThat(selector.priorityLevelCount()).isZero(); + } + @Test public void bulkRemovalRepairsInconsistentMembership() { List consumers = new CopyOnWriteArrayList<>(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedSubscriptionUnackedMessagesAccountingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedSubscriptionUnackedMessagesAccountingTest.java index ad9135b0df190..cafad7772e7d2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedSubscriptionUnackedMessagesAccountingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedSubscriptionUnackedMessagesAccountingTest.java @@ -19,17 +19,24 @@ package org.apache.pulsar.broker.service; import static org.assertj.core.api.Assertions.assertThat; +import java.util.Collections; +import java.util.Optional; import java.util.concurrent.TimeUnit; import org.apache.pulsar.broker.service.persistent.AbstractPersistentDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumersClassic; import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerConsumerBase; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; +import org.apache.pulsar.common.api.proto.CommandSubscribe.InitialPosition; +import org.apache.pulsar.common.api.proto.KeySharedMeta; +import org.apache.pulsar.common.api.proto.KeySharedMode; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; +import org.testng.annotations.DataProvider; import org.testng.annotations.Factory; import org.testng.annotations.Test; @@ -57,6 +64,7 @@ public SharedSubscriptionUnackedMessagesAccountingTest(boolean classic) { protected void doInitConf() throws Exception { super.doInitConf(); conf.setSubscriptionSharedUseClassicPersistentImplementation(classic); + conf.setSubscriptionKeySharedUseClassicPersistentImplementation(classic); conf.setMaxUnackedMessagesPerBroker(1000); } @@ -121,6 +129,63 @@ public void testRemovingSameConsumerTwiceDebitsUnackedMessagesOnce() throws Exce } } + @DataProvider + public Object[][] subscriptionTypes() { + return new Object[][] {{SubscriptionType.Shared}, {SubscriptionType.Key_Shared}}; + } + + @Test(dataProvider = "subscriptionTypes", timeOut = 60_000) + public void testStaleRemovalDoesNotRemoveEqualReplacement(SubscriptionType subscriptionType) throws Exception { + String topicName = newTopicName(); + try (Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName).enableBatching(false).create(); + org.apache.pulsar.client.api.Consumer client = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName).subscriptionName(SUBSCRIPTION).subscriptionType(subscriptionType) + .consumerName("replacement-test").subscribe()) { + producer.send("unacked"); + assertThat(client.receive(5, TimeUnit.SECONDS)).isNotNull(); + BrokerService brokerService = pulsar.getBrokerService(); + PersistentTopic topic = (PersistentTopic) brokerService.getTopicReference(topicName).orElseThrow(); + AbstractPersistentDispatcherMultipleConsumers dispatcher = + (AbstractPersistentDispatcherMultipleConsumers) topic.getSubscription(SUBSCRIPTION).getDispatcher(); + Consumer original = dispatcher.getConsumers().get(0); + synchronized (dispatcher) { + assertThat(original.getUnackedMessages()).isEqualTo(1); + dispatcher.removeConsumer(original); + assertUnackedMessagesCleared(dispatcher, brokerService, "original removal"); + } + // Use the production subscribe path and the same live connection/protocol identity. No mocked + // consumer or manually populated dispatcher collections are needed to create the replacement. + Consumer replacement = topic.subscribe(SubscriptionOption.builder() + .cnx(original.cnx()).consumerId(original.consumerId()).consumerName(original.consumerName()) + .subscriptionName(SUBSCRIPTION).subType(original.subType()).isDurable(true) + .startMessageId(MessageId.latest).initialPosition(InitialPosition.Latest) + .metadata(Collections.emptyMap()).subscriptionProperties(Optional.empty()) + .keySharedMeta(new KeySharedMeta().setKeySharedMode(KeySharedMode.AUTO_SPLIT)) + .build()).get(10, TimeUnit.SECONDS); + try { + assertThat(replacement).isNotSameAs(original).isEqualTo(original); + assertThat(replacement.hashCode()).isEqualTo(original.hashCode()); + synchronized (dispatcher) { + dispatcher.removeConsumer(original); + assertThat(dispatcher.getConsumers()).singleElement().isSameAs(replacement); + assertUnackedMessagesCleared(dispatcher, brokerService, "stale removal"); + } + // This also checks Key_Shared selector membership, not just the consumer list. + replacement.flowPermits(1); + assertThat(client.receive(5, TimeUnit.SECONDS)).as("replacement receives replay").isNotNull(); + synchronized (dispatcher) { + assertThat(replacement.getUnackedMessages()).isEqualTo(1); + assertThat(dispatcher.getTotalUnackedMessages()).isEqualTo(1); + assertThat(brokerService.getTotalUnackedMessages()).isEqualTo(1); + } + } finally { + replacement.close(); + } + assertUnackedMessagesCleared(dispatcher, brokerService, "replacement removal"); + } + } + private void assertUnackedMessagesCleared(AbstractPersistentDispatcherMultipleConsumers dispatcher, BrokerService brokerService, String removal) { assertThat(dispatcher.getTotalUnackedMessages()).as("subscription balance after %s", removal).isZero(); From cc0311f92a2bfe3eaefc12c4921a4292ee583431 Mon Sep 17 00:00:00 2001 From: Denovo1998 Date: Sun, 20 Sep 2026 21:37:48 +0800 Subject: [PATCH 2/2] Ensure sticky dispatcher cleanup handles unregistered consumers and selector removal uses identity. --- ...ngeExclusiveStickyKeyConsumerSelector.java | 2 +- ...tStickyKeyDispatcherMultipleConsumers.java | 4 +- ...KeyDispatcherMultipleConsumersClassic.java | 6 +- ...criptionUnackedMessagesAccountingTest.java | 95 +++++++++++++++++-- 4 files changed, 95 insertions(+), 12 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelector.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelector.java index 7fb9983819725..7ed3db4007d38 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelector.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelector.java @@ -81,7 +81,7 @@ private synchronized Optional internalAddConsumer(Consu @Override public synchronized Optional removeConsumer(Consumer consumer) { - rangeMap.entrySet().removeIf(entry -> entry.getValue().getRight().equals(consumer)); + rangeMap.entrySet().removeIf(entry -> entry.getValue().getRight() == consumer); return Optional.empty(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java index 4cd4fc6c07a33..cd33dd9cc3520 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java @@ -226,7 +226,9 @@ private synchronized void registerDrainingHashes(Consumer skipConsumer, @Override public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { - if (!containsConsumerInstance(consumer)) { + // A STICKY addition can finish its liveness check after dispatcher removal. Its selector removes by + // identity, so repeat its cleanup even for an unregistered consumer without affecting replacements. + if (keySharedMode != KeySharedMode.STICKY && !containsConsumerInstance(consumer)) { // Let the superclass repair stale list membership without touching a replacement's selector state. super.removeConsumer(consumer); return; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersClassic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersClassic.java index 41d4f8e54b057..67b82a08ec686 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersClassic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersClassic.java @@ -234,7 +234,9 @@ private void sortRecentlyJoinedConsumersIfNeeded() { @Override public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { - if (!containsConsumerInstance(consumer)) { + // A STICKY addition can finish its liveness check after dispatcher removal. Its selector removes by + // identity, so repeat its cleanup even for an unregistered consumer without affecting replacements. + if (keySharedMode != KeySharedMode.STICKY && !containsConsumerInstance(consumer)) { // Let the superclass repair stale list membership without touching a replacement's selector state. super.removeConsumer(consumer); return; @@ -248,7 +250,7 @@ public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceE selector.removeConsumer(consumer); super.removeConsumer(consumer); if (recentlyJoinedConsumers != null) { - recentlyJoinedConsumers.remove(consumer); + recentlyJoinedConsumers.keySet().removeIf(c -> c == consumer); if (consumerList.size() == 1) { recentlyJoinedConsumers.clear(); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedSubscriptionUnackedMessagesAccountingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedSubscriptionUnackedMessagesAccountingTest.java index cafad7772e7d2..4ceee7fd6f9b0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedSubscriptionUnackedMessagesAccountingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedSubscriptionUnackedMessagesAccountingTest.java @@ -19,21 +19,26 @@ package org.apache.pulsar.broker.service; import static org.assertj.core.api.Assertions.assertThat; +import io.netty.channel.Channel; import java.util.Collections; import java.util.Optional; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import org.apache.pulsar.broker.service.persistent.AbstractPersistentDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumersClassic; import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.client.api.ConsumerBuilder; +import org.apache.pulsar.client.api.KeySharedPolicy; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerConsumerBase; +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.Range; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.common.api.proto.CommandSubscribe.InitialPosition; -import org.apache.pulsar.common.api.proto.KeySharedMeta; -import org.apache.pulsar.common.api.proto.KeySharedMode; +import org.awaitility.Awaitility; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.DataProvider; @@ -131,17 +136,25 @@ public void testRemovingSameConsumerTwiceDebitsUnackedMessagesOnce() throws Exce @DataProvider public Object[][] subscriptionTypes() { - return new Object[][] {{SubscriptionType.Shared}, {SubscriptionType.Key_Shared}}; + return new Object[][] { + {SubscriptionType.Shared, null}, + {SubscriptionType.Key_Shared, KeySharedPolicy.autoSplitHashRange()}, + {SubscriptionType.Key_Shared, KeySharedPolicy.stickyHashRange().ranges(Range.of(0, 65535))}}; } @Test(dataProvider = "subscriptionTypes", timeOut = 60_000) - public void testStaleRemovalDoesNotRemoveEqualReplacement(SubscriptionType subscriptionType) throws Exception { + public void testStaleRemovalDoesNotRemoveEqualReplacement(SubscriptionType subscriptionType, + KeySharedPolicy keySharedPolicy) throws Exception { String topicName = newTopicName(); + ConsumerBuilder consumerBuilder = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName).subscriptionName(SUBSCRIPTION).subscriptionType(subscriptionType) + .consumerName("replacement-test"); + if (keySharedPolicy != null) { + consumerBuilder.keySharedPolicy(keySharedPolicy); + } try (Producer producer = pulsarClient.newProducer(Schema.STRING) .topic(topicName).enableBatching(false).create(); - org.apache.pulsar.client.api.Consumer client = pulsarClient.newConsumer(Schema.STRING) - .topic(topicName).subscriptionName(SUBSCRIPTION).subscriptionType(subscriptionType) - .consumerName("replacement-test").subscribe()) { + org.apache.pulsar.client.api.Consumer client = consumerBuilder.subscribe()) { producer.send("unacked"); assertThat(client.receive(5, TimeUnit.SECONDS)).isNotNull(); BrokerService brokerService = pulsar.getBrokerService(); @@ -161,7 +174,7 @@ public void testStaleRemovalDoesNotRemoveEqualReplacement(SubscriptionType subsc .subscriptionName(SUBSCRIPTION).subType(original.subType()).isDurable(true) .startMessageId(MessageId.latest).initialPosition(InitialPosition.Latest) .metadata(Collections.emptyMap()).subscriptionProperties(Optional.empty()) - .keySharedMeta(new KeySharedMeta().setKeySharedMode(KeySharedMode.AUTO_SPLIT)) + .keySharedMeta(original.getKeySharedMeta()) .build()).get(10, TimeUnit.SECONDS); try { assertThat(replacement).isNotSameAs(original).isEqualTo(original); @@ -186,6 +199,72 @@ public void testStaleRemovalDoesNotRemoveEqualReplacement(SubscriptionType subsc } } + @Test(timeOut = 60_000) + public void testCloseCleansUpStickyAdditionCompletedAfterRemoval() throws Exception { + String topicName = newTopicName(); + try (PulsarClient otherClient = newPulsarClient(pulsar.getBrokerServiceUrl(), 0); + Producer producer = otherClient.newProducer(Schema.STRING) + .topic(topicName).enableBatching(false).create(); + org.apache.pulsar.client.api.Consumer client = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName).subscriptionName(SUBSCRIPTION).subscriptionType(SubscriptionType.Key_Shared) + .keySharedPolicy(KeySharedPolicy.stickyHashRange().ranges(Range.of(0, 65535))).subscribe()) { + PersistentTopic topic = + (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).orElseThrow(); + AbstractPersistentDispatcherMultipleConsumers dispatcher = + (AbstractPersistentDispatcherMultipleConsumers) topic.getSubscription(SUBSCRIPTION).getDispatcher(); + StickyKeyConsumerSelector selector = ((StickyKeyDispatcher) dispatcher).getSelector(); + Consumer original = dispatcher.getConsumers().get(0); + ServerCnx originalCnx = (ServerCnx) original.cnx(); + Channel channel = originalCnx.ctx().channel(); + // Drain the initial FLOW so it cannot complete the liveness check instead of the held PONG. + Awaitility.await().atMost(10, TimeUnit.SECONDS).until(() -> original.getAvailablePermits() > 0); + // Hold the real liveness response without blocking the event loop or mocking the selector. + channel.eventLoop().submit(() -> channel.config().setAutoRead(false)).get(10, TimeUnit.SECONDS); + try { + CompletableFuture adding = topic.subscribe(SubscriptionOption.builder() + .cnx(topic.getProducers().values().iterator().next().getCnx()) + .consumerId(1234).consumerName("pending-sticky") + .subscriptionName(SUBSCRIPTION).subType(original.subType()).isDurable(true) + .startMessageId(MessageId.latest).initialPosition(InitialPosition.Latest) + .metadata(Collections.emptyMap()).subscriptionProperties(Optional.empty()) + .keySharedMeta(original.getKeySharedMeta()).build()); + // Wait for the liveness request on its owning event loop. The second consumer is already + // registered in the dispatcher, but its selector addition must still be outstanding. + originalCnx.ctx().executor().submit(() -> + assertThat(originalCnx.connectionCheckInProgress).isNotNull().isNotDone()) + .get(10, TimeUnit.SECONDS); + assertThat(adding).isNotDone(); + assertThat(dispatcher.getConsumers()).hasSize(2); + Consumer pending = dispatcher.getConsumers().stream() + .filter(c -> c != original).findFirst().orElseThrow(); + assertThat(pending.cnx()).isNotSameAs(originalCnx); + assertThat(selector.select(0)).isSameAs(original); + + // Consumer.close is also used when cursor reset disconnects consumers while subscribe is pending. + pending.close(); + original.close(); + assertThat(dispatcher.getConsumers()).isEmpty(); + assertThat(selector.select(0)).isNull(); + channel.eventLoop().submit(() -> channel.config().setAutoRead(true)).get(10, TimeUnit.SECONDS); + assertThat(adding.get(10, TimeUnit.SECONDS)).isSameAs(pending); + assertThat(selector.select(0)).as("late addition installed the removed consumer").isSameAs(pending); + + // ServerCnx repeats close when subscribe completes after the client has already timed out. + pending.close(); + assertThat(selector.getConsumerKeyHashRanges()).as("no orphan range after repeated close").isEmpty(); + try (org.apache.pulsar.client.api.Consumer successor = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName).subscriptionName(SUBSCRIPTION).subscriptionType(SubscriptionType.Key_Shared) + .keySharedPolicy(KeySharedPolicy.stickyHashRange().ranges(Range.of(0, 65535))) + .subscribeAsync().get(10, TimeUnit.SECONDS)) { + producer.send("after-cleanup"); + assertThat(successor.receive(5, TimeUnit.SECONDS)).isNotNull(); + } + } finally { + channel.eventLoop().submit(() -> channel.config().setAutoRead(true)).get(10, TimeUnit.SECONDS); + } + } + } + private void assertUnackedMessagesCleared(AbstractPersistentDispatcherMultipleConsumers dispatcher, BrokerService brokerService, String removal) { assertThat(dispatcher.getTotalUnackedMessages()).as("subscription balance after %s", removal).isZero();