From 5ecb222711c639e862a9fc410e90697c1d9853b3 Mon Sep 17 00:00:00 2001 From: Chimdumebi Nebolisa Date: Sun, 18 Jan 2026 19:32:29 -0600 Subject: [PATCH 1/7] Relax validation for defaultNumPartitions on non-partitioned autoTopicCreation override --- .../policies/data/impl/AutoTopicCreationOverrideImpl.java | 5 ----- 1 file changed, 5 deletions(-) diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/AutoTopicCreationOverrideImpl.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/AutoTopicCreationOverrideImpl.java index 52cf1f1829b8d..3e87795bccddd 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/AutoTopicCreationOverrideImpl.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/AutoTopicCreationOverrideImpl.java @@ -51,11 +51,6 @@ public static ValidateResult validateOverride(AutoTopicCreationOverride override if (override.getDefaultNumPartitions() <= 0) { return ValidateResult.fail("[defaultNumPartitions] cannot be less than 1 for partition type."); } - } else if (TopicType.NON_PARTITIONED.toString().equals(override.getTopicType())) { - if (override.getDefaultNumPartitions() != null) { - return ValidateResult.fail("[defaultNumPartitions] is not allowed to be" - + " set when the type is non-partition."); - } } } return ValidateResult.success(); From 7a13f4245b6dc134a946b6c9d925a9b48caacf56 Mon Sep 17 00:00:00 2001 From: Chimdumebi Nebolisa Date: Sun, 18 Jan 2026 19:32:44 -0600 Subject: [PATCH 2/7] Add broker test: accept but ignore defaultNumPartitions for non-partitioned autoTopicCreation override --- .../pulsar/broker/admin/AdminApi2Test.java | 26 +++++++++++++++++++ 1 file changed, 26 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java index eeac3e1c5e24d..5b115bc74988c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java @@ -2458,6 +2458,32 @@ public void testAutoTopicCreationOverrideWithMaxNumPartitionsLimit() throws Exce assertTrue(ex.getCause() instanceof NotAcceptableException); } } + + @Test + public void testAutoTopicCreationOverrideNonPartitionedDefaultNumPartitionsIgnored() throws Exception { + String tenantName = newUniqueName("prop-xyz2"); + String namespaceName = tenantName + "/ns-" + System.currentTimeMillis(); + TenantInfoImpl tenantInfo = new TenantInfoImpl(Set.of("role1", "role2"), Set.of("test")); + admin.tenants().createTenant(tenantName, tenantInfo); + admin.namespaces().createNamespace(namespaceName, Set.of("test")); + + AutoTopicCreationOverride overridePolicy = AutoTopicCreationOverride.builder() + .allowAutoTopicCreation(true) + .topicType(TopicType.NON_PARTITIONED.toString()) + .defaultNumPartitions(5) + .build(); + admin.namespaces().setAutoTopicCreation(namespaceName, overridePolicy); + + AutoTopicCreationOverride storedPolicy = admin.namespaces().getAutoTopicCreation(namespaceName); + assertTrue(storedPolicy.isAllowAutoTopicCreation()); + assertEquals(storedPolicy.getTopicType(), TopicType.NON_PARTITIONED.toString()); + assertEquals(storedPolicy.getDefaultNumPartitions(), Integer.valueOf(5)); + + String topicName = "persistent://" + namespaceName + "/auto-topic-" + UUID.randomUUID(); + pulsarClient.newProducer().topic(topicName).create().close(); + + assertEquals(admin.topics().getPartitionedTopicMetadata(topicName).partitions, 0); + } @Test public void testMaxTopicsPerNamespace() throws Exception { restartClusterAfterTest(); From 505efbe01183b80a8ddd3283b8bd0aa6b4590dec Mon Sep 17 00:00:00 2001 From: Chimdumebi Nebolisa Date: Mon, 19 Jan 2026 07:38:01 -0600 Subject: [PATCH 3/7] Fix validation for defaultNumPartitions on non-partitioned override --- .../impl/AutoTopicCreationOverrideImpl.java | 7 ++++++ .../data/AutoTopicCreationOverrideTest.java | 24 ++++++++++++++++++- 2 files changed, 30 insertions(+), 1 deletion(-) diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/AutoTopicCreationOverrideImpl.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/AutoTopicCreationOverrideImpl.java index 3e87795bccddd..79410b01fd5a3 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/AutoTopicCreationOverrideImpl.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/AutoTopicCreationOverrideImpl.java @@ -44,6 +44,7 @@ public static ValidateResult validateOverride(AutoTopicCreationOverride override if (!TopicType.isValidTopicType(override.getTopicType())) { return ValidateResult.fail(String.format("Unknown topic type [%s]", override.getTopicType())); } + if (TopicType.PARTITIONED.toString().equals(override.getTopicType())) { if (override.getDefaultNumPartitions() == null) { return ValidateResult.fail("[defaultNumPartitions] cannot be null when the type is partitioned."); @@ -51,6 +52,12 @@ public static ValidateResult validateOverride(AutoTopicCreationOverride override if (override.getDefaultNumPartitions() <= 0) { return ValidateResult.fail("[defaultNumPartitions] cannot be less than 1 for partition type."); } + } else if (TopicType.NON_PARTITIONED.toString().equals(override.getTopicType())) { + Integer p = override.getDefaultNumPartitions(); + if (p != null && p != 0 && p != 1) { + return ValidateResult.fail( + "[defaultNumPartitions] must be null, 0 or 1 when the type is non-partitioned."); + } } } return ValidateResult.success(); diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/AutoTopicCreationOverrideTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/AutoTopicCreationOverrideTest.java index 89d41547ba967..e7a34cf4505ad 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/AutoTopicCreationOverrideTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/AutoTopicCreationOverrideTest.java @@ -73,7 +73,28 @@ public void testNumPartitionsNotSet() { } @Test - public void testNumPartitionsOnNonPartitioned() { + public void testNumPartitionsOnNonPartitionedZeroAllowed() { + AutoTopicCreationOverride override = AutoTopicCreationOverride.builder() + .allowAutoTopicCreation(true) + .topicType(TopicType.NON_PARTITIONED.toString()) + .defaultNumPartitions(0) + .build(); + assertTrue(AutoTopicCreationOverrideImpl.validateOverride(override).isSuccess()); + } + + @Test + public void testNumPartitionsOnNonPartitionedOneAllowed() { + AutoTopicCreationOverride override = AutoTopicCreationOverride.builder() + .allowAutoTopicCreation(true) + .topicType(TopicType.NON_PARTITIONED.toString()) + .defaultNumPartitions(1) + .build(); + assertTrue(AutoTopicCreationOverrideImpl.validateOverride(override).isSuccess()); + } + + + @Test + public void testNumPartitionsOnNonPartitionedTooHighRejected() { AutoTopicCreationOverride override = AutoTopicCreationOverride.builder() .allowAutoTopicCreation(true) .topicType(TopicType.NON_PARTITIONED.toString()) @@ -81,4 +102,5 @@ public void testNumPartitionsOnNonPartitioned() { .build(); assertFalse(AutoTopicCreationOverrideImpl.validateOverride(override).isSuccess()); } + } From 95842a7a004db22a522965eb65f56b59355722f8 Mon Sep 17 00:00:00 2001 From: ChimdumebiNebolisa Date: Sat, 11 Jul 2026 19:21:52 -0500 Subject: [PATCH 4/7] [fix][test] Align non-partitioned auto topic creation tests Co-authored-by: Cursor --- .../pulsar/broker/admin/AdminApi2Test.java | 26 +++++++++++++++++-- 1 file changed, 24 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java index 211b7a8065a1c..14138d9f33d4b 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java @@ -2578,20 +2578,42 @@ public void testAutoTopicCreationOverrideNonPartitionedDefaultNumPartitionsIgnor AutoTopicCreationOverride overridePolicy = AutoTopicCreationOverride.builder() .allowAutoTopicCreation(true) .topicType(TopicType.NON_PARTITIONED.toString()) - .defaultNumPartitions(5) + .defaultNumPartitions(1) .build(); admin.namespaces().setAutoTopicCreation(namespaceName, overridePolicy); AutoTopicCreationOverride storedPolicy = admin.namespaces().getAutoTopicCreation(namespaceName); assertTrue(storedPolicy.isAllowAutoTopicCreation()); assertEquals(storedPolicy.getTopicType(), TopicType.NON_PARTITIONED.toString()); - assertEquals(storedPolicy.getDefaultNumPartitions(), Integer.valueOf(5)); + assertEquals(storedPolicy.getDefaultNumPartitions(), Integer.valueOf(1)); String topicName = "persistent://" + namespaceName + "/auto-topic-" + UUID.randomUUID(); pulsarClient.newProducer().topic(topicName).create().close(); assertEquals(admin.topics().getPartitionedTopicMetadata(topicName).partitions, 0); } + + @Test + public void testAutoTopicCreationOverrideNonPartitionedDefaultNumPartitionsRejected() throws Exception { + String tenantName = newUniqueName("prop-xyz2"); + String namespaceName = tenantName + "/ns-" + System.currentTimeMillis(); + TenantInfoImpl tenantInfo = new TenantInfoImpl(Set.of("role1", "role2"), Set.of("test")); + admin.tenants().createTenant(tenantName, tenantInfo); + admin.namespaces().createNamespace(namespaceName, Set.of("test")); + + AutoTopicCreationOverride overridePolicy = AutoTopicCreationOverride.builder() + .allowAutoTopicCreation(true) + .topicType(TopicType.NON_PARTITIONED.toString()) + .defaultNumPartitions(5) + .build(); + try { + admin.namespaces().setAutoTopicCreation(namespaceName, overridePolicy); + fail("Should have failed"); + } catch (PulsarAdminException e) { + assertEquals(e.getStatusCode(), Status.PRECONDITION_FAILED.getStatusCode()); + } + } + @Test public void testMaxTopicsPerNamespace() throws Exception { restartClusterAfterTest(); From 1ab6efbbc4c7060a973768a96cc6f24f4a4f2b8d Mon Sep 17 00:00:00 2001 From: ChimdumebiNebolisa Date: Sun, 2 Aug 2026 16:37:51 -0500 Subject: [PATCH 5/7] Improve consumer resume permit flushing --- .../pulsar/client/impl/ConsumerImpl.java | 14 ++- .../pulsar/client/impl/ConsumerImplTest.java | 100 ++++++++++++++++++ 2 files changed, 113 insertions(+), 1 deletion(-) 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 4b6f9660594c5..eb27a6d1a0bea 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 @@ -1955,6 +1955,18 @@ protected void increaseAvailablePermits(ClientCnx currentCnx, int delta) { } } + private void flushAvailablePermitsToBroker(ClientCnx currentCnx) { + int available = AVAILABLE_PERMITS_UPDATER.get(this); + while (available > 0 && !paused) { + if (AVAILABLE_PERMITS_UPDATER.compareAndSet(this, available, 0)) { + sendFlowPermitsToBroker(currentCnx, available); + break; + } else { + available = AVAILABLE_PERMITS_UPDATER.get(this); + } + } + } + public void increaseAvailablePermits(int delta) { increaseAvailablePermits(cnx(), delta); } @@ -1980,7 +1992,7 @@ public void pause() { public void resume() { if (paused) { paused = false; - increaseAvailablePermits(cnx(), 0); + flushAvailablePermitsToBroker(cnx()); } } diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java index 33732e56a5a44..c75e7f3a4a315 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java @@ -23,6 +23,7 @@ import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.Mockito.any; import static org.mockito.Mockito.atLeast; +import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; @@ -33,6 +34,8 @@ import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; +import io.netty.channel.ChannelHandlerContext; +import java.util.ArrayList; import java.util.Arrays; import java.util.BitSet; import java.util.List; @@ -45,6 +48,7 @@ import java.util.regex.Pattern; import java.util.stream.Collectors; import lombok.Cleanup; +import org.apache.pulsar.client.api.BatchReceivePolicy; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageIdAdv; @@ -58,6 +62,7 @@ import org.apache.pulsar.client.impl.conf.TopicConsumerConfigurationData; import org.apache.pulsar.client.util.ExecutorProvider; import org.apache.pulsar.client.util.ScheduledExecutorProvider; +import org.apache.pulsar.common.api.proto.BaseCommand; import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.util.Backoff; import org.awaitility.Awaitility; @@ -273,6 +278,81 @@ public void testMaxReceiverQueueSize() { Assert.assertEquals(consumer.getAvailablePermits(), permits + 100); } + @Test + public void testResumeFlushesSubThresholdPermitsOnlyOnce() { + ClientCnx cnx = ClientTestFixtures.mockClientCnx(); + List flowPermits = captureFlowPermits(cnx); + consumer.setClientCnx(cnx); + + int threshold = consumer.getCurrentReceiverQueueSize() / 2; + int owed = threshold - 1; + assertThat(owed).isPositive(); + + consumer.pause(); + consumer.increaseAvailablePermits(cnx, owed); + assertThat(flowPermits).isEmpty(); + assertThat(consumer.getAvailablePermits()).isEqualTo(owed); + + consumer.resume(); + assertThat(flowPermits).containsExactly(owed); + assertThat(consumer.getAvailablePermits()).isZero(); + + consumer.resume(); + assertThat(flowPermits).containsExactly(owed); + } + + @Test + public void testIncreaseAvailablePermitsStillUsesThreshold() { + ClientCnx cnx = ClientTestFixtures.mockClientCnx(); + List flowPermits = captureFlowPermits(cnx); + consumer.setClientCnx(cnx); + + int threshold = consumer.getCurrentReceiverQueueSize() / 2; + consumer.increaseAvailablePermits(cnx, threshold - 1); + assertThat(flowPermits).isEmpty(); + assertThat(consumer.getAvailablePermits()).isEqualTo(threshold - 1); + + consumer.increaseAvailablePermits(cnx, 1); + assertThat(flowPermits).containsExactly(threshold); + assertThat(consumer.getAvailablePermits()).isZero(); + } + + @Test + public void testResumeWithAutoScaledReceiverQueueFlushesOnlyPositivePermits() { + ConsumerConfigurationData configuration = new ConsumerConfigurationData<>(); + configuration.setAutoScaledReceiverQueueSizeEnabled(true); + configuration.setBatchReceivePolicy(BatchReceivePolicy.builder().maxNumMessages(3).build()); + executorProvider.shutdownNow(); + internalExecutor.shutdownNow(); + consumerConf = configuration; + createConsumer(configuration); + + ClientCnx cnx = ClientTestFixtures.mockClientCnx(); + List flowPermits = captureFlowPermits(cnx); + consumer.setClientCnx(cnx); + + int receiverQueueSize = consumer.getCurrentReceiverQueueSize(); + assertThat(receiverQueueSize).isEqualTo(4); + + consumer.pause(); + consumer.increaseAvailablePermits(cnx, 1); + consumer.resume(); + assertThat(flowPermits).containsExactly(1); + assertThat(consumer.getAvailablePermits()).isZero(); + + consumer.pause(); + consumer.setCurrentReceiverQueueSize(receiverQueueSize * 2); + consumer.resume(); + assertThat(flowPermits).containsExactly(1, receiverQueueSize); + assertThat(consumer.getAvailablePermits()).isZero(); + + consumer.pause(); + consumer.setCurrentReceiverQueueSize(receiverQueueSize / 2); + consumer.resume(); + assertThat(flowPermits).containsExactly(1, receiverQueueSize); + assertThat(consumer.getAvailablePermits()).isNegative(); + } + @Test public void testTopicPriorityLevel() { ConsumerConfigurationData consumerConf2 = new ConsumerConfigurationData<>(); @@ -320,6 +400,26 @@ public void testAutoGenerateConsumerName() { assertTrue(consumerNamePattern.matcher(consumer.getConsumerName()).matches()); } + private List captureFlowPermits(ClientCnx cnx) { + List flowPermits = new ArrayList<>(); + ChannelHandlerContext ctx = cnx.ctx(); + doAnswer(invocation -> { + ByteBuf command = invocation.getArgument(0); + try { + command.skipBytes(Integer.BYTES); + int commandSize = command.readInt(); + BaseCommand parsedCommand = new BaseCommand(); + parsedCommand.parseFrom(command, commandSize); + assertThat(parsedCommand.getType()).isEqualTo(BaseCommand.Type.FLOW); + flowPermits.add(parsedCommand.getFlow().getMessagePermits()); + } finally { + command.release(); + } + return null; + }).when(ctx).writeAndFlush(any(ByteBuf.class), any()); + return flowPermits; + } + @Test(invocationTimeOut = 1000) @SuppressWarnings({"rawtypes", "unchecked"}) public void testUpdateAutoScaleReceiverQueueHintRaceWithConcurrentDrain() { From e43450ce70ab0e8c823424d97060421e9fbd817e Mon Sep 17 00:00:00 2001 From: ChimdumebiNebolisa Date: Sun, 2 Aug 2026 17:24:07 -0500 Subject: [PATCH 6/7] Revert "Improve consumer resume permit flushing" This reverts commit 1ab6efbbc4c7060a973768a96cc6f24f4a4f2b8d. --- .../pulsar/client/impl/ConsumerImpl.java | 14 +-- .../pulsar/client/impl/ConsumerImplTest.java | 100 ------------------ 2 files changed, 1 insertion(+), 113 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 eb27a6d1a0bea..4b6f9660594c5 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 @@ -1955,18 +1955,6 @@ protected void increaseAvailablePermits(ClientCnx currentCnx, int delta) { } } - private void flushAvailablePermitsToBroker(ClientCnx currentCnx) { - int available = AVAILABLE_PERMITS_UPDATER.get(this); - while (available > 0 && !paused) { - if (AVAILABLE_PERMITS_UPDATER.compareAndSet(this, available, 0)) { - sendFlowPermitsToBroker(currentCnx, available); - break; - } else { - available = AVAILABLE_PERMITS_UPDATER.get(this); - } - } - } - public void increaseAvailablePermits(int delta) { increaseAvailablePermits(cnx(), delta); } @@ -1992,7 +1980,7 @@ public void pause() { public void resume() { if (paused) { paused = false; - flushAvailablePermitsToBroker(cnx()); + increaseAvailablePermits(cnx(), 0); } } diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java index c75e7f3a4a315..33732e56a5a44 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java @@ -23,7 +23,6 @@ import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.Mockito.any; import static org.mockito.Mockito.atLeast; -import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; @@ -34,8 +33,6 @@ import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; -import io.netty.channel.ChannelHandlerContext; -import java.util.ArrayList; import java.util.Arrays; import java.util.BitSet; import java.util.List; @@ -48,7 +45,6 @@ import java.util.regex.Pattern; import java.util.stream.Collectors; import lombok.Cleanup; -import org.apache.pulsar.client.api.BatchReceivePolicy; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageIdAdv; @@ -62,7 +58,6 @@ import org.apache.pulsar.client.impl.conf.TopicConsumerConfigurationData; import org.apache.pulsar.client.util.ExecutorProvider; import org.apache.pulsar.client.util.ScheduledExecutorProvider; -import org.apache.pulsar.common.api.proto.BaseCommand; import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.util.Backoff; import org.awaitility.Awaitility; @@ -278,81 +273,6 @@ public void testMaxReceiverQueueSize() { Assert.assertEquals(consumer.getAvailablePermits(), permits + 100); } - @Test - public void testResumeFlushesSubThresholdPermitsOnlyOnce() { - ClientCnx cnx = ClientTestFixtures.mockClientCnx(); - List flowPermits = captureFlowPermits(cnx); - consumer.setClientCnx(cnx); - - int threshold = consumer.getCurrentReceiverQueueSize() / 2; - int owed = threshold - 1; - assertThat(owed).isPositive(); - - consumer.pause(); - consumer.increaseAvailablePermits(cnx, owed); - assertThat(flowPermits).isEmpty(); - assertThat(consumer.getAvailablePermits()).isEqualTo(owed); - - consumer.resume(); - assertThat(flowPermits).containsExactly(owed); - assertThat(consumer.getAvailablePermits()).isZero(); - - consumer.resume(); - assertThat(flowPermits).containsExactly(owed); - } - - @Test - public void testIncreaseAvailablePermitsStillUsesThreshold() { - ClientCnx cnx = ClientTestFixtures.mockClientCnx(); - List flowPermits = captureFlowPermits(cnx); - consumer.setClientCnx(cnx); - - int threshold = consumer.getCurrentReceiverQueueSize() / 2; - consumer.increaseAvailablePermits(cnx, threshold - 1); - assertThat(flowPermits).isEmpty(); - assertThat(consumer.getAvailablePermits()).isEqualTo(threshold - 1); - - consumer.increaseAvailablePermits(cnx, 1); - assertThat(flowPermits).containsExactly(threshold); - assertThat(consumer.getAvailablePermits()).isZero(); - } - - @Test - public void testResumeWithAutoScaledReceiverQueueFlushesOnlyPositivePermits() { - ConsumerConfigurationData configuration = new ConsumerConfigurationData<>(); - configuration.setAutoScaledReceiverQueueSizeEnabled(true); - configuration.setBatchReceivePolicy(BatchReceivePolicy.builder().maxNumMessages(3).build()); - executorProvider.shutdownNow(); - internalExecutor.shutdownNow(); - consumerConf = configuration; - createConsumer(configuration); - - ClientCnx cnx = ClientTestFixtures.mockClientCnx(); - List flowPermits = captureFlowPermits(cnx); - consumer.setClientCnx(cnx); - - int receiverQueueSize = consumer.getCurrentReceiverQueueSize(); - assertThat(receiverQueueSize).isEqualTo(4); - - consumer.pause(); - consumer.increaseAvailablePermits(cnx, 1); - consumer.resume(); - assertThat(flowPermits).containsExactly(1); - assertThat(consumer.getAvailablePermits()).isZero(); - - consumer.pause(); - consumer.setCurrentReceiverQueueSize(receiverQueueSize * 2); - consumer.resume(); - assertThat(flowPermits).containsExactly(1, receiverQueueSize); - assertThat(consumer.getAvailablePermits()).isZero(); - - consumer.pause(); - consumer.setCurrentReceiverQueueSize(receiverQueueSize / 2); - consumer.resume(); - assertThat(flowPermits).containsExactly(1, receiverQueueSize); - assertThat(consumer.getAvailablePermits()).isNegative(); - } - @Test public void testTopicPriorityLevel() { ConsumerConfigurationData consumerConf2 = new ConsumerConfigurationData<>(); @@ -400,26 +320,6 @@ public void testAutoGenerateConsumerName() { assertTrue(consumerNamePattern.matcher(consumer.getConsumerName()).matches()); } - private List captureFlowPermits(ClientCnx cnx) { - List flowPermits = new ArrayList<>(); - ChannelHandlerContext ctx = cnx.ctx(); - doAnswer(invocation -> { - ByteBuf command = invocation.getArgument(0); - try { - command.skipBytes(Integer.BYTES); - int commandSize = command.readInt(); - BaseCommand parsedCommand = new BaseCommand(); - parsedCommand.parseFrom(command, commandSize); - assertThat(parsedCommand.getType()).isEqualTo(BaseCommand.Type.FLOW); - flowPermits.add(parsedCommand.getFlow().getMessagePermits()); - } finally { - command.release(); - } - return null; - }).when(ctx).writeAndFlush(any(ByteBuf.class), any()); - return flowPermits; - } - @Test(invocationTimeOut = 1000) @SuppressWarnings({"rawtypes", "unchecked"}) public void testUpdateAutoScaleReceiverQueueHintRaceWithConcurrentDrain() { From 1bdd1b6d0d61477b1b7313b68b1bf2b7903c9bb9 Mon Sep 17 00:00:00 2001 From: ChimdumebiNebolisa Date: Fri, 21 Aug 2026 16:50:42 -0500 Subject: [PATCH 7/7] Fix stale non-partitioned validation assertion --- .../broker/service/SetReplicationClustersValidationTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SetReplicationClustersValidationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SetReplicationClustersValidationTest.java index a4becebcec9fd..9bbd21c17ffaa 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SetReplicationClustersValidationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SetReplicationClustersValidationTest.java @@ -245,7 +245,7 @@ public void testSetReplicationClustersWithMismatchedDefaultNumPartitions() throw admin1.namespaces().setAutoTopicCreation(namespace, policy1); fail("Expected behaviour: Pulsar does not allow setting non-partitioned and a certain partition counts"); } catch (Exception ex) { - assertTrue(ex.getMessage().contains("is not allowed to be set when the type is non-partition")); + assertTrue(ex.getMessage().contains("must be null, 0 or 1 when the type is non-partitioned")); } }