From 39482bd035d393034554ef6e04a052dc09d8863e Mon Sep 17 00:00:00 2001 From: Jiwei Guo Date: Wed, 15 Jul 2026 20:29:22 +0800 Subject: [PATCH 1/3] Fix test --- .../pulsar/broker/admin/AdminApi2Test.java | 32 ++++++++++++------- 1 file changed, 20 insertions(+), 12 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 58fdd9f8a246b..1aade9bd60dc6 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 @@ -3493,41 +3493,49 @@ public void testCompactionPriority() throws Exception { final String namespace = newUniqueName(defaultTenant + "/ns"); admin.namespaces().createNamespace(namespace, Set.of("test")); final String topic = "persistent://" + namespace + "/topic" + UUID.randomUUID(); - pulsarClient.newProducer().topic(topic).create().close(); + Producer producer = pulsarClient.newProducer().topic(topic).create(); + producer.send("message".getBytes()); + producer.close(); TopicName topicName = TopicName.get(topic); PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService() .getTopicIfExists(topic).get().get(); PersistentTopic mockTopic = spy(persistentTopic); + org.mockito.Mockito.doReturn(CompletableFuture.completedFuture(null)) + .when(mockTopic).triggerCompactionWithCheckHasMoreMessages(); mockTopic.checkCompaction(); // Disabled by default - verify(mockTopic, times(0)).triggerCompaction(); + verify(mockTopic, times(0)).triggerCompactionWithCheckHasMoreMessages(); // Set namespace-level policy admin.namespaces().setCompactionThreshold(namespace, 1); Awaitility.await().untilAsserted(() -> assertNotNull(admin.namespaces().getCompactionThreshold(namespace))); - ManagedLedger managedLedger = persistentTopic.getManagedLedger(); - Field field = managedLedger.getClass().getDeclaredField("totalSize"); - field.setAccessible(true); - field.setLong(managedLedger, 1000L); + Awaitility.await().untilAsserted(() -> assertTrue(persistentTopic.isCompactionEnabled())); - mockTopic.checkCompaction(); - verify(mockTopic, times(1)).triggerCompaction(); + Awaitility.await().untilAsserted(() -> { + mockTopic.checkCompaction(); + verify(mockTopic, times(1)).triggerCompactionWithCheckHasMoreMessages(); + }); //Set topic-level policy admin.topics().setCompactionThreshold(topic, 0); Awaitility.await().untilAsserted(() -> assertNotNull(admin.topics().getCompactionThreshold(topic))); + Awaitility.await().untilAsserted(() -> assertFalse(persistentTopic.isCompactionEnabled())); mockTopic.checkCompaction(); - verify(mockTopic, times(1)).triggerCompaction(); + verify(mockTopic, times(1)).triggerCompactionWithCheckHasMoreMessages(); // Remove topic-level policy admin.topics().removeCompactionThreshold(topic); Awaitility.await().untilAsserted(() -> assertNull(admin.topics().getCompactionThreshold(topic))); - mockTopic.checkCompaction(); - verify(mockTopic, times(2)).triggerCompaction(); + Awaitility.await().untilAsserted(() -> assertTrue(persistentTopic.isCompactionEnabled())); + Awaitility.await().untilAsserted(() -> { + mockTopic.checkCompaction(); + verify(mockTopic, times(2)).triggerCompactionWithCheckHasMoreMessages(); + }); // Remove namespace-level policy admin.namespaces().removeCompactionThreshold(namespace); Awaitility.await().untilAsserted(() -> assertNull(admin.namespaces().getCompactionThreshold(namespace))); + Awaitility.await().untilAsserted(() -> assertFalse(persistentTopic.isCompactionEnabled())); mockTopic.checkCompaction(); - verify(mockTopic, times(2)).triggerCompaction(); + verify(mockTopic, times(2)).triggerCompactionWithCheckHasMoreMessages(); } @Test From 04a65c2d18724ef7c155217e3d9c64e1b30b1aae Mon Sep 17 00:00:00 2001 From: Jiwei Guo Date: Fri, 17 Jul 2026 10:34:37 +0800 Subject: [PATCH 2/3] fix checkstyle --- .../test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java | 2 -- 1 file changed, 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 1aade9bd60dc6..5e41c05df5c27 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 @@ -44,7 +44,6 @@ import io.grpc.netty.shaded.io.netty.util.concurrent.FastThreadLocal; import jakarta.ws.rs.NotAcceptableException; import jakarta.ws.rs.core.Response.Status; -import java.lang.reflect.Field; import java.net.URI; import java.net.URL; import java.net.http.HttpClient; @@ -72,7 +71,6 @@ import lombok.CustomLog; import lombok.Data; import org.apache.bookkeeper.mledger.AsyncCallbacks; -import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; From 82a788a29c8d1e0288c3250ad926e808187cfc41 Mon Sep 17 00:00:00 2001 From: Jiwei Guo Date: Fri, 17 Jul 2026 16:33:33 +0800 Subject: [PATCH 3/3] address comments --- .../java/org/apache/pulsar/broker/admin/AdminApi2Test.java | 5 +++-- 1 file changed, 3 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 5e41c05df5c27..1d0a8e4cdf93a 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 @@ -28,6 +28,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -3491,14 +3492,14 @@ public void testCompactionPriority() throws Exception { final String namespace = newUniqueName(defaultTenant + "/ns"); admin.namespaces().createNamespace(namespace, Set.of("test")); final String topic = "persistent://" + namespace + "/topic" + UUID.randomUUID(); + @Cleanup Producer producer = pulsarClient.newProducer().topic(topic).create(); producer.send("message".getBytes()); - producer.close(); TopicName topicName = TopicName.get(topic); PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService() .getTopicIfExists(topic).get().get(); PersistentTopic mockTopic = spy(persistentTopic); - org.mockito.Mockito.doReturn(CompletableFuture.completedFuture(null)) + doReturn(CompletableFuture.completedFuture(null)) .when(mockTopic).triggerCompactionWithCheckHasMoreMessages(); mockTopic.checkCompaction(); // Disabled by default