Repository navigation
[fix][broker] Don't serve topic policies from a cache whose init future has not completed - #26513
Conversation
…re has not completed getTopicPoliciesAsync awaits prepareInitPoliciesCacheAsync, hops threads and re-derives that guarantee from policyCacheInitMap as "existingFuture != null". A namespace-bundle bounce landing inside the hop wipes both policy caches and installs a new, still-loading init future, so the resumed read served the wiped cache: Optional.empty() for a topic that has policies, which a topic load then turned into the namespace retention burned into its ManagedLedgerConfig. Read the caches only from a complete, non-failed init future and retry otherwise; add a bounce-injecting test seam, a unit test and an end-to-end retention test. Assisted-by: Claude Code (Fable 5.1)
lhotari
left a comment
There was a problem hiding this comment.
The completion guard addresses the stale-generation read. Please make the regression fixture control the replacement reader until the read reaches the decision under test; its current counter does not establish that the stale-read window was exercised. The reflection cleanup is nonblocking.
Addresses the review on apache#26513. The regression fixture no longer samples the replacement generation after the fact. It now holds each read until the replacement is installed with its reader gated -- initPolicesCache calls hasMoreEventsAsync() first, so that reader has read nothing and the caches hold nothing for the namespace, the old reader having been closed by the wipe -- and releases the gate only once every read the window holds has been observed to reach its decision: either the retry, which re-enters getTopicPoliciesAsync with the same topic and GetType, or the completion of the read's own future. A read resuming from the same generation completion as the one that opened the window joins it, and one that arrives after that window has closed re-awaits the replacement generation instead of running past it. Nothing depends on reader speed or on how the common pool schedules the continuation, and the counters count decisions rather than sampled state. The gated reader's two synchronous methods throw rather than block a broker thread on the gate; the service never calls them. Unpatched, the unit test reports one window, zero retries and one read answered from the wiped cache; with the fix, four windows, four retries and none answered. The end-to-end test reports four windows covering all seven policy reads of the topic load, all seven answered and zero retries unpatched, leaving the namespace retention live; eight windows, eight retries and none answered with the fix. The two tests install the replacement service through a new typed @VisibleForTesting PulsarService#setTopicPoliciesService instead of writing the private field reflectively, so a rename or a type change is a compile error rather than a runtime one. Assisted-by: Claude Code (Fable 5.1)
|
Addressed the review in 6f142a8. Both points addressed: the fixture now gates the replacement generation's reader and releases it only once every read the window holds has been observed to retry or to answer, and both tests install the service through a new typed |
lhotari
left a comment
There was a problem hiding this comment.
LGTM. The replacement reader now stays gated until the held reads retry or complete, and both tests use the typed service setter. This addresses both review points.
…re has not completed (apache#26513)
Related: #26137, #20763
Motivation
SystemTopicBasedTopicPoliciesService.getTopicPoliciesAsynccan returnOptional.empty()for a topic that has topic-level policies, without any exception or log line, when a namespace-bundle bounce (unload followed by reload of the namespace's last bundle on the same broker) lands inside the thread hop the method performs between awaiting the namespace's policy-cache initialization and reading the cache. A topic load that hits this window builds itsManagedLedgerConfigfrom the namespace retention instead of the topic retention, and nothing repairs it afterwards. The consequence -- a topic silently enforcing its namespace retention instead of its own while the admin API keeps reporting the topic value -- is what surfaced in a production cluster; the interleaving itself is not claimed to have been captured there. The tests in this PR reproduce it without sleeps or timing tuning: they hold the replacement generation's reader until the read under test has reached its decision, so the interleaving is exercised exactly rather than merely made likely.All line numbers below refer to master at the merge base of this PR, before the change.
The mechanism
Annotated excerpt of
getTopicPoliciesAsync(pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java, lines 580-606, quoted verbatim; the annotations follow the excerpt, keyed by line number):prepareInitPoliciesCacheAsync(namespace)is a real await.policyCacheInitMapholds oneCompletableFuture<Void>per namespace -- one "generation" of that namespace's policy cache -- installed withputIfAbsent(line 649) and completed once the__change_eventsreader has drained (line 667). Theinsertedit resolves to istrueboth when this call installed the future (line 697) and when it joined an existing one (line 699); it isfalseonly when the service is closed (line 636) or the namespace is missing or deleted (line 644).thenComposeAsyncwith no executor hands the continuation toCompletableFuture's default async executor (the commonForkJoinPoolwhen its parallelism is above 1, otherwise a new thread per task), so it runs on another thread afterpreparedFuturecompletes. Between the completion ofpreparedFutureand the execution of the lambda,policyCacheInitMapcan change.existingFuture != nullasks whether some init future is present, never whether it is complete. A future installed a moment ago that has not read a single event yet passes.The interleaving
getTopicPoliciesAsync(topic, LOCAL_ONLY). Generation N of the namespace is complete, sopreparedFuturecompletes withinserted == trueand the continuation of line 582 is queued on the common pool.removeOwnedNamespaceBundleAsync(line 724) callscleanPoliciesCacheInitMap(namespace)(line 935), which, insidepolicyCacheInitMap.compute, removes everypoliciesCacheandglobalPoliciesCacheentry of the namespace and returnsnull, i.e. removes the map entry (lines 944-950).addOwnedNamespaceBundleAsync(line 610) callsprepareInitPoliciesCacheAsync(namespace), whichputIfAbsents a new, incomplete future -- generation N+1 -- and starts a new reader (lines 647-658).insertedistrueandexistingFutureis generation N+1, non-null, so line 587 takes the read branch on a cache that was wiped in step 2 and not yet repopulated.Optional.empty()is returned as a final answer for a topic that has policies. No exception, no log line.Steps 2 and 3 are the ordinary bundle ownership callbacks registered in
start()(lines 740-752); the defect is only their ordering against a queued continuation.A wipe without a new generation is harmless: with the map entry gone,
existingFuture == nullselects the retry branch,getTopicPoliciesAsynccallsprepareInitPoliciesCacheAsyncagain, installs and awaits a fresh generation and reads a populated cache. It is the newly installed, still-loading generation that turns a retry into a silent empty read. The injection seam used by the tests therefore replays exactly the two production calls of a bounce --cleanPoliciesCacheInitMap, thenprepareInitPoliciesCacheAsync-- in that window, and nothing else.The other done-checks in this file all ask for "present AND complete":
existingFuture == null || existingFuture.isDone()(line 250),chainedFuture != null && chainedFuture.isDone()(line 301),initFuture != null && !initFuture.isDone()(line 972). Line 587 is the outlier. Before PIP-376 (#23319) this read was guarded byinitialized == null || !initialized.isDone()(in the lines that commit removed); #23319 replaced it with the current predicate, evaluated synchronously in athenAccept, and #23786 moved the continuation behind the executor-lessthenComposeAsynchop ("switch thread to avoid potential metadata thread cost and recursive deadlock"), which widened the gap between the generation's completion and thecomputefrom a few instructions into a thread-pool hand-off.Consequence: the retention is burned in and the repair is lost
BrokerService.getManagedLedgerConfig(line 2405; policy resolution at lines 2434-2475) reads the topic policies withGetType.LOCAL_ONLY(line 2418) andGetType.GLOBAL_ONLY(line 2422) throughgetTopicPoliciesBypassSystemTopic, which delegates togetTopicPoliciesAsync(line 1494). AnOptional.empty()there is indistinguishable from "this topic has no retention policy":retentionPoliciesstaysnull(lines 2446-2450), the namespaceretention_policies-- or the broker defaults -- are taken instead (lines 2465-2475) and written into theManagedLedgerConfigwithsetRetentionTime/setRetentionSizeInMB(lines 2557-2559) before the managed ledger is opened. A failed read, by contrast, fails the topic load loudly; the remedies discussed in #26137 key off that failure signal and cannot see this silent empty.AbstractTopic.initTopicPolicy()(lines 629-670) reads the same policies again (GLOBAL_ONLYat lines 648-649,LOCAL_ONLYat lines 652-653) and hands them toTopicPolicyListenerWrapper.completeInitialization(global.orElse(null), local.orElse(null))(lines 657-658).TopicPolicyListenerWrapper.emitInitialPolicies(lines 145-152) emits nothing when the loaded value isnulland nothing was received during initialization, so an empty read there repairs nothing. One topic load performs seven policy reads in total.BrokerService.getTopic(lines 1379-1382) issues the first one, aLOCAL_ONLYread whose result is discarded, purely to await the namespace's policy-cache initialization before the load proceeds.getManagedLedgerConfigthen runs twice, readingLOCAL_ONLYandGLOBAL_ONLYconcurrently each time: once throughgetManagedLedgerFactoryForTopic(line 1619) whilefetchTopicPropertiesAsync(line 2133) picks the managed-ledger factory, and once fromcreatePersistentTopic0(line 2208) for the config the managed ledger is actually opened with.initTopicPolicy()reads both once more. The retention is wrong if thecreatePersistentTopic0pair is stale, and it stays wrong if theinitTopicPolicy()pair is stale too.With
topicPolicyListenerReplayEnabled=false(the default since #26134,ServiceConfigurationline 2057) no later broadcast re-notifies the topic once generation N+1 has finished loading, so the wrong retention persists for the life of the topic instance whileadmin.topicPolicies().getRetention(topic, true)keeps returning the topic value. Data is trimmed to the namespace retention.Related changes
computeblocks; [fix][broker] fix prepareInitPoliciesCacheAsync in SystemTopicBasedTopicPoliciesService #24980 madeprepareInitPoliciesCacheAsyncgenuinely await the initialization; [fix][broker] Don't let a closing topic-policies reader abort a concurrent cache-init reload #26132 identity-guards the cleanup side of a concurrent reload. None of them checks whether the future found after the hop is complete.PersistentTopic.initialize(), which runs aftergetManagedLedgerConfighas already decided the retention, so it does not cover this path.branch-4.0,branch-4.2andbranch-5.0-M1.Modifications
SystemTopicBasedTopicPoliciesService.getTopicPoliciesAsync: read the caches only when the init future currently owning the namespace has completed successfully.(Excerpt of the hunk: the six comment lines the patch adds above the predicate are omitted here.)
A future still loading, or one already dropped from the map, falls through to the existing retry branch, which re-awaits through
prepareInitPoliciesCacheAsyncbefore reading again. A comment above the predicate explains whyisDone()is required (the bounce-inside-the-hop interleaving).prepareInitPoliciesCacheAsync, the!insertedshort-circuit (service closed, or namespace being deleted -- unchanged behaviour: it still answers from the caches as they are) and the retry branch itself are untouched.One narrow behaviour change is worth naming: a generation that completed exceptionally is now routed to the retry as well, and that retry re-enters
prepareInitPoliciesCacheAsync, whoseexistingFuture.thenApply(__ -> true)(line 699) propagates the failure to the caller. The window is short -- betweeninitNamespacePolicyFuture.completeExceptionally(ex)(lines 680/683, and the timeout path at line 778) and the identity-guardedpolicyCacheInitMap.remove(namespace, initFuture)insidecleanupFailedPolicyCacheInit(line 826) -- and a read landing in it previously returned the still-populated cache, since that cleanup deliberately leaves cached policies in place (lines 812-814). Failing the topic load is the safe direction here: it is the signal [Bug] Topic policy loading errors are ignored which can result in data loss #26137 is about, and it is loud, where the previous behaviour was to serve a cache nobody was awaiting.The retry log line.
"The future of has been removed from cache, retry getTopicPolicies again"lost its subject when the{}placeholder was dropped in the slog migration ([improve] PIP-467: Convert pulsar-broker module logging from SLF4J to slog #25535), and "removed" describes only one of the three cases the branch now covers. It now reads"Policy cache init future is missing, failed or not yet complete, retry getTopicPolicies again", same level (INFO, an existing line) and same structurednamespaceattribute. The level is kept deliberately: a bundle bounce is a low-frequency ownership event, and the line fires at most once per policy read in flight at that instant, since reads that start while the replacement generation is still loading await it instead of retrying.PulsarService: a typed@VisibleForTesting public void setTopicPoliciesService(TopicPoliciesService), added next tosetTransactionBufferProviderand following its shape. The class carries@Setter(AccessLevel.PROTECTED), so the tests below could not install a replacement service through the generated setter and would otherwise have had to write the private field reflectively; an explicit method makes Lombok skip its protected one and turns a rename or a type change into a compile error. Nothing in production calls it.Tests, all in
pulsar-broker/src/test/java/org/apache/pulsar/broker/service/:StaleCacheGenerationInjectingTopicPoliciesService(new, package-private helper): aSystemTopicBasedTopicPoliciesServicethat opens a controlled stale-read window and holds the read inside it. When a read resumes from a generation it awaited successfully -- the production precondition -- and the namespace is armed, the helper replays a bundle bounce (cleanPoliciesCacheInitMap(namespace), thensuper.prepareInitPoliciesCacheAsync(namespace)fired and forgotten, asaddOwnedNamespaceBundleAsyncdoes) and overridescreateSystemTopicClientso the replacement generation's reader is wrapped in one whosehasMoreEventsAsync()andreadNextAsync()return nothing until a gate is released.initPolicesCachecallshasMoreEventsAsync()first, so that reader has not read a single event and the caches hold nothing for the namespace: the wipe closed the previous reader and the replacement is gated. The read is released into the continuation of line 582 only once that replacement generation is inpolicyCacheInitMapand its reader creation has been gated (the reader is wrapped on arrival, whether or not it has connected yet), and the gate is released only once every read the window holds has been observed to reach its decision -- either the retry, which is the re-entrantgetTopicPoliciesAsynccall with the same topic andGetTypethat line 604 makes, or the completion of the read's own future. A read that resumes from the same generation completion, such as the sibling of the concurrentLOCAL_ONLY/GLOBAL_ONLYpairs a topic load issues, joins that window if it gets there before the window closes, and otherwise re-awaits the replacement generation and opens or joins the next one, so no read of an armed namespace slips past a window while budget remains; once the budget is spent the helper passes every read through untouched. Nothing therefore depends on how fast the replacement reader is or on how the common pool schedules the continuation, andstaleWindowCount(),retriesInsideStaleWindow()andreadsServedInsideStaleWindow()count decisions rather than sampled state. Windows open one at a time per namespace and at mostbudgettimes (arm(namespace, budget)/disarm()), so a broker that retries its way out of the stale read converges; a gate is held only until the reads its window holds have decided, far insidetopicPoliciesCacheInitTimeoutSeconds, anddisarm()andclose()release any window still open. No cache is edited by hand, nothing is stubbed and no failure is injected: the helper only chooses when two production methods run and when the replacement reader is allowed to read.SystemTopicBasedTopicPoliciesServiceTest#testGetTopicPoliciesWhenBundleBounceReplacesPolicyCacheGeneration(new method in the existing class): the direct regression test ongetTopicPoliciesAsync.TopicPoliciesStaleCacheGenerationRetentionTest(new class): the end-to-end effect on the liveManagedLedgerConfig. It extendsMockedPulsarServiceBaseTestrather thanSharedPulsarBaseTestbecause it closes and replaces the broker'sTopicPoliciesService, which would leak into every other class sharing a runtime.Verifying this change
This change added tests and can be verified as follows:
SystemTopicBasedTopicPoliciesServiceTest#testGetTopicPoliciesWhenBundleBounceReplacesPolicyCacheGenerationcloses the current service, installs the injecting one withpulsar.setTopicPoliciesService(...)and callsstart(pulsar)so the replacement receives bundle-ownership callbacks like the production service. It creates a topic with a local policy, waits until generation 1 of the namespace is complete and the policy is visible in the cache, arms a budget of four windows, issues a singlegetTopicPoliciesAsync(topic, LOCAL_ONLY), and asserts that a window was opened (so the next assertions cannot hold vacuously), that the result is present and carries the policy that was set, that no read was answered inside a window (readsServedInsideStaleWindow()is zero) and that at least one read retried inside one.On unpatched master it fails with:
The helper's own log lines say why: one window opened, one read answered from the wiped cache of the still-loading replacement, no retry.
With the fix it passes: four windows opened, four retries, zero reads answered inside a window, and the retry line
Policy cache init future is missing, failed or not yet complete, retry getTopicPolicies againappears four times -- once per window, i.e. every window became a retry instead of an empty read. Each retry re-enters the read and can open the next window, which is what spends the budget of four. The whole class (20 tests) passes with the fix.TopicPoliciesStaleCacheGenerationRetentionTest#testTopicRetentionSurvivesPolicyCacheGenerationReplacementDuringLoadruns a broker withtopicLevelPoliciesEnabled=true,systemTopicEnabled=true,defaultNumberOfNamespaceBundles=1,defaultRetentionTimeInMinutes=0,defaultRetentionSizeInMB=0andtopicPolicyListenerReplayEnabled=false(all defaults, pinned because the test depends on them: a listener replay would push the retention onto the live managed ledger out of band), a namespace retention of 30 min and a topic-level retention of 300 min (synthetic values; only which of the two ends up live matters). A control load asserts that the livetopic.getManagedLedger().getConfig().getRetentionTimeMillis()equals 300 min. The topic is then unloaded (waiting untilgetTopicReferenceis empty), the helper is armed with a budget of 8 windows, the topic is loaded again and the helper is disarmed. The test asserts that a window was opened, that both thePersistentTopicand theManagedLedgerinstance differ from the control leg (ManagedLedgerFactoryImplcaches managed ledgers by name and ignores the config passed on a hit, so a reused instance would be asserting on the control leg's config), that the policy store still reports 300 min, and -- inside a bounded 10 s Awaitility, so that an out-of-band repair would still turn it green -- that the live retention equals 300 min.On unpatched master it fails with:
(1 800 000 ms is the 30 min namespace value; 18 000 000 ms is the 300 min topic value.)
On unpatched master the helper's log lines account for every read of that load: four windows covering all seven reads -- the preliminary read alone in the first, then each of the three concurrent
LOCAL_ONLY/GLOBAL_ONLYpairs joining one window -- all seven answered from the wiped cache, zero retries. That is why nothing is left to repair the retention.With the fix it passes: the same load opens eight windows and takes eight retries, with zero reads answered inside a window, and the retry line appears eight times -- the window budget. All eight land on the preliminary
LOCAL_ONLYread inBrokerService.getTopic, which is where the load spends the whole budget: each retry re-enters the helper, so the first read consumes it all and the six retention-deciding reads run against a settled generation. The topic keeps its 300 min retention. The predicate itself is exercised directly by the unit test above.Both tests fail on the unpatched code for the real reason (the empty read, the namespace retention), not because the helper forces internal state, and neither uses a sleep or a tuned delay: the window stays open until the read under test has decided, so the outcome cannot depend on reader speed or executor scheduling. In both, the vacuity assertion on
staleWindowCount()precedes the failing assertion and passes, so a run that never opened a window would fail there instead of reporting a false negative. Each test was run twice with the fix and twice on unpatched code; the window, retry and answered-read counts were identical across repeats../gradlew :pulsar-broker:checkstyleMain :pulsar-broker:checkstyleTestand./gradlew quickCheckpass on the change.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes
None of these apply: no configuration default changes, no threading change, and no change to the client or admin API; the fix only routes one more case into a retry path that already existed. The one added method,
PulsarService#setTopicPoliciesService, is@VisibleForTestingand mirrors the existingsetTransactionBufferProvider.