Skip to content

[fix][broker] Fix ownership-generation races in OwnershipCache removeOwnership and lock-expiry cleanup - #26197

Merged
lhotari merged 5 commits into
apache:masterfrom
SongOf:fix/broker-ownership-cache-generation-race
Sep 10, 2026
Merged

lhotari merged 5 commits into
apache:masterfrom
SongOf:fix/broker-ownership-cache-generation-race

Conversation

@SongOf

@SongOf SongOf commented Jul 16, 2026 •

Copy link
Copy Markdown
Contributor

Motivation

OwnershipCache.removeOwnership(bundle) removes whichever ResourceLock happens to be in
locallyAcquiredLocks at that moment, without checking that it belongs to the ownership
"generation" the caller was working on. Three defects follow from this:

  1. A stale unload chain can destroy a re-acquired ownership. When a bundle's resource lock
    expires, the lock-expiry callback in the cache loader fires unloadNamespaceBundle(...)
    asynchronously (dropping its future) and invalidates the local owner cache immediately. A
    concurrent lookup can then re-acquire the bundle (new lock, new OwnedBundle) while the stale
    unload chain is still running. When that chain reaches its final removeOwnership(bundle), it
    removes and releases the newly acquired lock, deleting the new owner's znode while the local
    cache still claims active ownership. This opens a transient dual-serving window (fenced-ledger
    errors on persistent topics, silent message loss on non-persistent topics) and causes ownership
    flapping. The same stale chain also deactivates the newer OwnedBundle through the cache lookup
    in updateBundleState(bundle, false), and its cleanUnloadedTopicFromCache(bundle) step can
    evict the newer generation's topics from the broker topic cache.

  2. removeOwnership reports success while an acquisition is in flight, leaving zombie
    ownership.
    The lock is only put into locallyAcquiredLocks after acquireLock completes, so
    a removeOwnership call racing with an in-flight acquisition finds the map empty, returns a
    completed future ("released"), and the lock installed afterwards is never cleaned up: the broker
    silently keeps the lock and the cache entry for a bundle the caller believes released (e.g. a
    deleted bundle/namespace), polluting load reports until restart.

  3. The lock-expiry callback never removes the expired lock from locallyAcquiredLocks
    (asymmetric add/remove bookkeeping), so a dead lock handle lingers in the map and drifts from
    the owned-bundle cache. The callback also drops the future returned by
    unloadNamespaceBundle(...), so failures of the triggered unload are silent.

Which load manager each part applies to

ExtensibleLoadManagerImpl.tryAcquiringOwnership does not go through OwnershipCache (it routes to
ServiceUnitStateChannel.assign), and NamespaceService.unloadNamespaceBundle /
removeOwnedServiceUnitAsync branch away from OwnershipCache when the extension is enabled. So:

  • Legacy load balancer only (OwnershipCache / OwnedBundle): the ownership-generation
    binding, removeOwnership(OwnedBundle), the per-bundle acquire/release barrier, and all of the
    lock-expiry listener changes. Defects 1 (lock/cache half), 2 and 3 above are legacy-only.
  • Both load balancers (BrokerService): the shared topic-future snapshot and the identity-safe
    cleanUnloadedTopicFromCache(bundle, snapshot), used by both OwnedBundle.handleUnloadRequest
    and ServiceUnitStateChannelImpl.closeServiceUnit. This is the topic-cache half of defect 1.

Modifications

Ownership generation (OwnershipCache, OwnedBundle)

  • Bind each OwnedBundle to the ResourceLock acquired for it, so an OwnedBundle instance
    represents a single ownership generation. Add removeOwnership(OwnedBundle) which releases only
    that generation via ConcurrentHashMap.remove(key, value); if the bundle has since been
    re-acquired, the newer ownership is left untouched. OwnedBundle.handleUnloadRequest now uses it
    and no longer calls updateBundleState (its own isActive flag is already flipped via CAS
    earlier in the method; the cache lookup could deactivate a newer generation).
  • Add a per-bundle async barrier (serialize(...)) that orders tryAcquiringOwnership against both
    removeOwnership(NamespaceBundle) and removeOwnership(OwnedBundle) for the same bundle, so a
    release can neither run in the gap between an acquire's cache publication and its lock install,
    nor report ownership fully released while a fresh acquisition lands moments later. The barrier is
    handed to the next queued operation on PulsarService.getExecutor() rather than inline, so the
    drain depth does not depend on the queue length (a nested inline drain overflowed the stack at
    roughly 1,800 queued operations and left the bundle permanently unacquirable on that broker). If
    the executor rejects the hand-off (broker shutting down), the barrier is released inline instead so
    the queued releases still run.
  • The lock-expiry listener is now fully generation-aware: it removes the expired lock from
    locallyAcquiredLocks with a two-arg remove, triggers unloadNamespaceBundle only if the cache
    still holds the listener's own OwnedBundle (a late listener must not close a re-acquired
    generation's topics or release its lock), invalidates the owned-bundle cache entry only if it
    still holds this generation, delivers onNamespaceBundleUnload to the ownership listeners only
    when that invalidation actually happened (a generation that was never published produced no owned
    event, and a newer generation in the cache is still owned; listeners such as the topic-policies
    service keep per-namespace counts on these events, so an unmatched unload event is not harmless),
    and logs failures of the triggered unload instead of dropping the future.
  • If the lock's expiry future is already done by the time the loader is about to publish the
    OwnedBundle, the load fails (IllegalStateException) instead of publishing an already-dead
    ownership; Caffeine drops a failed load, so no zombie cache entry remains.
  • removeOwnership(OwnedBundle) releases nothing for an OwnedBundle that is not bound to a lock
    (not created by the cache) instead of falling back to the generation-blind release; the two public
    lock-less OwnedBundle constructors, which have no callers, are deprecated.
  • OwnedBundle.handleUnloadRequest guards the topic snapshot and the start of the unload against
    synchronous exceptions and feeds them into the existing non-fatal path, so a failure there (for
    example RejectedExecutionException during shutdown) still releases the ownership instead of
    leaving the bundle inactive with its lock registered forever. The cleanup step keeps using the
    snapshot even when unloadServiceUnit threw, since the topic closes may already have been started
    from it.

Topic cache (BrokerService, ServiceUnitStateChannelImpl)

  • getTopicFuturesInBundle(bundle) captures one snapshot of the bundle's topic futures at unload
    start; the new unloadServiceUnit(..., snapshot) overload closes exactly that set, and
    cleanUnloadedTopicFromCache(bundle, snapshot) removes a topic only if its cached future is
    identical to the captured one. Both unload call sites thread the same snapshot through both
    steps. A topic that lands in the cache after the snapshot is deliberately left alone (see the
    comments at both call sites). The one-arg cleanUnloadedTopicFromCache(NamespaceBundle) is kept
    as a deprecated overload that captures the snapshot itself.

Split path (NamespaceService)

  • splitAndOwnBundleOnceAndRetry composes on the future returned by removeOwnership(bundle)
    instead of dropping it (a release failure is logged, not turned into a split failure), and the
    adjacent exceptionally handler now reports its own exception instead of the enclosing, always
    null one. The wait is bounded by metadataStoreOperationTimeoutSeconds: the release is queued
    behind any in-flight acquire of the same bundle, which has no timeout of its own, so on timeout
    the split completes and the release keeps running in the background.

Known follow-ups not addressed here

Raised during review and left out to keep this change scoped; neither changes the invariant above.

  • Failed acquires are no longer coalesced. When another broker holds the lock, concurrent
    tryAcquiringOwnership calls for the same bundle used to share one Caffeine load; they now run
    sequentially and each issues its own metadata write. Bounded by the number of distinct lookup-key
    variants in flight (a handful).
  • Residual concurrent-expiry window. If the lock expires right after the loader's isDone()
    check, the OwnedBundle is published and then dropped as soon as the listener's deferred
    invalidation runs on publication. In that window a lookup can be answered with this broker for an
    ownership that is already gone (clients re-lookup and converge), and since the deferred
    invalidation and the acquire's onNamespaceBundleOwned are both dependents of the same cache
    future, onNamespaceBundleUnload can be delivered before onNamespaceBundleOwned for that one
    generation.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • OwnershipCacheTest.testStaleUnloadDoesNotReleaseReacquiredOwnership: simulates the
    expiry-callback interleaving (cache invalidated, bundle re-acquired, then the stale unload chain
    completes) and asserts the re-acquired lock, znode, and active state survive.
  • OwnershipCacheTest.testTryAcquiringOwnershipWaitsForInFlightOwnedBundleRelease: races an
    acquire against an in-flight removeOwnership(OwnedBundle) and asserts the acquire is queued
    behind it and returns a fresh generation.
  • OwnershipCacheTest.testRemoveOwnershipWithAcquisitionInFlight: gates acquireLock behind a
    controllable future to call removeOwnership deterministically inside the in-flight window, then
    asserts no lock/cache entry/znode remains once both settle.
  • OwnershipCacheTest.testExpiredLockIsRemovedFromLocallyAcquiredLocks: releases the lock without
    going through removeOwnership and asserts locallyAcquiredLocks is cleaned up together with the
    cache and exactly one unload event is delivered.
  • OwnershipCacheTest.testExpiryBeforePublicationDoesNotLeaveActiveZombieOwnership: an
    already-expired lock handed to an in-flight load makes the acquisition fail and leaves no trace,
    and neither an owned nor an unload event is delivered for it.
  • OwnershipCacheTest.testDelayedExpiryListenerOfOldGenerationLeavesReacquiredOwnershipUntouched:
    completes the old generation's expiry future only after a newer generation was acquired, with the
    listener's unload routed through the real by-name lookup; asserts the newer generation stays
    cached, active, its lock is never released, and no unload event is delivered for it. Pins both
    the listener's unload guard and the identity check in invalidateLocalOwnerCache.
  • OwnershipCacheTest.testSerializedOperationsDrainDeepBacklogWithoutWedgingBundle: queues 5,000
    acquires behind a gated first one and asserts all complete and the barrier map is empty.
  • OwnershipCacheTest.testUnloadRequestStillReleasesOwnershipWhenTopicUnloadThrowsSynchronously:
    unloadServiceUnit throws synchronously; handleUnloadRequest must return a future, still
    release the ownership, and still run the cleanup with the original snapshot.
  • OwnershipCacheTest.testSerializedOperationsStillDrainWhenExecutorRejectsTasks: with a shut-down
    executor, a release queued behind an acquire still completes and the barrier map ends up empty.
  • OwnershipCacheTest.testRemoveOwnershipOfDetachedOwnedBundleDoesNotReleaseCurrentOwnership:
    releasing through an OwnedBundle not created by the cache leaves the current lock and znode alone.
  • NamespaceServiceTest.testSplitCompletesWhenReleasingTheOldBundleStalls: with the old bundle's
    release never completing, the split still completes once the metadata operation timeout elapses.
  • BrokerServiceTest.testCleanUnloadedTopicFromCacheIsGenerationSafe,
    ...IsGenerationSafeForBookkeeping, ...RemovesMatchingSnapshot: against a real embedded
    broker, a stale cleanup with an old snapshot leaves a reloaded topic and its bookkeeping intact,
    while a matching snapshot is removed.
  • Each new test was watched to fail before its corresponding change. OwnershipCacheTest (18), TopicPoliciesTest, MetadataStoreTopicPoliciesTest,
    BrokerServiceTest, NamespaceUnloadingTest, NamespaceOwnershipListenerTest and the split /
    unload cases of NamespaceServiceTest pass locally.

Does this pull request potentially affect one of the following parts:

If the box was checked, please highlight the changes

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

Threading model: OwnershipCache operations for the same bundle are now serialized per bundle, and
each queued operation starts on PulsarService.getExecutor() once the previous one settles.

@SongOf
SongOf force-pushed the fix/broker-ownership-cache-generation-race branch from 1f7e054 to 8f4d0b6 Compare July 22, 2026 14:38
@lhotari

lhotari commented Jul 22, 2026

Copy link
Copy Markdown
Member

I ran an AI-assisted review of this PR (Claude as the local reviewer plus OpenAI Codex gpt-5.6-sol as a second independent pass; both sets of findings were cross-verified against the actual sources before posting). Overall assessment first: this is a solid, well-tested fix for real races — the generation-bound OwnedBundle, the per-bundle acquire/release barrier, and the symmetric expiry-callback cleanup each address a genuine defect, and the new tests are deterministic and fail before the fix. Two findings are worth addressing before merge; the rest are optional.

1. cleanUnloadedTopicFromCache is only partially generation-safe (BrokerService.removeTopicFromCache)

In removeTopicFromCache(String, NamespaceBundle, CompletableFuture), only the final topics.remove(topic, createTopicFuture) is identity-guarded. The TopicEvent.UNLOAD BEFORE/SUCCESS notifications, the multiLayerTopicsMap removal (including replication-metrics teardown when a namespace empties), the compactor-stats removal, and forgetSegmentLoad all run unconditionally. A stale cleanup whose captured future no longer matches therefore still strips the newer generation's stats/load-report bookkeeping and fires unload events for a topic that stays cached and serving. This guarded-remove-with-unconditional-bookkeeping shape is admittedly pre-existing (removeTopicFromCache(AbstractTopic) has the same issue), but the new javadoc claims a stale call "can never evict a newer generation's topic", which currently only holds for the topics map itself, and testCleanUnloadedTopicFromCacheIsGenerationSafe only asserts getTopicReference() presence. Suggestion: perform the conditional topics.remove first, and only run the bookkeeping/events when that remove actually won.

2. The "race-free" thenRun claim in the expiry-before-publication guard isn't a JDK guarantee (OwnershipCache cache loader)

The comment on expiredBeforePublication asserts thenRun "only runs its action inline, before returning, when the future it is attached to is already complete." When the expiry future completes concurrently with registration, the JDK allows the action to be claimed and run by the completing thread after thenRun returns — so the flag can still read false and the loader publishes an OwnedBundle whose lock is already dead; tryAcquiringOwnership reports success and onNamespaceBundleOwned fires. Tracing the fallback shows the damage is tightly bounded: the listener has already done the two-arg locallyAcquiredLocks.remove, and the deferred invalidateLocalOwnerCache(bundle, ownedBundle) fires as part of the load future's completion, so the cache entry is dropped essentially at publication time. The residual effect is indistinguishable from a lock expiring immediately after a legitimately successful acquire, and state converges with no zombie — so low severity in practice, but the comment (and the matching claim in the invalidateLocalOwnerCache javadoc) should be corrected, and could optionally be hardened by re-checking rl.getLockExpiredFuture().isDone() after registering the listener.

3. (optional) serialize() invokes the operation inside ConcurrentHashMap.compute

In the common uncontended case (previous == null), thenCompose(ignore -> operation.get()) runs the supplier synchronously inside the compute lambda, under the map's bin lock. The suppliers only initiate async work today, but if any synchronous completion chain ever reaches serialize() for the same bundle on the same thread (plausible with already-completed metadata futures, e.g. a fast-failing acquire, or mocks in tests), that is a reentrant compute — behavior ConcurrentHashMap leaves undefined, and it could silently break the very mutual exclusion the barrier provides. Safer shape: inside compute, only capture the previous barrier and install the new one; chain precedingOp → operation.get() → completion after compute returns.

4. (optional) Straggler-topic semantics change

A topic installed in topics for the bundle after the snapshot is captured is now neither closed by the unload nor evicted by the cleanup. The old behavior was arguably worse (force-evicting live topics without closing them, enabling a second instance on reload), so this looks like a deliberate trade — but convergence for stragglers now relies entirely on per-connection ownership checks and lookup redirects. Worth a sentence in the PR description confirming it's intended.

5. (optional) Barrier latency notes, for awareness only

removeOwnership(NamespaceBundle) callers (bundle deletion, split, shutdown) now wait behind an in-flight acquire (bounded by metadata operation timeouts); tryAcquiringOwnership queues per bundle even on a cache hit while another operation is pending; and concurrent acquires that previously coalesced onto one shared failed load now retry sequentially. All per-bundle and bounded.

Also: nice catch on the ex → e fix in the split path — the old code could NPE by dereferencing the enclosing handler's null throwable.

@SongOf
SongOf force-pushed the fix/broker-ownership-cache-generation-race branch from 8f4d0b6 to 724a40a Compare July 23, 2026 19:05
@SongOf

SongOf commented Jul 23, 2026

Copy link
Copy Markdown
Contributor Author

@lhotari
Thanks for the thorough, cross-verified review — all addressed:

  1. Fixed. removeTopicFromCache now runs the identity-guarded topics.remove first and only proceeds to the multiLayerTopicsMap/replication-metrics/compactor-stats/segment-load bookkeeping and UNLOAD events if that remove actually wins. Added testCleanUnloadedTopicFromCacheIsGenerationSafeForBookkeeping to cover the case the existing test missed.
  2. Fixed. Corrected both the thenRun comment and the matching invalidateLocalOwnerCache javadoc — they no longer claim synchronous execution on a concurrent completion. Also added the suggested rl.getLockExpiredFuture().isDone() recheck alongside the flag to narrow the window further.
  3. Fixed. serialize() now only captures/installs the barrier inside compute(); the precedingOp → operation.get() → completion chain is built after compute() returns, so a fully synchronous, self-reentrant operation only ever produces ordinary, non-nested compute() calls.
  4. Documented. Added a comment at both snapshot sites (OwnedBundle and ServiceUnitStateChannelImpl, which share the same pattern) confirming this is intentional, why it's safer than the old force-evict behavior, and that BookKeeper ledger fencing — not this unload path — is the actual backstop against a stale writer.
  5. Noted, no action needed — though while addressing the above we found and fixed a related gap in the same spirit: removeOwnership(OwnedBundle), the path every normal unload actually goes through, wasn't wrapped in the barrier at all, so a concurrent acquire could still observe and report success for a generation mid-release. It's now serialized too, with a regression test (testTryAcquiringOwnershipWaitsForInFlightOwnedBundleRelease) — so this latency note extends there as well.

@lhotari

lhotari commented Jul 26, 2026

Copy link
Copy Markdown
Member

please merge origin/master (don't rebase) to the PR branch and resolve conflicts

# Conflicts:
#	pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java
#	pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java
@SongOf

SongOf commented Jul 27, 2026

Copy link
Copy Markdown
Contributor Author

please merge origin/master (don't rebase) to the PR branch and resolve conflicts

merged @lhotari

@alexandrebrg

Copy link
Copy Markdown
Contributor

Hello, we recently went through a production issue where two brokers seemed to hold ownership of the same topic.
This PR seems to be a good candidate to determine what happened, we are wondering if @SongOf would be willing to merge master against the branch, so maybe we could merge this ?

@lhotari lhotari added the triage/lhotari/important lhotari's triaging label for important issues or PRs label Sep 8, 2026
@lhotari

lhotari commented Sep 8, 2026

Copy link
Copy Markdown
Member

Hello, we recently went through a production issue where two brokers seemed to hold ownership of the same topic. This PR seems to be a good candidate to determine what happened, we are wondering if @SongOf would be willing to merge master against the branch, so maybe we could merge this ?

@alexandrebrg Which broker version did it happen with?
btw. The merge conflicts are already resolved. I'll perform a review.

@alexandrebrg

Copy link
Copy Markdown
Contributor

The broker version was pulsar-4.1.0 (a bit old I confess). We suspect that this occurred during a network outage we had, where connections with Zookeepers was heavily disturbed

Thank you

@lhotari

lhotari commented Sep 8, 2026

Copy link
Copy Markdown
Member

The broker version was pulsar-4.1.0 (a bit old I confess). We suspect that this occurred during a network outage we had, where connections with Zookeepers was heavily disturbed

@alexandrebrg In that case, you might also be affected by the issues fixed in #25910 and #25913. There's also #26000, which fixes an issue related to leader broker loss that could make other issues more likely, such as the bug that the current PR under discussion will fix.

@alexandrebrg

Copy link
Copy Markdown
Contributor

Good to know, we'll try to proceed to an upgrade then

@lhotari lhotari 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.

I read this at b08c6f94 in a worktree and exercised it rather than only reading it: OwnershipCacheTest passes 13/13 (7.7s) at the head commit, and I mutation-tested the production change to find out which parts the new tests actually pin.

The invariant is clear and, as far as I can tell, right. An OwnedBundle instance is now one ownership generation, because it carries the ResourceLock that created it; removeOwnership(OwnedBundle) only takes effect if locallyAcquiredLocks still maps the bundle to that lock; serialize() orders acquire against both release overloads per bundle, so an acquire can no longer be reported successful for a generation that is mid-release; and invalidateLocalOwnerCache(bundle, expected) and cleanUnloadedTopicFromCache(bundle, snapshot) extend the same identity rule to the owned-bundle cache and the broker topic cache. The generation-blind holes in the motivation are real, and this addresses them at the right level rather than papering over the symptom.

Four of the five parts of the change I mutated are each killed by a specific test, which is a good result for a change of this kind:

mutation of the production change killed by
removeOwnership(OwnedBundle) → generation-blind locallyAcquiredLocks.remove(bundle) testStaleUnloadDoesNotReleaseReacquiredOwnership
serialize(...) → return operation.get() (no barrier at all) testRemoveOwnershipWithAcquisitionInFlight, testTryAcquiringOwnershipWaitsForInFlightOwnedBundleRelease
lock-expiry listener → drop locallyAcquiredLocks.remove(namespaceBundle, rl) testExpiredLockIsRemovedFromLocallyAcquiredLocks
cache loader → never throw when expiry won the publication race testExpiryBeforePublicationDoesNotLeaveActiveZombieOwnership
ownedBundle == expectedOwnedBundle in invalidateLocalOwnerCache → if (ex == null) nothing — still 13/13 green

Three comments below. The first is the one place the generation rule was not applied, and it happens to be the most destructive of the three paths; the second is a defect in the new barrier itself; the third is that last table row.

Two of them turn on the same interleaving, which is worth stating once: the lock-expiry listener is registered with a plain thenRun, so when expiredFuture completes concurrently with registration the listener body can run on the completing thread, racing the loader rather than strictly before or after it. Under that ordering the loader fails the publication check, Caffeine drops the failed entry, a lookup installs a newer generation, and the still-running listener then operates on it. That is the same window testExpiryBeforePublicationDoesNotLeaveActiveZombieOwnership already builds a controllable ResourceLock for, so it should be cheap to drive deterministically.

Two things that are not defects but change how a reader should scope the fix:

  • The ownership-generation half only affects the legacy load balancer. ExtensibleLoadManagerImpl.tryAcquiringOwnership does not go through OwnershipCache at all — it routes to assign(...) on ServiceUnitStateChannel (ExtensibleLoadManagerImpl.java:578-592) — and NamespaceService.unloadNamespaceBundle / removeOwnedServiceUnitAsync branch away from OwnershipCache when the extension is enabled. What the ServiceUnitStateChannelImpl change contributes here is the topic-cache snapshot, which does apply to both paths. Saying which of the three motivating defects is fixed on which path would save the next reviewer the trip.
  • cleanUnloadedTopicFromCache(NamespaceBundle) was removed rather than deprecated. Every in-repo caller is updated, so this is only about out-of-tree callers of a broker-internal class — keeping a deprecated one-arg overload that captures the snapshot itself would cost nothing.

A few things I checked that turned out fine, in case it saves re-arguing them: closing a pre-taken snapshot instead of iterating the live topics map does not meaningfully widen the straggler window, because the old topics.forEach ran after updateBundleState, which completes synchronously (ownedBundlesCache uses a direct executor), so the two observation points are effectively the same instant; dropping updateBundleState(bundle, false) from handleUnloadRequest loses nothing, since the CAS above it already flipped isActive and the removed call resolved through the cache and could deactivate a newer generation; locallyAcquiredLocks.remove(bundle, lock) really is identity-based, because ResourceLockImpl overrides hashCode() but not equals(); and the release path does not deadlock against the new barrier, because ResourceLockImpl.release() completes expiredFuture before its own result future, so a re-entrant serialize() call only queues behind the in-flight barrier. The ex.getMessage() → e.getMessage() correction at NamespaceService.java:1102 is a real fix too — the old message reported the wrong exception.

One meta point, since it affects how much the green run is worth: a clean local run is weak evidence for a change like this, so I have based the assessment on the invariant and on what survives mutation instead.

@SongOf
SongOf force-pushed the fix/broker-ownership-cache-generation-race branch from b9ae6cc to 0adab8c Compare September 9, 2026 16:50
@SongOf
SongOf force-pushed the fix/broker-ownership-cache-generation-race branch from 0adab8c to 2166893 Compare September 9, 2026 16:51

@lhotari lhotari 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.

The follow-up in 2166893 addresses my three earlier comments: the delayed expiry listener checks its ownership generation, the normal queue drain hands off through the executor, and the new delayed-listener test exercises both generation guards. I have one nonblocking test-maintenance comment below.

@lhotari lhotari added this to the 5.0.0-M2 milestone Sep 10, 2026
@lhotari
lhotari merged commit 37d64fe into apache:master Sep 10, 2026
43 checks passed
lhotari added a commit that referenced this pull request Sep 10, 2026
…Ownership and lock-expiry cleanup (#26197)

Co-authored-by: maxlisongsong <[email protected]>
Co-authored-by: Lari Hotari <[email protected]>
(cherry picked from commit 37d64fe)
lhotari added a commit that referenced this pull request Sep 11, 2026
…Ownership and lock-expiry cleanup (#26197)

Co-authored-by: maxlisongsong <[email protected]>
Co-authored-by: Lari Hotari <[email protected]>
(cherry picked from commit 37d64fe)
Radiancebobo pushed a commit to Radiancebobo/pulsar that referenced this pull request Oct 8, 2026
…Ownership and lock-expiry cleanup (apache#26197)

Co-authored-by: maxlisongsong <[email protected]>
Co-authored-by: Lari Hotari <[email protected]>
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.

4 participants