This is an automated email from the ASF dual-hosted git repository.
lhotari pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new 0629dc5add9 [improve][broker] Load topic policies on non-persistent
topic load and gate the policy replay (#26134)
0629dc5add9 is described below
commit 0629dc5add92c638aa37a01eacb7933723452a5c
Author: Lari Hotari <[email protected]>
AuthorDate: Thu Jul 2 21:16:19 2026 +0300
[improve][broker] Load topic policies on non-persistent topic load and gate
the policy replay (#26134)
---
.../apache/pulsar/broker/ServiceConfiguration.java | 9 ++
.../pulsar/broker/service/AbstractTopic.java | 106 +++++++++++++++++++--
.../SystemTopicBasedTopicPoliciesService.java | 39 ++++++--
.../broker/service/TopicPolicyListenerWrapper.java | 68 +++++++++++--
.../service/nonpersistent/NonPersistentTopic.java | 25 ++---
.../broker/service/persistent/PersistentTopic.java | 58 +----------
.../admin/MetadataStoreTopicPoliciesTest.java | 8 ++
.../pulsar/broker/admin/TopicPoliciesTest.java | 67 +++++++++++++
.../SystemTopicBasedTopicPoliciesServiceTest.java | 63 ++++++++++++
.../service/TopicPolicyListenerWrapperTest.java | 102 ++++++++++++++++++--
.../service/persistent/PersistentTopicTest.java | 7 ++
11 files changed, 450 insertions(+), 102 deletions(-)
diff --git
a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
index 67112f19a77..9b38c6127b0 100644
---
a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
+++
b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
@@ -1992,6 +1992,15 @@ public class ServiceConfiguration implements
PulsarConfiguration {
+ "please enable the system topic first.")
private boolean topicLevelPoliciesEnabled = true;
+ @FieldContext(
+ category = CATEGORY_SERVER,
+ doc = "When enabled, all registered topic-policy listeners in a
namespace are re-notified with the current"
+ + " topic policies after the namespace's topic-policy
cache finishes its initial load. Topics load"
+ + " and apply their own policies when they are loaded, so
this broadcast is normally redundant; it"
+ + " is only needed for custom plugins that register
TopicPolicyListeners and depend on it for"
+ + " backwards compatibility. Disabled by default.")
+ private boolean topicPolicyListenerReplayEnabled = false;
+
@FieldContext(
category = CATEGORY_SERVER,
doc = """
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java
index ff31d66ef9f..3975b180cf6 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java
@@ -50,6 +50,7 @@ import java.util.concurrent.atomic.AtomicIntegerFieldUpdater;
import java.util.concurrent.atomic.AtomicLongFieldUpdater;
import java.util.concurrent.atomic.LongAdder;
import java.util.concurrent.locks.ReentrantReadWriteLock;
+import java.util.function.Function;
import java.util.function.ToLongFunction;
import lombok.Getter;
import lombok.Setter;
@@ -60,6 +61,7 @@ import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.pulsar.broker.PulsarService;
import org.apache.pulsar.broker.ServiceConfiguration;
+import
org.apache.pulsar.broker.loadbalance.extensions.ExtensibleLoadManagerImpl;
import org.apache.pulsar.broker.resourcegroup.ResourceGroup;
import org.apache.pulsar.broker.resourcegroup.ResourceGroupPublishLimiter;
import org.apache.pulsar.broker.resources.NamespaceResources;
@@ -119,6 +121,10 @@ public abstract class AbstractTopic implements Topic,
TopicPolicyListener {
protected final BrokerService brokerService;
+ // Wraps this topic as a TopicPolicyListener so topic-policy updates
received while the initial policy is still
+ // loading are buffered and applied in order once initTopicPolicy()
completes initialization.
+ protected final TopicPolicyListenerWrapper topicPolicyListener = new
TopicPolicyListenerWrapper(this);
+
// Prefix for replication cursors
protected final String replicatorPrefix;
@@ -180,7 +186,17 @@ public abstract class AbstractTopic implements Topic,
TopicPolicyListener {
AtomicLongFieldUpdater.newUpdater(AbstractTopic.class,
"usageCount");
private volatile long usageCount = 0;
- private Map<String/*subscription*/, SubscriptionPolicies>
subscriptionPolicies = Collections.emptyMap();
+ // Effective per-subscription policies, merged from the local and global
topic policies with local precedence.
+ // Unlike the PolicyHierarchyValue-backed fields, subscriptionPolicies is
a plain map. It is kept as the merge of
+ // the two scopes below (rather than assigned directly) because the
local-before-global initialization order
+ // (see TopicPolicyListenerWrapper) would otherwise let the global map --
empty by default in TopicPolicies --
+ // overwrite and clear the local per-subscription policies that were
applied just before it.
+ // subscriptionPolicies is volatile because it is read on dispatch threads
(getSubscriptionDispatchRate) while it
+ // is updated on the policy-update thread.
localSubscriptionPolicies/globalSubscriptionPolicies are only ever read
+ // and written on the (single) policy-update thread, so they don't need to
be volatile.
+ private volatile Map<String/*subscription*/, SubscriptionPolicies>
subscriptionPolicies = Collections.emptyMap();
+ private Map<String/*subscription*/, SubscriptionPolicies>
localSubscriptionPolicies = Collections.emptyMap();
+ private Map<String/*subscription*/, SubscriptionPolicies>
globalSubscriptionPolicies = Collections.emptyMap();
protected final LongAdder msgOutFromRemovedSubscriptions = new LongAdder();
protected final LongAdder bytesOutFromRemovedSubscriptions = new
LongAdder();
@@ -320,11 +336,36 @@ public abstract class AbstractTopic implements Topic,
TopicPolicyListener {
topicPolicies.getEntryFilters().updateTopicValue(data.getEntryFilters(),
isGlobalPolicies);
topicPolicies.getDispatcherPauseOnAckStatePersistentEnabled()
.updateTopicValue(data.getDispatcherPauseOnAckStatePersistentEnabled(),
isGlobalPolicies);
- this.subscriptionPolicies = data.getSubscriptionPolicies();
+
+ // Merge instead of assigning directly: keep the local and global
per-subscription policies separately and
+ // recompute the effective map with local precedence, so applying the
(default-empty) global map does not
+ // clear the local per-subscription policies during the
local-before-global initialization.
+ if (isGlobalPolicies) {
+ globalSubscriptionPolicies = data.getSubscriptionPolicies();
+ } else {
+ localSubscriptionPolicies = data.getSubscriptionPolicies();
+ }
+ subscriptionPolicies =
mergeSubscriptionPolicies(globalSubscriptionPolicies,
localSubscriptionPolicies);
updateEntryFilters();
}
+ // Merges the global and local per-subscription policies with local
precedence: a subscription present in the
+ // local policies keeps its local value; otherwise the global value (if
any) is used.
+ private static Map<String, SubscriptionPolicies> mergeSubscriptionPolicies(
+ Map<String, SubscriptionPolicies> globalSubscriptionPolicies,
+ Map<String, SubscriptionPolicies> localSubscriptionPolicies) {
+ if (globalSubscriptionPolicies.isEmpty()) {
+ return localSubscriptionPolicies;
+ }
+ if (localSubscriptionPolicies.isEmpty()) {
+ return globalSubscriptionPolicies;
+ }
+ Map<String, SubscriptionPolicies> merged = new
HashMap<>(globalSubscriptionPolicies);
+ merged.putAll(localSubscriptionPolicies);
+ return merged;
+ }
+
protected void updateTopicPolicyByNamespacePolicy(Policies
namespacePolicies) {
log.debug()
.attr("namespacePolicies", namespacePolicies)
@@ -566,12 +607,7 @@ public abstract class AbstractTopic implements Topic,
TopicPolicyListener {
}
protected TopicPolicyListener getTopicPolicyListener() {
- return this;
- }
-
- protected void registerTopicPolicyListener() {
- brokerService.getPulsar().getTopicPoliciesService()
-
.registerListenerAsync(TopicName.getPartitionedTopicName(topic),
getTopicPolicyListener());
+ return topicPolicyListener;
}
protected void unregisterTopicPolicyListener() {
@@ -579,6 +615,60 @@ public abstract class AbstractTopic implements Topic,
TopicPolicyListener {
.unregisterListener(TopicName.getPartitionedTopicName(topic),
getTopicPolicyListener());
}
+ /**
+ * Registers the topic-policy listener and applies the topic's initial
policies (global and local) to this topic.
+ * Shared by {@link
org.apache.pulsar.broker.service.persistent.PersistentTopic} and
+ * {@link
org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic} so both load
their own policies on
+ * topic load, which removes the need to broadcast every topic's policy
when a namespace's policy cache finishes
+ * loading (see {@code topicPolicyListenerReplayEnabled}).
+ *
+ * <p>Each call re-initializes the listener wrapper and, whatever the
outcome, always completes its initialization
+ * afterwards, so the wrapper never stays in the buffering phase (dropping
updates) even if policy loading fails.
+ * This makes the method safe to run again (e.g. a future retry); runs are
expected to be serialized.
+ */
+ protected CompletableFuture<Void> initTopicPolicy() {
+ final var topicPoliciesService =
brokerService.getPulsar().getTopicPoliciesService();
+ final var partitionedTopicName =
TopicName.getPartitionedTopicName(topic);
+
+ // Begin a fresh initialization phase: updates are buffered until
initialization completes below. This resets
+ // any previous phase so the method can be run again.
+ topicPolicyListener.startInitialization();
+ CompletableFuture<Void> initTopicPolicyFuture =
+
topicPoliciesService.registerListenerAsync(partitionedTopicName,
topicPolicyListener)
+ .thenCompose(registered -> {
+ if (!registered) {
+ return CompletableFuture.completedFuture(null);
+ }
+ if
(ExtensibleLoadManagerImpl.isInternalTopic(topic)) {
+ // Internal topics don't load topic-level
policies
+ return CompletableFuture.completedFuture(null);
+ }
+ // future for fetching global topic policies
+ CompletableFuture<Optional<TopicPolicies>>
globalPoliciesFuture =
+
topicPoliciesService.getTopicPoliciesAsync(partitionedTopicName,
+
TopicPoliciesService.GetType.GLOBAL_ONLY);
+ // future for fetching local topic policies
+ CompletableFuture<Optional<TopicPolicies>>
localPoliciesFuture =
+
topicPoliciesService.getTopicPoliciesAsync(partitionedTopicName,
+
TopicPoliciesService.GetType.LOCAL_ONLY);
+ return
globalPoliciesFuture.thenCombine(localPoliciesFuture, (global, local) -> {
+ // finally update the topic policies with the
latest value or loaded value
+ return CompletableFuture.runAsync(() ->
+
topicPolicyListener.completeInitialization(global.orElse(null),
+ local.orElse(null)),
+ getPoliciesNotifyThread());
+ }).thenCompose(Function.identity());
+ });
+ // Whatever the outcome -- success, failure, or the listener not being
registered -- make sure the wrapper
+ // leaves the initialization (buffering) phase, so it forwards any
buffered value plus all future live updates
+ // instead of dropping them. This is a no-op when the loaded policies
were already applied above. Return the
+ // whenComplete stage (not initTopicPolicyFuture) so the returned
future completes only after this has run, and
+ // whenComplete's pass-through semantics carry the original success or
failure to the caller's initialize().
+ return initTopicPolicyFuture.whenCompleteAsync((v, ex) -> {
+ topicPolicyListener.completeInitializationUnlessAlreadyCompleted();
+ }, getPoliciesNotifyThread());
+ }
+
protected boolean isSameAddressProducersExceeded(Producer producer) {
if (isSystemTopic() || producer.isRemote()) {
return false;
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java
index f565eb026f4..8b4cae18bc4 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java
@@ -892,18 +892,43 @@ public class SystemTopicBasedTopicPoliciesService
implements TopicPoliciesServic
.attr("topic", reader.getSystemTopic().getTopicName())
.log("Reach the end of the system topic.");
- // replay policy message
- List<CompletableFuture<Void>> notifyFutures = new
ArrayList<>();
- for (Map.Entry<TopicName, TopicPolicies> entry :
policiesCache.entrySet()) {
- TopicName topicName = entry.getKey();
- TopicPolicies policies = entry.getValue();
- notifyFutures.add(notifyListenersForTopicAsync(topicName,
policies));
+ // Optionally re-notify the topic-policy listeners for this
namespace with the cached policies.
+ // Topics load their own policies on load
(AbstractTopic#initTopicPolicy), so this replay is only
+ // needed for custom plugins that register
TopicPolicyListeners and rely on the broadcast; it is
+ // off by default (topicPolicyListenerReplayEnabled).
+ if
(pulsarService.getConfiguration().isTopicPolicyListenerReplayEnabled()) {
+ NamespaceName namespaceObject =
reader.getSystemTopic().getTopicName().getNamespaceObject();
+ FutureUtil.completeAfter(future,
replayTopicPolicyListeners(namespaceObject));
+ } else {
+ future.complete(null);
}
- FutureUtil.completeAfter(future,
FutureUtil.waitForAll(notifyFutures));
}
});
}
+ /**
+ * Re-notifies the registered topic-policy listeners for every topic in
{@code namespace} with the currently
+ * cached policies (both local and global). Topics apply their own
policies on load, so this is only needed for
+ * custom plugins that register {@link TopicPolicyListener}s and rely on
the broadcast when a namespace's policy
+ * cache finishes loading (see {@code topicPolicyListenerReplayEnabled}).
+ */
+ @VisibleForTesting
+ CompletableFuture<Void> replayTopicPolicyListeners(NamespaceName
namespace) {
+ List<CompletableFuture<Void>> notifyFutures = new ArrayList<>();
+ addNamespacePolicyNotifications(policiesCache, namespace,
notifyFutures);
+ addNamespacePolicyNotifications(globalPoliciesCache, namespace,
notifyFutures);
+ return FutureUtil.waitForAll(notifyFutures);
+ }
+
+ private void addNamespacePolicyNotifications(Map<TopicName, TopicPolicies>
cache, NamespaceName namespace,
+ List<CompletableFuture<Void>>
notifyFutures) {
+ for (Map.Entry<TopicName, TopicPolicies> entry : cache.entrySet()) {
+ if (Objects.equals(entry.getKey().getNamespaceObject(),
namespace)) {
+ notifyFutures.add(notifyListenersForTopicAsync(entry.getKey(),
entry.getValue()));
+ }
+ }
+ }
+
// Full teardown of a namespace's topic-policies state: removes and closes
the reader, the message-handler
// tracker, the cached policies and the init future. Used when the whole
namespace is unloaded.
@VisibleForTesting
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TopicPolicyListenerWrapper.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TopicPolicyListenerWrapper.java
index 1e2f4336a3c..da22b44a632 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TopicPolicyListenerWrapper.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TopicPolicyListenerWrapper.java
@@ -26,10 +26,15 @@ import org.apache.pulsar.common.policies.data.TopicPolicies;
/**
* This TopicPolicyListener is used as a wrapper for the real
TopicPolicyListener.
- * This prevents a race condition in initialization where the topic policy
state can change while the
- * topic policy state is being applied to the topic in
PersistentTopic#initTopicPolicy() method or in
- * NonPersistentTopic#initialize method. The impact of the race conditions is
that the topic policy state would
- * be left in an inconsistent state until another update arrives. This is a
rare corner case, but possible.
+ * This prevents a race condition in initialization where the topic policy
state can change while the topic policy
+ * state is being applied to the topic in AbstractTopic#initTopicPolicy(). The
impact of the race condition is that the
+ * topic policy state would be left inconsistent until another update arrives.
This is a rare corner case, but possible.
+ *
+ * <p>Updates received while initializing are buffered (only the latest per
scope is kept) and applied by
+ * {@link #completeInitialization}; updates received afterwards are forwarded
immediately. The wrapper is reusable so
+ * that AbstractTopic#initTopicPolicy() can be run again -- for example to
retry it, which this enables but does not
+ * implement. {@link #startInitialization()} begins a new buffering phase and
initialization completes at most once per
+ * phase. Concurrent initialization phases are not supported; {@code
initTopicPolicy} runs are serialized by the caller.
*/
@CustomLog
public class TopicPolicyListenerWrapper implements TopicPolicyListener {
@@ -42,12 +47,29 @@ public class TopicPolicyListenerWrapper implements
TopicPolicyListener {
private Optional<TopicPolicies> latestGlobalPolicies;
private Optional<TopicPolicies> latestLocalPolicies;
private boolean initialized;
- private final long createdTimestampNanos = System.nanoTime();
+ // Timestamp when the current initialization phase started, set by
startInitialization(). Used only to warn if the
+ // phase takes too long (i.e. completeInitialization was never called
after policy loading started).
+ private long initializationStartedNanos;
private static final long INITIALIZATION_WARNING_LOG_INTERVAL_NANOS =
TimeUnit.SECONDS.toNanos(30);
private int lastIntervalLogged;
public TopicPolicyListenerWrapper(TopicPolicyListener realTopicListener) {
this.realTopicListener = realTopicListener;
+ startInitialization();
+ }
+
+ /**
+ * Starts (or restarts) the initialization phase: {@link #onUpdate}
buffers updates (keeping only the latest per
+ * scope) instead of forwarding them, until {@link
#completeInitialization} applies them. Called at the start of
+ * every {@code initTopicPolicy} run so the method can be re-run cleanly
(which is what would let a future change
+ * retry it; no retry is implemented here). Runs before the listener is
registered, so no update can arrive before
+ * the phase (and its warning timer) has started.
+ */
+ public synchronized void startInitialization() {
+ initialized = false;
+ latestGlobalPolicies = null;
+ latestLocalPolicies = null;
+ initializationStartedNanos = System.nanoTime();
}
@Override
@@ -77,21 +99,49 @@ public class TopicPolicyListenerWrapper implements
TopicPolicyListener {
/**
* Complete initialization of the TopicPolicyListenerWrapper and emit the
latest policies to the real listener.
+ *
* @param loadedGlobalPolicies the loaded global policies
- * @param loadedLocalPolicies the loaded local policies
+ * @param loadedLocalPolicies the loaded local policies
*/
public synchronized void completeInitialization(TopicPolicies
loadedGlobalPolicies,
TopicPolicies
loadedLocalPolicies) {
+ // Idempotent: an initialization phase completes at most once.
initTopicPolicy runs a terminal
+ // completeInitializationUnlessAlreadyCompleted() after applying the
loaded policies, so a later call must be a
+ // no-op and must not re-emit policies. A new phase is started
explicitly via startInitialization().
+ if (initialized) {
+ return;
+ }
+
// The listener might have received a newer value (or a delete) than
the loaded one while the loading
// was happening; prefer the latest value received during
initialization, falling back to the loaded
// value only when nothing was received for that scope.
- emitInitialPolicies(latestGlobalPolicies, loadedGlobalPolicies);
+ //
+ // Emit the local policy before the global policy. A local topic
policy takes precedence over a global one,
+ // so applying the local value first means that by the time the global
value is applied the local override is
+ // already in place and the merged (local-wins) result is what takes
effect. Emitting the global value first
+ // would briefly apply it on its own and let a global-only setting act
before the local policy overrides it --
+ // e.g. a compaction subscription being created for a global
compaction policy even though the local policy
+ // disables compaction. This does not fully solve such ordering
hazards, but it removes them whenever a local
+ // policy exists. When no local policy exists nothing is emitted for
the local scope (see emitInitialPolicies),
+ // so this ordering does not change behavior for topics that only have
a global policy.
emitInitialPolicies(latestLocalPolicies, loadedLocalPolicies);
+ emitInitialPolicies(latestGlobalPolicies, loadedGlobalPolicies);
+
latestGlobalPolicies = null;
latestLocalPolicies = null;
initialized = true;
}
+ /**
+ * Completes initialization with no loaded policies, unless it has already
completed. Used as a safety net at the
+ * end of {@code initTopicPolicy} so the wrapper always leaves the
buffering phase -- even when the listener was not
+ * registered or policy loading failed -- and therefore stops dropping
updates: it emits any buffered value and
+ * forwards all future live updates. A no-op once initialization has
completed (e.g. with loaded policies).
+ */
+ public synchronized void completeInitializationUnlessAlreadyCompleted() {
+ completeInitialization(null, null);
+ }
+
private void emitInitialPolicies(Optional<TopicPolicies> latestReceived,
TopicPolicies loaded) {
if (latestReceived != null) {
// A value (or a delete) was received during initialization; it
supersedes the loaded value.
@@ -104,11 +154,11 @@ public class TopicPolicyListenerWrapper implements
TopicPolicyListener {
// warn if the initialization takes too long and updates have been received
// this helps detect issues where completeInitialization didn't get called
after loading policies
private void maybeLogWarning() {
- long durationNanos = System.nanoTime() - createdTimestampNanos;
+ long durationNanos = System.nanoTime() - initializationStartedNanos;
int warningLogIntervalCount = (int) (durationNanos /
INITIALIZATION_WARNING_LOG_INTERVAL_NANOS);
if (warningLogIntervalCount > lastIntervalLogged) {
log.warn().attr("topicPolicyListener", realTopicListener)
- .attr("sinceCreationMs",
TimeUnit.NANOSECONDS.toMillis(durationNanos))
+ .attr("sinceInitializationStartedMs",
TimeUnit.NANOSECONDS.toMillis(durationNanos))
.log("TopicPolicyUpdate buffered.
TopicPolicyListenerWrapper initialization phase took too long. "
+ "completeInitialization should have been called
to complete the phase.");
lastIntervalLogged = warningLogIntervalCount;
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java
index 6cc66f2c18e..8377162a7fe 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java
@@ -67,7 +67,6 @@ import org.apache.pulsar.broker.service.SubscriptionOption;
import org.apache.pulsar.broker.service.Topic;
import org.apache.pulsar.broker.service.TopicAttributes;
import org.apache.pulsar.broker.service.TopicPolicyListener;
-import org.apache.pulsar.broker.service.TopicPolicyListenerWrapper;
import org.apache.pulsar.broker.service.TransportCnx;
import
org.apache.pulsar.broker.service.schema.exceptions.IncompatibleSchemaException;
import
org.apache.pulsar.broker.service.schema.exceptions.NotExistSchemaException;
@@ -131,9 +130,6 @@ public class NonPersistentTopic extends AbstractTopic
implements Topic, TopicPol
TOPIC_ATTRIBUTES_FIELD_UPDATER =
AtomicReferenceFieldUpdater.newUpdater(
NonPersistentTopic.class, TopicAttributes.class,
"topicAttributes");
- // prevents race conditions in topic policy initialization
- private final TopicPolicyListenerWrapper topicPolicyListener = new
TopicPolicyListenerWrapper(this);
-
private static class TopicStats {
public double averageMsgSize;
public double aggMsgRateIn;
@@ -161,7 +157,6 @@ public class NonPersistentTopic extends AbstractTopic
implements Topic, TopicPol
super(topic, brokerService);
this.log = LOG.with().ctx(super.log).build();
this.isFenced = false;
- registerTopicPolicyListener();
}
private CompletableFuture<Void> updateClusterMigrated() {
@@ -190,13 +185,15 @@ public class NonPersistentTopic extends AbstractTopic
implements Topic, TopicPol
updateResourceGroupLimiter(policies);
return updateClusterMigrated();
}, getPoliciesNotifyThread())
- // Complete the topic-policy listener wrapper so buffered and
future topic-level policy
- // updates are forwarded to this topic. Without this the
wrapper stays uninitialized forever
- // and all topic-level policy updates are silently dropped.
Unlike PersistentTopic,
- // non-persistent topics don't load initial topic policies
(matching the previous behavior),
- // so the loaded values are passed as null.
- .thenRunAsync(() ->
topicPolicyListener.completeInitialization(null, null),
- getPoliciesNotifyThread());
+ // Load the topic's initial policies (global and local) and
register the policy listener, so a
+ // non-persistent topic applies its own policies on load, the
same as a persistent topic does.
+ .thenCompose(ignore -> initTopicPolicy())
+ // a failure to load the initial topic policies must not fail
topic loading.
+ .exceptionally(ex -> {
+ log.warn().attr("topic", topic).exception(ex)
+ .log("Error loading topic policies during
initialization. Ignoring the failure.");
+ return null;
+ });
}
@Override
@@ -1348,8 +1345,4 @@ public class NonPersistentTopic extends AbstractTopic
implements Topic, TopicPol
old -> old != null ? old : new
TopicAttributes(TopicName.get(topic)));
}
- @Override
- public TopicPolicyListener getTopicPolicyListener() {
- return topicPolicyListener;
- }
}
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
index c2f656c776d..5eba52c8da5 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
@@ -59,7 +59,6 @@ import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.atomic.AtomicReferenceFieldUpdater;
import java.util.function.BiConsumer;
import java.util.function.BiFunction;
-import java.util.function.Function;
import lombok.Getter;
import lombok.Value;
import org.apache.bookkeeper.client.BKException.BKNoSuchLedgerExistsException;
@@ -136,8 +135,6 @@ import org.apache.pulsar.broker.service.Subscription;
import org.apache.pulsar.broker.service.SubscriptionOption;
import org.apache.pulsar.broker.service.Topic;
import org.apache.pulsar.broker.service.TopicPoliciesService;
-import org.apache.pulsar.broker.service.TopicPolicyListener;
-import org.apache.pulsar.broker.service.TopicPolicyListenerWrapper;
import org.apache.pulsar.broker.service.TransportCnx;
import org.apache.pulsar.broker.service.schema.BookkeeperSchemaStorage;
import
org.apache.pulsar.broker.service.schema.exceptions.IncompatibleSchemaException;
@@ -300,9 +297,6 @@ public class PersistentTopic extends AbstractTopic
implements Topic, AddEntryCal
@Getter
private volatile long lastMaxReadPositionMovedForwardTimestamp = 0;
- // prevents race conditions in topic policy initialization
- private final TopicPolicyListenerWrapper topicPolicyListener = new
TopicPolicyListenerWrapper(this);
-
@Getter
private final ExecutorService orderedExecutor;
@@ -515,8 +509,10 @@ public class PersistentTopic extends AbstractTopic
implements Topic, AddEntryCal
.thenCompose(ignore -> removeOrphanReplicationCursors())
.exceptionally(ex -> {
log.warn()
+ .attr("topic", topic)
.exceptionMessage(ex)
- .log("Error getting policies and
isEncryptionRequired will be set to false");
+ .log("Error loading topic policies during
initialization. Ignoring the failure. "
+ + "isEncryptionRequired will be set to
false.");
isEncryptionRequired = false;
return null;
}));
@@ -4878,50 +4874,6 @@ public class PersistentTopic extends AbstractTopic
implements Topic, AddEntryCal
});
}
- protected CompletableFuture<Void> initTopicPolicy() {
- final var topicPoliciesService =
brokerService.pulsar().getTopicPoliciesService();
- final var partitionedTopicName =
TopicName.getPartitionedTopicName(topic);
-
- return
topicPoliciesService.registerListenerAsync(partitionedTopicName,
topicPolicyListener)
- .thenCompose(registered -> {
- if (!registered) {
- return CompletableFuture.completedFuture(null);
- }
- if (ExtensibleLoadManagerImpl.isInternalTopic(topic)) {
- // Internal topics don't load topic-level policies,
but the listener wrapper must
- // still be initialized so any buffered/future updates
are forwarded to the topic
- // instead of being silently dropped.
- return CompletableFuture.runAsync(
- () ->
topicPolicyListener.completeInitialization(null, null),
- getPoliciesNotifyThread());
- }
- // future for fetching global topic policies
- CompletableFuture<Optional<TopicPolicies>>
globalPoliciesFuture =
-
topicPoliciesService.getTopicPoliciesAsync(partitionedTopicName,
- TopicPoliciesService.GetType.GLOBAL_ONLY);
- // future for fetching local topic policies
- CompletableFuture<Optional<TopicPolicies>>
localPoliciesFuture =
-
topicPoliciesService.getTopicPoliciesAsync(partitionedTopicName,
- TopicPoliciesService.GetType.LOCAL_ONLY);
- CompletableFuture<Void> initialPoliciesFuture =
-
globalPoliciesFuture.thenCombine(localPoliciesFuture, (global, local) -> {
- // finally update the topic policies with the
latest value or loaded value
- return CompletableFuture.runAsync(() ->
-
topicPolicyListener.completeInitialization(global.orElse(null),
- local.orElse(null)),
- getPoliciesNotifyThread());
- }).thenCompose(Function.identity());
- return initialPoliciesFuture.exceptionallyCompose(ex ->
- // The topic load path logs and continues when
initial policy loading fails. Make sure the
- // already-registered wrapper is not left
buffering future live updates forever.
- CompletableFuture.runAsync(
- () ->
topicPolicyListener.completeInitialization(null, null),
- getPoliciesNotifyThread())
- .thenCompose(__ -> FutureUtil.failedFuture(
-
FutureUtil.unwrapCompletionException(ex))));
- });
- }
-
@VisibleForTesting
public MessageDeduplication getMessageDeduplication() {
return messageDeduplication;
@@ -5097,8 +5049,4 @@ public class PersistentTopic extends AbstractTopic
implements Topic, AddEntryCal
return future;
}
- @Override
- public TopicPolicyListener getTopicPolicyListener() {
- return topicPolicyListener;
- }
}
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/MetadataStoreTopicPoliciesTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/MetadataStoreTopicPoliciesTest.java
index e7fefa16497..b6927cd5ee8 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/MetadataStoreTopicPoliciesTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/MetadataStoreTopicPoliciesTest.java
@@ -44,6 +44,14 @@ public class MetadataStoreTopicPoliciesTest extends
TopicPoliciesTest {
// Not applicable to MetadataStoreTopicPoliciesService.
}
+ @Test(enabled = false)
+ @Override
+ public void testNonPersistentTopicAppliesTopicPolicyOnLoad() throws
Exception {
+ // This test is specific to SystemTopicBasedTopicPoliciesService
(casts the service and uses
+ // getPoliciesCacheInit). The non-persistent load-path behavior itself
is backend-agnostic and is
+ // covered against the default SystemTopicBasedTopicPoliciesService in
TopicPoliciesTest.
+ }
+
@Test(enabled = false)
@Override
public void testSystemTopicShouldBeCompacted() throws Exception {
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java
index 17ae9d343e4..f1d5ad47547 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java
@@ -251,6 +251,73 @@ public class TopicPoliciesTest extends
MockedPulsarServiceBaseTest {
assertEquals(topic1.getHierarchyTopicPolicies().getMaxSubscriptionsPerTopic().get(),
Integer.valueOf(10));
}
+ @Test
+ public void testNonPersistentTopicAppliesTopicPolicyOnLoad() throws
Exception {
+ // Non-persistent topics now load and apply their own topic policies
on load (like persistent topics), so a
+ // freshly loaded non-persistent topic must already reflect its
topic-level policy without waiting for a
+ // namespace-wide broadcast. Before this change a non-persistent topic
never applied its policies on load.
+ TopicName topicName = TopicName.get(
+ TopicDomain.non_persistent.value(),
+ NamespaceName.get(myNamespace),
+ "test-np-" + UUID.randomUUID()
+ );
+ String topic = topicName.toString();
+
+ SystemTopicBasedTopicPoliciesService policyService =
+ (SystemTopicBasedTopicPoliciesService)
pulsar.getTopicPoliciesService();
+
+ admin.topics().createNonPartitionedTopic(topic);
+ admin.topicPolicies().setMaxSubscriptionsPerTopicAsync(topic,
10).get();
+
+ //wait until topic loaded with right policy value.
+ Awaitility.await().untilAsserted(() -> {
+ AbstractTopic loaded = (AbstractTopic)
pulsar.getBrokerService().getTopic(topic, true).get().get();
+
assertEquals(loaded.getHierarchyTopicPolicies().getMaxSubscriptionsPerTopic().get(),
Integer.valueOf(10));
+ });
+
+ //unload the topic
+
pulsar.getNamespaceService().unloadNamespaceBundle(pulsar.getNamespaceService().getBundle(topicName)).get();
+ assertFalse(pulsar.getBrokerService().getTopics().containsKey(topic));
+
+ //re-own the namespace bundle without loading the topic
+ log.info().attr("lookup",
admin.lookups().lookupTopic(topic)).log("lookup");
+
assertTrue(pulsar.getBrokerService().isTopicNsOwnedByBrokerAsync(topicName).join());
+ assertFalse(pulsar.getBrokerService().getTopics().containsKey(topic));
+ //make sure namespace policy reader is fully started.
+ Awaitility.await().untilAsserted(() ->
+
assertTrue(policyService.getPoliciesCacheInit(topicName.getNamespaceObject()).isDone()));
+
+ //load the topic: it must already reflect the topic policy, proving it
was applied on load, not via a
+ //later broadcast.
+ AbstractTopic loaded = (AbstractTopic)
pulsar.getBrokerService().getTopic(topic, true).get().get();
+
assertEquals(loaded.getHierarchyTopicPolicies().getMaxSubscriptionsPerTopic().get(),
Integer.valueOf(10));
+ }
+
+ @Test
+ public void testGlobalPolicyUpdateDoesNotClearLocalSubscriptionPolicy()
throws Exception {
+ // A global topic policy carries no subscription-level overrides by
default (an empty subscriptionPolicies
+ // map). Because AbstractTopic keeps the local and global
per-subscription policies separately and merges them
+ // with local precedence, applying such a global policy must not clear
the local per-subscription dispatch-rate
+ // policy -- which the previous direct assignment would do under the
local-before-global ordering.
+ final String topic = "persistent://" + myNamespace +
"/test-sub-policy-merge-" + UUID.randomUUID();
+ final String subName = "sub-1";
+ admin.topics().createNonPartitionedTopic(topic);
+ admin.topics().createSubscription(topic, subName, MessageId.earliest);
+
+ DispatchRate localRate = DispatchRateImpl.builder()
+
.dispatchThrottlingRateInMsg(100).dispatchThrottlingRateInByte(2048).ratePeriodInSecond(1).build();
+ admin.topicPolicies().setSubscriptionDispatchRate(topic, subName,
localRate);
+
+ AbstractTopic topicRef = (AbstractTopic)
pulsar.getBrokerService().getTopic(topic, false).get().orElseThrow();
+ Awaitility.await().untilAsserted(() ->
+
assertEquals(topicRef.getSubscriptionDispatchRate(subName).getDispatchThrottlingRateInMsg(),
100));
+
+ // Simulate a global topic-policy update that carries no
subscription-level overrides.
+ topicRef.onUpdate(TopicPolicies.builder().isGlobal(true).build());
+
+
assertEquals(topicRef.getSubscriptionDispatchRate(subName).getDispatchThrottlingRateInMsg(),
100);
+ }
+
@Test
public void testSetSizeBasedBacklogQuota() throws Exception {
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java
index 96ccdee50cc..90e0db5a2f4 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java
@@ -37,6 +37,7 @@ import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@@ -45,6 +46,7 @@ import java.util.concurrent.TimeoutException;
import lombok.Cleanup;
import lombok.CustomLog;
import org.apache.commons.lang3.reflect.FieldUtils;
+import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest;
import org.apache.pulsar.broker.systopic.SystemTopicClient;
import org.apache.pulsar.client.admin.PulsarAdminException;
@@ -53,6 +55,7 @@ import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.common.events.PulsarEvent;
import org.apache.pulsar.common.naming.NamespaceName;
+import org.apache.pulsar.common.naming.SystemTopicNames;
import org.apache.pulsar.common.naming.TopicName;
import org.apache.pulsar.common.policies.data.ClusterData;
import org.apache.pulsar.common.policies.data.TenantInfoImpl;
@@ -836,4 +839,64 @@ public class SystemTopicBasedTopicPoliciesServiceTest
extends MockedPulsarServic
assertSame(reloadInitFuture,
spyService.getPoliciesCacheInit(namespace));
Mockito.verify(reloadReader, Mockito.never()).closeAsync();
}
+
+ @Test
+ public void
testReplayTopicPolicyListenersNotifiesOnlyNamespaceScopedLocalAndGlobalPolicies()
throws Exception {
+ SystemTopicBasedTopicPoliciesService service =
+ (SystemTopicBasedTopicPoliciesService)
pulsar.getTopicPoliciesService();
+ final NamespaceName namespaceA = NamespaceName.get(NAMESPACE1);
+ final NamespaceName namespaceB = NamespaceName.get(NAMESPACE2);
+ final TopicName localTopicA = TopicName.get("persistent", namespaceA,
"replay-local-a");
+ final TopicName globalTopicA = TopicName.get("persistent", namespaceA,
"replay-global-a");
+ final TopicName topicB = TopicName.get("persistent", namespaceB,
"replay-b");
+
+ // Seed the caches: namespace A has one topic with a cached local
policy and another with a cached global
+ // policy; namespace B has a topic with a cached local policy that
must not be replayed for namespace A.
+ service.policiesCache.put(localTopicA,
TopicPolicies.builder().isGlobal(false).build());
+ service.globalPoliciesCache.put(globalTopicA,
TopicPolicies.builder().isGlobal(true).build());
+ service.policiesCache.put(topicB,
TopicPolicies.builder().isGlobal(false).build());
+
+ final Map<TopicName, List<TopicPolicies>> received = new
ConcurrentHashMap<>();
+ for (TopicName topicName : List.of(localTopicA, globalTopicA, topicB))
{
+ final List<TopicPolicies> updates = new CopyOnWriteArrayList<>();
+ received.put(topicName, updates);
+ service.registerListenerAsync(topicName, updates::add).get();
+ }
+
+ service.replayTopicPolicyListeners(namespaceA).get(30,
TimeUnit.SECONDS);
+
+ // Only namespace A's topics are notified, once each, and both the
local and the global cache are replayed.
+ Assertions.assertThat(received.get(localTopicA)).hasSize(1);
+ Assertions.assertThat(received.get(globalTopicA)).hasSize(1);
+ // Namespace B is left untouched. The pre-fix code iterated the whole
cache and replayed every namespace.
+ Assertions.assertThat(received.get(topicB)).isEmpty();
+ }
+
+ @Test
+ public void testTopicPolicyListenerReplayDisabledByDefault() {
+ Assertions.assertThat(new
ServiceConfiguration().isTopicPolicyListenerReplayEnabled()).isFalse();
+ }
+
+ @Test
+ public void testChangeEventsTopicPolicyLoadDoesNotRecurse() throws
Exception {
+ SystemTopicBasedTopicPoliciesService service =
+ (SystemTopicBasedTopicPoliciesService)
pulsar.getTopicPoliciesService();
+ final String namespaceStr = "system-topic/change-events-recursion";
+ admin.namespaces().createNamespace(namespaceStr);
+ final NamespaceName namespace = NamespaceName.get(namespaceStr);
+ final TopicName changeEvents =
+ TopicName.get("persistent", namespace,
SystemTopicNames.NAMESPACE_EVENTS_LOCAL_NAME);
+
+ // The __change_events system topic must not load topic-level
policies: that would create a policy-cache
+ // reader on __change_events while __change_events is still loading --
a recursive, deadlocking dependency.
+ // isSelf() guards getTopicPoliciesAsync so it returns empty for the
__change_events topic without ever
+ // creating a reader. AbstractTopic#initTopicPolicy (now called for
persistent AND non-persistent topics)
+ // relies on this short-circuit when a __change_events topic itself is
loaded.
+ Assertions.assertThat(service.getTopicPoliciesAsync(changeEvents,
TopicPoliciesService.GetType.LOCAL_ONLY)
+ .get(30, TimeUnit.SECONDS)).isEmpty();
+ Assertions.assertThat(service.getTopicPoliciesAsync(changeEvents,
TopicPoliciesService.GetType.GLOBAL_ONLY)
+ .get(30, TimeUnit.SECONDS)).isEmpty();
+ // No policy-cache reader was created as a side effect, which is what
would recurse.
+
Assertions.assertThat(service.getReaderCaches()).doesNotContainKey(namespace);
+ }
}
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicPolicyListenerWrapperTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicPolicyListenerWrapperTest.java
index 0e755359274..9d32342aab5 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicPolicyListenerWrapperTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicPolicyListenerWrapperTest.java
@@ -56,15 +56,15 @@ public class TopicPolicyListenerWrapperTest {
assertThat(real.updates).isEmpty();
// On completion, the buffered local value wins over the loaded local
value; the loaded global value
- // is applied since none was buffered.
+ // is applied since none was buffered. The local policy is emitted
before the global one.
TopicPolicies loadedGlobal = globalPolicies();
wrapper.completeInitialization(loadedGlobal, localPolicies());
- assertThat(real.updates).containsExactly(loadedGlobal, bufferedLocal);
+ assertThat(real.updates).containsExactly(bufferedLocal, loadedGlobal);
// After initialization, updates are forwarded immediately.
TopicPolicies liveUpdate = localPolicies();
wrapper.onUpdate(liveUpdate);
- assertThat(real.updates).containsExactly(loadedGlobal, bufferedLocal,
liveUpdate);
+ assertThat(real.updates).containsExactly(bufferedLocal, loadedGlobal,
liveUpdate);
}
@Test
@@ -78,7 +78,7 @@ public class TopicPolicyListenerWrapperTest {
wrapper.onUpdate(bufferedLocal);
wrapper.completeInitialization(globalPolicies(), localPolicies());
- assertThat(real.updates).containsExactly(bufferedGlobal,
bufferedLocal);
+ assertThat(real.updates).containsExactly(bufferedLocal,
bufferedGlobal);
}
@Test
@@ -86,10 +86,12 @@ public class TopicPolicyListenerWrapperTest {
RecordingListener real = new RecordingListener();
TopicPolicyListenerWrapper wrapper = new
TopicPolicyListenerWrapper(real);
+ // The local policy is emitted before the global policy, so a local
topic policy takes precedence over a
+ // global one once both have been applied.
TopicPolicies loadedGlobal = globalPolicies();
TopicPolicies loadedLocal = localPolicies();
wrapper.completeInitialization(loadedGlobal, loadedLocal);
- assertThat(real.updates).containsExactly(loadedGlobal, loadedLocal);
+ assertThat(real.updates).containsExactly(loadedLocal, loadedGlobal);
}
@Test
@@ -118,7 +120,93 @@ public class TopicPolicyListenerWrapperTest {
wrapper.onUpdate(newerGlobal);
wrapper.completeInitialization(globalPolicies(), localPolicies());
- // Global: the newer update wins; Local: the delete (null) wins over
the loaded local value.
- assertThat(real.updates).containsExactly(newerGlobal, null);
+ // Local (emitted first): the delete (null) wins over the loaded local
value; Global: the newer update wins.
+ assertThat(real.updates).containsExactly(null, newerGlobal);
+ }
+
+ @Test
+ public void shouldNotEmitLocalScopeWhenNoLocalPolicyExists() {
+ RecordingListener real = new RecordingListener();
+ TopicPolicyListenerWrapper wrapper = new
TopicPolicyListenerWrapper(real);
+
+ // With no local policy, only the global policy is emitted; no local
onUpdate happens, so the
+ // local-before-global ordering leaves behavior unchanged for topics
that only have a global policy.
+ TopicPolicies loadedGlobal = globalPolicies();
+ wrapper.completeInitialization(loadedGlobal, null);
+ assertThat(real.updates).containsExactly(loadedGlobal);
+ }
+
+ @Test
+ public void shouldIgnoreCompleteInitializationAfterAlreadyCompleted() {
+ RecordingListener real = new RecordingListener();
+ TopicPolicyListenerWrapper wrapper = new
TopicPolicyListenerWrapper(real);
+ wrapper.startInitialization();
+
+ TopicPolicies loadedLocal = localPolicies();
+ wrapper.completeInitialization(null, loadedLocal);
+ assertThat(real.updates).containsExactly(loadedLocal);
+
+ // Completing again (e.g. from initTopicPolicy's terminal handler)
must be a no-op and must not re-emit.
+ wrapper.completeInitialization(globalPolicies(), localPolicies());
+ wrapper.completeInitializationUnlessAlreadyCompleted();
+ assertThat(real.updates).containsExactly(loadedLocal);
+ }
+
+ @Test
+ public void
shouldEmitBufferedValueAndForwardLiveUpdatesWhenCompletedWithoutLoadedPolicies()
{
+ RecordingListener real = new RecordingListener();
+ TopicPolicyListenerWrapper wrapper = new
TopicPolicyListenerWrapper(real);
+ wrapper.startInitialization();
+
+ // A policy update arrives while initializing and is buffered.
+ TopicPolicies buffered = localPolicies();
+ wrapper.onUpdate(buffered);
+ assertThat(real.updates).isEmpty();
+
+ // initTopicPolicy's terminal handler completes initialization with no
loaded policies -- the path taken after
+ // a policy-load error or when the listener was not registered. The
buffered value is emitted and the wrapper
+ // leaves the buffering phase.
+ wrapper.completeInitializationUnlessAlreadyCompleted();
+ assertThat(real.updates).containsExactly(buffered);
+
+ // Subsequent live updates now flow through instead of being dropped.
+ TopicPolicies live = globalPolicies();
+ wrapper.onUpdate(live);
+ assertThat(real.updates).containsExactly(buffered, live);
+ }
+
+ @Test
+ public void shouldForwardLiveUpdatesAfterCompletingWithNothingBuffered() {
+ RecordingListener real = new RecordingListener();
+ TopicPolicyListenerWrapper wrapper = new
TopicPolicyListenerWrapper(real);
+ wrapper.startInitialization();
+
+ // Completed with nothing buffered and no loaded policies (e.g. after
a failed load): nothing is emitted, but
+ // the wrapper still leaves the buffering phase so later live updates
are forwarded rather than dropped.
+ wrapper.completeInitializationUnlessAlreadyCompleted();
+ assertThat(real.updates).isEmpty();
+
+ TopicPolicies live = localPolicies();
+ wrapper.onUpdate(live);
+ assertThat(real.updates).containsExactly(live);
+ }
+
+ @Test
+ public void
shouldRebufferAndReapplyAfterStartInitializationIsCalledAgain() {
+ RecordingListener real = new RecordingListener();
+ TopicPolicyListenerWrapper wrapper = new
TopicPolicyListenerWrapper(real);
+ wrapper.startInitialization();
+ TopicPolicies firstLocal = localPolicies();
+ wrapper.completeInitialization(null, firstLocal);
+ assertThat(real.updates).containsExactly(firstLocal);
+
+ // A new initialization phase (e.g. re-running initTopicPolicy):
updates are buffered again until it completes,
+ // and a value buffered during the phase is applied on completion.
+ wrapper.startInitialization();
+ TopicPolicies bufferedGlobal = globalPolicies();
+ wrapper.onUpdate(bufferedGlobal);
+ assertThat(real.updates).containsExactly(firstLocal);
+ wrapper.completeInitialization(null, null);
+ assertThat(real.updates).containsExactly(firstLocal, bufferedGlobal);
}
}
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java
index 535b0ae75ad..7c9ce90180a 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java
@@ -876,6 +876,13 @@ public class PersistentTopicTest extends BrokerTestBase {
public void onUpdate(TopicPolicies policies) {
receivedUpdates.add(policies);
}
+
+ // initTopicPolicy() moved to AbstractTopic (a different package),
so widen it to public here to keep
+ // this same-package test able to invoke it directly.
+ @Override
+ public CompletableFuture<Void> initTopicPolicy() {
+ return super.initTopicPolicy();
+ }
}
final String topic =
"persistent://prop/ns-abc/testTopicPolicyInitFailure-" + UUID.randomUUID();