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 decc80ff036 [fix][broker] Don't let a closing topic-policies reader 
abort a concurrent cache-init reload (#26132)
decc80ff036 is described below

commit decc80ff0369c0387c6cd02745ebe993dbf71ff1
Author: Lari Hotari <[email protected]>
AuthorDate: Thu Jul 2 03:02:34 2026 +0300

    [fix][broker] Don't let a closing topic-policies reader abort a concurrent 
cache-init reload (#26132)
---
 .../SystemTopicBasedTopicPoliciesService.java      | 108 ++++++++-------------
 .../SystemTopicBasedTopicPoliciesServiceTest.java  |  93 ++++++++++++++----
 2 files changed, 116 insertions(+), 85 deletions(-)

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 398c7c813e7..f565eb026f4 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
@@ -466,7 +466,7 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
                     // The cached writer will be closed when an exception 
happens
                     // This is potentially not a great idea since we should be 
able to rely on the Pulsar client's
                     // behavior for restoring a Producer after a failure.
-                    
writerCaches.synchronous().invalidate(topicName.getNamespaceObject());
+                    cleanWriterCache(topicName.getNamespaceObject());
                     throw FutureUtil.wrapToCompletionException(t);
                 });
     }
@@ -662,7 +662,7 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
                                     initPolicesCache(reader, stageFuture);
                                     return stageFuture
                                             // Read policies in background
-                                            .thenAccept(__ -> 
readMorePoliciesAsync(reader));
+                                            .thenAccept(__ -> 
readMorePoliciesAsync(reader, initNamespacePolicyFuture));
                                 }).thenApply(__ -> {
                                     initNamespacePolicyFuture.complete(null);
                                     return null;
@@ -702,13 +702,7 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
     }
 
     private CompletableFuture<SystemTopicClient.Reader<PulsarEvent>> 
newReader(NamespaceName ns) {
-        return readerCaches.compute(ns, (__, existingFuture) -> {
-            if (existingFuture == null) {
-                return createSystemTopicClient(ns);
-            }
-
-            return existingFuture;
-        });
+        return readerCaches.computeIfAbsent(ns, __ -> 
createSystemTopicClient(ns));
     }
 
     protected CompletableFuture<SystemTopicClient.Reader<PulsarEvent>> 
createSystemTopicClient(
@@ -734,7 +728,7 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
         }
         AtomicInteger bundlesCount = 
ownedBundlesCountPerNamespace.get(namespace);
         if (bundlesCount == null || bundlesCount.decrementAndGet() <= 0) {
-            cleanPoliciesCacheInitMap(namespace, true);
+            cleanPoliciesCacheInitMap(namespace);
             cleanWriterCache(namespace);
             cleanOwnedBundlesCount(namespace);
         }
@@ -804,17 +798,20 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
 
     /**
      * Identity-guarded cleanup for an initialization that failed (it timed 
out, the {@code __change_events} reader
-     * could not be created, or reading the topic threw). Unlike {@link 
#cleanPoliciesCacheInitMap}, which
+     * could not be created, or reading the topic threw), or whose background 
reader was later closed (an
+     * {@code AlreadyClosedException} surfaced in {@link 
#readMorePoliciesAsync} after a namespace unload closed the
+     * reader). Unlike {@link #cleanPoliciesCacheInitMap}, which
      * removes/closes by namespace key unconditionally, this only tears down 
state that still belongs to
      * {@code initFuture}. By the time the failure is observed, a concurrent 
retry — or a namespace-bundle unload that
      * left the init future orphaned — may already own the namespace with a 
fresh future and reader; removing by key
      * would drop that newer future and close its reader, pinning the 
namespace again. Guarding on identity ensures a
      * late failure never clobbers a newer initialization.
      *
-     * @param closeReader when {@code true}, also clears the cached policies 
and closes the reader that belongs to this
-     *                    initialization; when {@code false}, only the init 
future is dropped, leaving the reader cached
-     *                    for the retry to reuse (mirrors the transient 
read-error path of
-     *                    {@link #cleanPoliciesCacheInitMap}).
+     * @param closeReader when {@code true}, also closes the reader and 
message-handler tracker that belong to this
+     *                    initialization; when {@code false}, only the init 
future is dropped, leaving the reader
+     *                    cached for the retry to reuse. The cached policies 
are intentionally left in place; they
+     *                    are cleared only when the whole namespace is 
unloaded, so this cleanup cannot race a
+     *                    concurrent re-initialization.
      */
     @VisibleForTesting
     void cleanupFailedPolicyCacheInit(@NonNull NamespaceName namespace,
@@ -823,23 +820,29 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
         // initialization, never one a concurrent retry creates immediately 
afterwards.
         CompletableFuture<SystemTopicClient.Reader<PulsarEvent>> readerFuture =
                 closeReader ? readerCaches.get(namespace) : null;
+        TopicPolicyMessageHandlerTracker tracker = 
topicPolicyMessageHandlerTrackers.get(namespace);
+
+        // Identity guard: only proceed while this initialization still owns 
the namespace's init future.
         if (!policyCacheInitMap.remove(namespace, initFuture)) {
             // Superseded by a retry or an unload; that owner is responsible 
for its own reader/state.
             return;
         }
+
         // Complete the dropped future (a no-op if the caller already 
completed it) outside any map remapping function,
         // so awaiting topic loads fail fast and retry instead of hanging 
until the broker restarts (issue #25294).
         failPendingPolicyCacheInit(namespace, initFuture);
         if (!closeReader) {
             return;
         }
-        policiesCache.entrySet().removeIf(entry -> 
Objects.equals(entry.getKey().getNamespaceObject(), namespace));
-        globalPoliciesCache.entrySet()
-                .removeIf(entry -> 
Objects.equals(entry.getKey().getNamespaceObject(), namespace));
-        TopicPolicyMessageHandlerTracker tracker = 
topicPolicyMessageHandlerTrackers.remove(namespace);
-        if (tracker != null) {
+
+        // Close the tracker captured above only if it is still the one 
installed for this namespace, so a
+        // concurrent re-initialization that installed a newer tracker is left 
untouched.
+        if (tracker != null && 
topicPolicyMessageHandlerTrackers.remove(namespace, tracker)) {
             tracker.close();
         }
+
+        // Remove and close the reader captured above only if it is still the 
current one, so a reader
+        // created by a later initialization is never closed by this stale 
cleanup.
         if (readerFuture != null && readerCaches.remove(namespace, 
readerFuture)
                 && !readerFuture.isCompletedExceptionally()) {
             readerFuture.thenCompose(SystemTopicClient.Reader::closeAsync)
@@ -853,7 +856,7 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
     private void initPolicesCache(SystemTopicClient.Reader<PulsarEvent> 
reader, CompletableFuture<Void> future) {
         if (closed.get()) {
             future.completeExceptionally(new 
BrokerServiceException(getClass().getName() + " is closed."));
-            
cleanPoliciesCacheInitMap(reader.getSystemTopic().getTopicName().getNamespaceObject(),
 true);
+            
cleanPoliciesCacheInitMap(reader.getSystemTopic().getTopicName().getNamespaceObject());
             return;
         }
         reader.hasMoreEventsAsync().whenComplete((hasMore, ex) -> {
@@ -901,13 +904,10 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
         });
     }
 
+    // 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
-    void cleanPoliciesCacheInitMap(@NonNull NamespaceName namespace, boolean 
closeReader) {
-        if (!closeReader) {
-            failPendingPolicyCacheInit(namespace, 
policyCacheInitMap.remove(namespace));
-            return;
-        }
-
+    void cleanPoliciesCacheInitMap(@NonNull NamespaceName namespace) {
         TopicPolicyMessageHandlerTracker topicPolicyMessageHandlerTracker =
                 topicPolicyMessageHandlerTrackers.remove(namespace);
         if (topicPolicyMessageHandlerTracker != null) {
@@ -959,52 +959,17 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
         ownedBundlesCountPerNamespace.remove(namespace);
     }
 
-
-    private void cleanCacheAndCloseReader(@NonNull NamespaceName namespace, 
boolean cleanOwnedBundlesCount,
-                                          boolean cleanWriterCache) {
-        if (cleanWriterCache) {
-            writerCaches.synchronous().invalidate(namespace);
-        }
-        CompletableFuture<SystemTopicClient.Reader<PulsarEvent>> readerFuture 
= readerCaches.remove(namespace);
-
-        TopicPolicyMessageHandlerTracker topicPolicyMessageHandlerTracker =
-                topicPolicyMessageHandlerTrackers.remove(namespace);
-        if (topicPolicyMessageHandlerTracker != null) {
-            topicPolicyMessageHandlerTracker.close();
-        }
-
-        if (cleanOwnedBundlesCount) {
-            ownedBundlesCountPerNamespace.remove(namespace);
-        }
-        if (readerFuture != null && !readerFuture.isCompletedExceptionally()) {
-            readerFuture.thenCompose(SystemTopicClient.Reader::closeAsync)
-                    .exceptionally(ex -> {
-                        log.warn().attr("namespace", 
namespace).exception(ex).log("Close change_event reader fail.");
-                        return null;
-                    });
-        }
-
-        policyCacheInitMap.compute(namespace, (k, v) -> {
-            policiesCache.entrySet().removeIf(entry -> 
Objects.equals(entry.getKey().getNamespaceObject(), namespace));
-            globalPoliciesCache.entrySet()
-                    .removeIf(entry -> 
Objects.equals(entry.getKey().getNamespaceObject(), namespace));
-            return null;
-        });
-    }
-
-
-
-
     /**
      * This is an async method for the background reader to continue syncing 
new messages.
      *
      * Note: You should not do any blocking call here. because it will affect
      * #{@link SystemTopicBasedTopicPoliciesService#getTopicPoliciesAsync} 
method to block loading topic.
      */
-    private void readMorePoliciesAsync(SystemTopicClient.Reader<PulsarEvent> 
reader) {
+    private void readMorePoliciesAsync(SystemTopicClient.Reader<PulsarEvent> 
reader,
+                                       CompletableFuture<Void> initFuture) {
         NamespaceName namespaceObject = 
reader.getSystemTopic().getTopicName().getNamespaceObject();
         if (closed.get()) {
-            cleanPoliciesCacheInitMap(namespaceObject, true);
+            cleanupFailedPolicyCacheInit(namespaceObject, initFuture, true);
             return;
         }
         reader.readNextAsync()
@@ -1023,16 +988,23 @@ public class SystemTopicBasedTopicPoliciesService 
implements TopicPoliciesServic
                 })
                 .whenComplete((__, ex) -> {
                     if (ex == null) {
-                        readMorePoliciesAsync(reader);
+                        readMorePoliciesAsync(reader, initFuture);
                     } else {
                         if (isAlreadyClosedException(ex)) {
                             log.info()
                                     .attr("topic", 
reader.getSystemTopic().getTopicName())
                                     .log("Closing the topic policies reader 
for");
-                            cleanPoliciesCacheInitMap(namespaceObject, true);
+                            // Tear down by init-future identity, not by 
namespace key: this reader may have been
+                            // closed by a namespace unload while a concurrent 
reload already installed a fresh
+                            // reader and init future for the same namespace 
(the close only surfaces here, on the
+                            // client executor, afterwards). A namespace-keyed 
cleanup would clobber that newer
+                            // generation and abort its init with "...aborted 
because the cached state was cleared",
+                            // failing the reloading topic. 
cleanupFailedPolicyCacheInit only tears down state that
+                            // still belongs to this initialization, so a 
superseded reader's late close is a no-op.
+                            cleanupFailedPolicyCacheInit(namespaceObject, 
initFuture, true);
                         } else {
                             log.warn().exception(ex).log("Read more topic 
polices exception, read again.");
-                            readMorePoliciesAsync(reader);
+                            readMorePoliciesAsync(reader, initFuture);
                         }
                     }
                 });
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 b6dbffede81..96ccdee50cc 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
@@ -599,17 +599,19 @@ public class SystemTopicBasedTopicPoliciesServiceTest 
extends MockedPulsarServic
         });
 
         // Cleanup must run exactly once per trigger and not repeat 
recursively (in older code it ran 3 times).
-        // Two failures are triggered here: the reader.close() above drives 
readMorePoliciesAsync into
-        // cleanPoliciesCacheInitMap (1x), and the second 
prepareInitPoliciesCacheAsync fails in initPolicesCache and
-        // is torn down by the identity-guarded cleanupFailedPolicyCacheInit 
(1x).
+        // Two failures are triggered here, and both tear down through the 
identity-guarded cleanupFailedPolicyCacheInit
+        // (2x): the reader.close() above drives readMorePoliciesAsync's 
AlreadyClosed branch into it, and the second
+        // prepareInitPoliciesCacheAsync fails in initPolicesCache and is torn 
down by it as well. The namespace-keyed
+        // cleanPoliciesCacheInitMap must not be reached from the reader-close 
path, otherwise a superseded reader could
+        // clobber a newer generation's init future.
         boolean logFound = 
testLogAppender.getEvents().stream().anyMatch(logEvent ->
                 logEvent.getMessage().toString().contains("Failed to create 
reader on __change_events topic"));
         assertFalse(logFound);
         boolean logFound2 = 
testLogAppender.getEvents().stream().anyMatch(logEvent ->
                 logEvent.getMessage().toString().contains("Failed to check the 
move events for the system topic"));
         assertTrue(logFound2);
-        verify(spyService, times(1)).cleanPoliciesCacheInitMap(any(), 
anyBoolean());
-        verify(spyService, times(1)).cleanupFailedPolicyCacheInit(any(), 
any(), anyBoolean());
+        verify(spyService, times(0)).cleanPoliciesCacheInitMap(any());
+        verify(spyService, times(2)).cleanupFailedPolicyCacheInit(any(), 
any(), anyBoolean());
 
         // make sure not occur Recursive update
         boolean logFound3 = 
testLogAppender.getEvents().stream().anyMatch(logEvent ->
@@ -678,7 +680,7 @@ public class SystemTopicBasedTopicPoliciesServiceTest 
extends MockedPulsarServic
                         || logEvent.getMessage().toString().contains("Failed 
to read event from the system topic"));
         assertFalse(logFound2);
         verify(spyService, times(1)).cleanupFailedPolicyCacheInit(any(), 
any(), anyBoolean());
-        verify(spyService, times(0)).cleanPoliciesCacheInitMap(any(), 
anyBoolean());
+        verify(spyService, times(0)).cleanPoliciesCacheInitMap(any());
     }
 
     @Test(timeOut = 60_000)
@@ -727,22 +729,16 @@ public class SystemTopicBasedTopicPoliciesServiceTest 
extends MockedPulsarServic
 
         // Dropping the cached init future (e.g. on a namespace-bundle unload) 
must complete it so the topic loads
         // awaiting it fail fast and retry, instead of hanging until the 
broker restarts (issue #25294).
-        CompletableFuture<Void> pendingWithReaderClose = new 
CompletableFuture<>();
-        service.policyCacheInitMap.put(namespace, pendingWithReaderClose);
-        service.cleanPoliciesCacheInitMap(namespace, true);
-        assertTrue(pendingWithReaderClose.isCompletedExceptionally());
-        assertNull(service.getPoliciesCacheInit(namespace));
-
-        CompletableFuture<Void> pendingWithoutReaderClose = new 
CompletableFuture<>();
-        service.policyCacheInitMap.put(namespace, pendingWithoutReaderClose);
-        service.cleanPoliciesCacheInitMap(namespace, false);
-        assertTrue(pendingWithoutReaderClose.isCompletedExceptionally());
+        CompletableFuture<Void> pendingInitFuture = new CompletableFuture<>();
+        service.policyCacheInitMap.put(namespace, pendingInitFuture);
+        service.cleanPoliciesCacheInitMap(namespace);
+        assertTrue(pendingInitFuture.isCompletedExceptionally());
         assertNull(service.getPoliciesCacheInit(namespace));
 
         // An already-completed init future must not be overwritten/disturbed.
         CompletableFuture<Void> alreadyDone = 
CompletableFuture.completedFuture(null);
         service.policyCacheInitMap.put(namespace, alreadyDone);
-        service.cleanPoliciesCacheInitMap(namespace, true);
+        service.cleanPoliciesCacheInitMap(namespace);
         assertFalse(alreadyDone.isCompletedExceptionally());
     }
 
@@ -777,4 +773,67 @@ public class SystemTopicBasedTopicPoliciesServiceTest 
extends MockedPulsarServic
         assertNull(service.getReaderCaches().get(namespace));
         Mockito.verify(newerReader, Mockito.times(1)).closeAsync();
     }
+
+    @Test
+    @SuppressWarnings("unchecked")
+    public void testClosedSupersededReaderDoesNotAbortReloadedInit() throws 
Exception {
+        // Reproduces the race behind the flaky 
AdminApi2Test.testGetInternalStatsWithProperties: a namespace unload
+        // closes the __change_events reader while a reload (e.g. getTopic 
right after unload) installs a fresh reader
+        // and init future for the same namespace. The old reader's close only 
surfaces later, on the pulsar-client
+        // executor, as an AlreadyClosedException in readMorePoliciesAsync. 
That late cleanup must NOT clobber the newer
+        // generation and abort its init future ("...aborted because the 
cached state was cleared"), which would fail
+        // the reloading topic load.
+        @Cleanup
+        TestLogAppender testLogAppender = TestLogAppender.create(log);
+
+        pulsar.getTopicPoliciesService().close();
+        SystemTopicBasedTopicPoliciesService spyService =
+                Mockito.spy(new SystemTopicBasedTopicPoliciesService(pulsar));
+        FieldUtils.writeField(pulsar, "topicPoliciesService", spyService, 
true);
+
+        final NamespaceName namespace = NamespaceName.get(NAMESPACE5);
+        admin.namespaces().createNamespace(NAMESPACE5);
+
+        // A real reader, spied so its background read loop is fully 
controllable: it reports "no more events" so the
+        // initialization completes and readMorePoliciesAsync starts, then 
parks on a read future we complete by hand.
+        SystemTopicClient.Reader<PulsarEvent> oldReader =
+                
Mockito.spy(spyService.createSystemTopicClient(namespace).get(30, 
TimeUnit.SECONDS));
+        CompletableFuture<Message<PulsarEvent>> parkedRead = new 
CompletableFuture<>();
+        
Mockito.doReturn(CompletableFuture.completedFuture(false)).when(oldReader).hasMoreEventsAsync();
+        Mockito.doReturn(parkedRead).when(oldReader).readNextAsync();
+        Mockito.doReturn(CompletableFuture.completedFuture(oldReader))
+                .when(spyService).createSystemTopicClient(namespace);
+        spyService.getReaderCaches().put(namespace, 
CompletableFuture.completedFuture(oldReader));
+
+        // Drive initialization: readMorePoliciesAsync(oldReader, <old init 
future>) is now looping, parked on
+        // parkedRead, having registered its whenComplete callback.
+        assertTrue(spyService.prepareInitPoliciesCacheAsync(namespace).get(30, 
TimeUnit.SECONDS));
+        Mockito.verify(oldReader, Mockito.atLeastOnce()).readNextAsync();
+
+        // Simulate the concurrent unload+reload having already replaced the 
generation: a fresh reader and a fresh,
+        // still-pending init future that a reloading topic is awaiting.
+        SystemTopicClient.Reader<PulsarEvent> reloadReader = 
Mockito.mock(SystemTopicClient.Reader.class);
+        
Mockito.doReturn(CompletableFuture.completedFuture(null)).when(reloadReader).closeAsync();
+        CompletableFuture<SystemTopicClient.Reader<PulsarEvent>> 
reloadReaderFuture =
+                CompletableFuture.completedFuture(reloadReader);
+        CompletableFuture<Void> reloadInitFuture = new CompletableFuture<>();
+        spyService.getReaderCaches().put(namespace, reloadReaderFuture);
+        spyService.policyCacheInitMap.put(namespace, reloadInitFuture);
+
+        // The old reader finally observes it was closed; this runs 
readMorePoliciesAsync's AlreadyClosed cleanup
+        // synchronously on this thread.
+        parkedRead.completeExceptionally(new 
PulsarClientException.AlreadyClosedException("reader is already closed"));
+
+        // The cleanup ran (it logged), but being identity-guarded on the init 
future it left the newer generation
+        // untouched. Before the fix it cleared 
readerCaches/policyCacheInitMap by namespace key and aborted the reload.
+        assertTrue(testLogAppender.getEvents().stream().anyMatch(e ->
+                e.getMessage().toString().contains("Closing the topic policies 
reader for")));
+        assertFalse("the reload's init future must not be aborted by the 
superseded reader's late close",
+                reloadInitFuture.isCompletedExceptionally());
+        assertFalse(reloadInitFuture.isDone());
+        assertSame("the reload's reader must remain cached", 
reloadReaderFuture,
+                spyService.getReaderCaches().get(namespace));
+        assertSame(reloadInitFuture, 
spyService.getPoliciesCacheInit(namespace));
+        Mockito.verify(reloadReader, Mockito.never()).closeAsync();
+    }
 }

Reply via email to