Skip to content

[feat][broker] PIP-469: Legacy-aware topic policies backend routing and metadata-store topic policies - #25707

Merged
BewareMyPower merged 20 commits into
apache:masterfrom
BewareMyPower:bewaremypower/pip-469-impl
May 18, 2026
Merged

BewareMyPower merged 20 commits into
apache:masterfrom
BewareMyPower:bewaremypower/pip-469-impl

Conversation

@BewareMyPower

Copy link
Copy Markdown
Contributor

pip: #25547

This PR reuses TopicPoliciesTest to test most functionality of the new metadata store based implementation, so it also introduces improvements on TopicPoliciesTest, which requires too much time to run before (100+ tests, where each test calls internalSetup and internalCleanup).

@BewareMyPower BewareMyPower self-assigned this May 7, 2026
@BewareMyPower BewareMyPower added this to the 5.0.0-M1 milestone May 7, 2026
@BewareMyPower
BewareMyPower requested a review from Copilot May 7, 2026 05:22

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Note

Copilot was unable to run its full agentic suite in this review.

Implements PIP-469 topic-policies backend routing to preserve legacy __change_events behavior while enabling a new metadata-store-backed topic policies implementation, and refactors existing topic policies tests to run significantly faster.

Changes:

  • Add LegacyAwareTopicPoliciesService to route per-namespace operations to either legacy system-topic or configured backend.
  • Introduce MetadataStoreTopicPoliciesService implementation backed by metadata stores, with listener support.
  • Refactor TopicPoliciesTest lifecycle to reduce repeated broker setup/teardown and add a derived test suite for the metadata-store backend.

Reviewed changes

Copilot reviewed 9 out of 9 changed files in this pull request and generated 6 comments.

Show a summary per file
File Description
pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicPolicyTestUtils.java Extends bypass-cache helper to work with legacy-aware routing and the metadata-store backend.
pulsar-broker/src/test/java/org/apache/pulsar/broker/service/LegacyAwareTopicPoliciesServiceTest.java Adds upgrade/downgrade and listener behavior coverage for legacy-aware routing and metadata-store storage.
pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java Refactors test lifecycle to reduce runtime and adapts tests to multiple policies backends.
pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/MetadataStoreTopicPoliciesTest.java Reuses TopicPoliciesTest against the metadata-store backend, disabling inapplicable system-topic tests.
pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java Loosens visibility on bundle-ownership cleanup to enable reuse/integration.
pulsar-broker/src/main/java/org/apache/pulsar/broker/service/MetadataStoreTopicPoliciesService.java Adds new metadata-store-backed implementation for topic policies with store notifications and listeners.
pulsar-broker/src/main/java/org/apache/pulsar/broker/service/LegacyAwareTopicPoliciesService.java Adds per-namespace routing layer between legacy system-topic backend and configured backend.
pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java Wraps configured topic policies service with legacy-aware routing when system topics are enabled.
pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java Updates configuration documentation to describe built-in topic policies service implementations and legacy behavior.

Comment thread pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java Outdated
Comment thread pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java Outdated
@BewareMyPower
BewareMyPower marked this pull request as draft May 7, 2026 09:05
@BewareMyPower

Copy link
Copy Markdown
Contributor Author

We might need to consider whether to change the registerListener API because we cannot register listener for both LegacyAwareTopicPoliciesService and SystemTopicBasedTopicPoliciesService.

@BewareMyPower
BewareMyPower marked this pull request as ready for review May 8, 2026 08:06

@dao-jun dao-jun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM, just 2 minor comments

  1. Metadata-store listeners are never unregistered on close() (resource leak)

MetadataStoreTopicPoliciesService.start() line 354-355:

localStore.registerListener(notification -> handleNotification(notification, false));
configurationStore.registerListener(notification -> handleNotification(notification, true));

The return value of registerListener (a registration handle) is discarded. close() only clears the local listeners map and invalidates caches, but never unregisters these metadata-store listeners. After close(), the service can still receive and process notifications (the closed flag guards most paths, but the listener
lambdas themselves are still held by the metadata store). Consider capturing and closing the registration handles in close().

  1. Path structure deviates from PIP-469 design document

The PIP design doc specifies:

  • Global: /admin/topic-policies/{tenant}/{namespace}/{domain}/{encodedTopic}
  • Local: /admin/local-policies/topic-policies/{tenant}/{namespace}/{domain}/{encodedTopic}

The implementation uses:

  • Global: /admin/topic-policies/global/{tenant}/{namespace}/{domain}/{encodedTopic}
  • Local: /admin/topic-policies/local/{tenant}/{namespace}/{domain}/{encodedTopic}

The code's approach is arguably cleaner (both scopes under one root), but the discrepancy from the approved design should be explicitly called out and either the PIP or the code should be updated for consistency.

@BewareMyPower
BewareMyPower marked this pull request as draft May 14, 2026 16:00
@BewareMyPower

Copy link
Copy Markdown
Contributor Author

Let me address these two comments before merging

@BewareMyPower

Copy link
Copy Markdown
Contributor Author

@dao-jun I've updated the document for why the implementation adopts a different path.

The return value of registerListener (a registration handle) is discarded. close() only clears the local listeners map and invalidates caches, but never unregisters these metadata-store listeners. After close(), the service can still receive and process notifications (the closed flag guards most paths, but the listener
lambdas themselves are still held by the metadata store). Consider capturing and closing the registration handles in close().

Regarding the previous comment, it seems to be wrong. Firstly, registerListener does not return any handle, see

void registerListener(Consumer<Notification> listener);

Secondly, it follows the similar pattern of

pulsar.getLocalMetadataStore().registerListener(this::handleMetadataChanges);
if (pulsar.getConfigurationMetadataStore() != pulsar.getLocalMetadataStore()) {
pulsar.getConfigurationMetadataStore().registerListener(this::handleMetadataChanges);
}

Though this PR has registered two listeners with a parameter that indicates if the policy is global so that we don't need to check if the notification path starts with GLOBAL_POLICIES_ROOT or LOCAL_POLICIES_ROOT.

Could you double check the review comments if they are generated from LLM?

@BewareMyPower
BewareMyPower marked this pull request as ready for review May 18, 2026 03:21
@shibd
shibd self-requested a review May 18, 2026 05:14

@nodece nodece left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@BewareMyPower
BewareMyPower merged commit 8652efa into apache:master May 18, 2026
80 of 82 checks passed
@BewareMyPower

Copy link
Copy Markdown
Contributor Author

@nodece Let me merge it first because this topic is out of the scope of the topic policies service itself. I agree that your comments [1] and [2] make sense, but it's more like an issue that should be resolved by structured concurrency, which needs careful design.

It would be better add a common abstraction for listeners and skip all of them after closing the topic policies service.

[1] #25707 (comment)
[2] #25707 (comment)

@BewareMyPower
BewareMyPower deleted the bewaremypower/pip-469-impl branch May 18, 2026 11:31
@BewareMyPower

Copy link
Copy Markdown
Contributor Author

I tried to write a test to reproduce and then found that even if a listener is registered due to the race, there still might be a real issue.

    private static final CountDownLatch realCloseLatch = new CountDownLatch(1);
    private static final CountDownLatch closeLatch = new CountDownLatch(1);

    public static class MockTopicPoliciesService extends MetadataStoreTopicPoliciesService {

        public void close() {
            super.close();
            realCloseLatch.countDown();
            try {
                if (!closeLatch.await(10, TimeUnit.SECONDS)) {
                    log.info().log("Failed to close TopicPoliciesService within the timeout");
                }
            } catch (InterruptedException e) {
                log.info().log("Interrupted when closing TopicPoliciesService");
            }
        }
    }

The test can be easily written based on the mocked implementation above.

MetadataStoreTopicPoliciesService#close is a very simple operation that changes closed to true. Even if a different broker receives a topic policies update request and updates the policies, handleNotification will skip calling notifyListeners in

if (closed.get()
|| (notification.getType() != NotificationType.Created
&& notification.getType() != NotificationType.Modified
&& notification.getType() != NotificationType.Deleted)) {
return;
}

Therefore, basically these zombie listeners won't be notified.

In another case, if the CAS of the closed flag happened before listeners were notified, where asynchronous operations had been started, these listeners would be cleared and would not be notified anymore.

BewareMyPower added a commit that referenced this pull request Jun 8, 2026
…nd metadata-store topic policies (#25707)

(cherry picked from commit 8652efa)
BewareMyPower added a commit that referenced this pull request Jun 8, 2026
…nd metadata-store topic policies (#25707)

(cherry picked from commit 8652efa)
(cherry picked from commit 712372f)
priyanshu-ctds pushed a commit to datastax/pulsar that referenced this pull request Jun 8, 2026
…nd metadata-store topic policies (apache#25707)

(cherry picked from commit 8652efa)
(cherry picked from commit 712372f)
(cherry picked from commit c333504)
priyanshu-ctds pushed a commit to datastax/pulsar that referenced this pull request Jun 9, 2026
…nd metadata-store topic policies (apache#25707)

(cherry picked from commit 8652efa)
(cherry picked from commit 712372f)
(cherry picked from commit c333504)
nodece pushed a commit to ascentstream/pulsar that referenced this pull request Aug 28, 2026
…nd metadata-store topic policies (apache#25707)

(cherry picked from commit 8652efa)
(cherry picked from commit 712372f)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants