Conversation
PulsarRegistrationManager never delivered the bookkeeper RegistrationManager contract - addRegistrationListener was a no-op stub - and never rebuilt the /ledgers/available ephemeral registration after a metadata store session loss: the bookie disappeared from the placement view until process restart, on every backend used through the pulsar metadata driver. On SessionLost, re-validate every tracked registration with a read / disambiguate / version-compared write / reacquire loop on a single-thread mutation executor, at most one target per registration path, so that a session lost while a previous loop is still draining is never swallowed. The write is always version-compared, a non-ephemeral occupant is never written, and the loop only stops after the lock-handle reacquire, which is the leg that moves the record onto the live session on stores where a put does not transfer the ephemeral ownership. The registration-expired listener fires only after repeated confirmation of a genuine foreign owner; public registration calls keep the synchronous exception contract through a submit-and-block adapter. Signed-off-by: dao-jun <daojun@apache.org>
Run the real lock stack on an in-memory store that models the identity, session and version semantics of the metadata store backends, covering: purge-driven recreation, retry through store unavailability, grace-window rebinds with a real write, invisible version substitutions, adopt-then- tick interleavings, three-strike escalation and counter resets, close races, readonly conversions around in-flight revalidations, the blocking adapter contract, new-session wake-up, dead-lock-handle recovery, non-ephemeral occupants, and the zk ownership semantics carried by the reacquire leg. Signed-off-by: dao-jun <daojun@apache.org>
lhotari
left a comment
There was a problem hiding this comment.
Thanks for digging into bookie registration after a session loss. Wiring addRegistrationListener to actually do something, and removing the handle in doUnregisterBookie instead of releasing it again at close(), are both real fixes. Not notifying listeners on every SessionLost is also the right call, since BookKeeper's listener shuts the bookie down if re-registration fails.
I have concerns about the approach, though. On ZooKeeper the lock layer already re-creates these registrations: LockManagerImpl revalidates every lock on SessionReestablished, and ZKSessionTest.testReacquireLocksAfterSessionLost covers it. This loop becomes a second writer on the same paths. In the current form it can end with a registration handle that never revalidates, or with a torn-down registration coming back. Details are inline. The description also needs a few corrections: there is no etcd store on master since PIP-462, and the new test class has 19 tests.
| * the mutation executor. | ||
| */ | ||
| private void reacquireRegistrationLock(RegistrationRevalidation target) { | ||
| lockManager.acquireLock(registrationPath(target.bookieId, target.readOnly), target.cachedValue) |
There was a problem hiding this comment.
[BUG] The reacquired handle normally stays in Init and never revalidates
The loop writes the registration before reacquiring it, so acquireLock here always finds the node already present. ResourceLockImpl.acquire then goes through revalidate, whose same-session adoption branch updates version/value and returns (ResourceLockImpl.java:297-302) without ever setting state = Valid. The new handle stays in Init, and silentRevalidateOnce returns immediately for anything that isn't Valid (ResourceLockImpl.java:225-227).
So after the loop converges, if the node is deleted again without a session loss, nothing re-creates it and no expiry is signalled. The adoption gap is pre-existing in ResourceLockImpl, but this PR makes it the normal recovery path. Could the reacquire avoid the pre-write, or could adoption mark the lock Valid (with a test in LockManagerTest)? A regression test that deletes the node after convergence and expects it back would pin it.
| log.info().attr("bookieId", target.bookieId) | ||
| .log("Registration lock is no longer valid, re-acquiring it"); | ||
| } | ||
| registrationMap(target.readOnly).put(target.bookieId, newLock); |
There was a problem hiding this comment.
[BUG] The superseded lock handle stays live and can resurrect a torn-down registration
onReacquireCompleted swaps in the new handle but leaves the old ResourceLockImpl in state Valid. LockManagerImpl.locks.put drops it from the map (LockManagerImpl.java:92), but not its scheduled retry (ResourceLockImpl.java:252-253, backoff up to a minute). When that retry fires, revalidate() re-creates a missing path (ResourceLockImpl.java:269-276).
Sequence: a long outage leaves the old handle with a pending retry; the loop converges; the bookie goes readonly or unregisters; the old retry then writes back a writable registration nothing tracks, until the session ends. discardUnreleasableLock doesn't help since it only completes the expired future. Could the old handle be retired (state Released, retry cancelled) through a lock-layer hook, or could the loop keep and revalidate the existing handle instead of replacing it?
| private void addRevalidationTarget(List<RegistrationRevalidation> newTargets, | ||
| BookieId bookieId, boolean readOnly, ResourceLock<BookieServiceInfo> lock) { | ||
| RegistrationRevalidation existing = activeRevalidations.get(registrationPath(bookieId, readOnly)); | ||
| if (existing == null || existing.lock != lock) { |
There was a problem hiding this comment.
[BUG] A SessionLost handled while the write or reacquire is in flight is dropped
addRevalidationTarget skips a path whose active target holds the same handle, and target.lock only changes at :880, after the reacquire. If the write landed on a session that is lost while the reacquire is in flight, the reacquire adopts the record without writing (createdBySelf is identity-scoped on Oxia, OxiaMetadataStore.java:150) and :885 stops the loop with the record bound to the dead session — the case the description says is never swallowed. Could a loss generation or a rerun flag be checked before stopRevalidation?
| log.warn().attr("bookieId", bookieId) | ||
| .log("Cannot release the registration lock through its handle," | ||
| + " removing the registration record directly"); | ||
| removeOwnRegistrationRecord(registrationPath(bookieId, readOnly)); |
There was a problem hiding this comment.
[BUG] Registration cleanup isn't sequenced with later registrations, and a failed unregister can't be retried
Two related lifecycle gaps:
- After a
BadVersion,releaseRegistrationLockreturns onceremoveOwnRegistrationRecordhas been issued; its get→delete runs later on store callback threads. AnunregisterBookiefollowed byregisterBookie(or a readonly/writable flap) can let the new handle adopt the record — adoption doesn't bump the version — and the stale cleanup then deletes it with the version it read. The cleanup needs to complete inside the mutation, not just move threads. doUnregisterBookienow removes the handle before releasing it (:287-294). If the release fails transiently, the exception propagates, and a laterunregisterBookiefinds no handle and returns success without deleting anything. Keeping the handle (or an explicit cleanup obligation) until the delete succeeds would fix that.
| return; | ||
| } | ||
|
|
||
| // Transient failure, for example the store being unavailable: keep retrying. |
There was a problem hiding this comment.
[MINOR] During a ZK→Oxia migration the loop retries writes the store deliberately refuses
Every ZK store is wrapped in DualMetadataStore, which emits a synthetic SessionLost in PREPARATION precisely so components stop writing (DualMetadataStore.java:203-205) and rejects puts until the migration completes (DualMetadataStore.java:372-373). This loop reacts by writing; the IllegalStateException is treated as transient, so it logs an INFO and a WARN with stack trace every few seconds per registration for the whole phase. Could non-retriable rejections pause the loop (resuming after the migration) rather than retry, with repeated failures logged once per streak?
| return registrationManagerLogs.countOf(ESCALATE); | ||
| } | ||
|
|
||
| /** Waits longer than the maximum backoff to observe that no further ring ticks happen. */ |
There was a problem hiding this comment.
[MINOR] The tests rely on ~63 s of fixed sleeps and on another class's log text
A few requests on the tests:
assertRingIsQuiet()sleeps 6 s and is called nine times, plus about 9.5 s of other sleeps, because the 1–5 sBackoffis hard-coded. A package-private constructor taking theBackoffwould make these fast and deterministic.- Several checks count log lines, including
ResourceLockImpl's own messages, so rewording a log in another class breaks this test. Where the fake already exposes versions and sessions, asserting state is sturdier. - The logger level set at
:120isn't restored in teardown, the last assertion of T6 (:450-453) checks the fake's purge rather than the manager's cleanup, and references like "round3" or "Review item 6/6" shouldn't be in committed tests. AssertJ would also be preferred overassertTrue(x >= n).
The adoption branch of ResourceLockImpl.revalidate updated the version without marking the lock valid, so a handle that adopted an already present record stayed in Init and was never revalidated again: nothing re-created the record when it was later deleted, and no expiry was signalled. On stores where the creator identity outlives the session, the adoption also left the record bound to the expired session. The adoption is now a version-compared write that re-binds the record to the live session and leaves the lock valid, and its completion records the fresh version also when the lock was released while the write was in flight, because the release deletes with that expectation. A BadVersion from any version-compared write of a revalidation means a concurrent revalidation or acquire moved the record, not a conflict: the decision is made again from a fresh read, once, instead of failing or expiring the lock. A genuine foreign holder still surfaces as a LockBusy. A revalidation in flight while the lock is released could re-create the path behind the release: the release is now serialized with the revalidations, and a released lock never re-creates a missing path. A lock handle replaced in the lock manager is retired, cancelling its pending revalidation, so it cannot resurrect a torn-down resource. Signed-off-by: dao-jun <daojun@apache.org>
Replace the manager-side revalidation loop with the coordination layer's own recovery: after the session is re-established, the lock manager already revalidates every tracked lock, and with the adoption fix it re-creates a swept record and re-binds a record left by an expired session on every backend. The manager maps the expiry of a registration lock to the bookkeeper registration-expired contract, with the still-tracked handle owning the notification, so a voluntary unregister does not drive a re-registration of what was just torn down, and the escalation is the lock layer's single server-authoritative expiry observation instead of a three-strike counter. The unregister path keeps the handle tracked until its release succeeded, so a failed unregister can be retried, and the direct-record cleanup of an unreleasable lock completes inside the mutation, so it cannot delete a record a subsequent registration adopted. Signed-off-by: dao-jun <daojun@apache.org>
The package-private revalidation wrapper has no caller outside ResourceLockImpl: restore the original private visibility. Also collapse the blank-line debris the fake store pruning left behind. Signed-off-by: dao-jun <daojun@apache.org>
| // The previous handle for this path is superseded: retire it, so that | ||
| // a revalidation of it that is still scheduled or in flight cannot | ||
| // re-create the path and resurrect a torn-down resource. | ||
| replaced.retire(); |
There was a problem hiding this comment.
Potential deadlock on the ZooKeeper path. For example:
- Thread A enters synchronized
old.release(), sendsstore.delete(), and is descheduled before attaching its.thenRuncompletion. It still holds the old lock monitor. - ZooKeeper completes the delete. Another caller acquires the same path; its create callback holds the
LockManagerImplmonitor and blocks here inreplaced.retire()waiting for A. - A resumes. Because the delete future is already complete,
.thenRunruns inline on A. ItsexpiredFuture.complete()invokes the manager cleanup callback, which waits for the manager monitor held by the create callback.
Both threads are now blocked. ZooKeeper's single event executor does not prevent this because A's late completion runs on A's thread. Avoid taking the old lock monitor while holding the manager monitor, and preserve the map replacement ordering.
| Thread.currentThread().interrupt(); | ||
| // Stop accepting mutations. Store operations that were already in flight may still | ||
| // complete after this: their results are discarded silently by the completion handlers. | ||
| mutationExecutor.shutdownNow(); |
There was a problem hiding this comment.
shutdownNow() can strand a caller whose mutation was already accepted. runBlockingMutation() waits on a future completed only by the queued task. For example, one registerBookie() is blocked on a slow metadata operation while a concurrent unregisterBookie() has queued behind it; close() discards the queued task, so that caller waits until the 65-second timeout. Complete the futures for discarded tasks exceptionally so accepted calls fail promptly during shutdown.
Motivation
Two long-standing gaps in
PulsarRegistrationManager, which serves every bookie whose metadata service URI uses themetadata-store:scheme (zk and oxia alike):addRegistrationListeneris an empty stub — the comment reads "Not implemented. Does not seem to map into MetadataStoreExtended." The bookkeeperRegistrationManagercontract is never delivered:BookieStateManagerregisters a listener so it can re-register the bookie when the registration expires, and that listener can never fire./ledgers/available/<bookie>ephemeral registration is swept server-side and never rebuilt. The bookie silently disappears from the placement view until the process restarts.Modifications
The recovery is carried by the coordination layer, which already owns these paths as resource locks — not by a second writer inside the manager. When the session is re-established,
LockManagerImplalready revalidates every tracked lock, andResourceLockImpl.revalidatealready re-creates a swept record, replaces a record left by an expired session, or fails withLockBusyagainst a foreign holder (ZKSessionTest.testReacquireLocksAfterSessionLostcovers that machinery on zk). Two gaps made that path incomplete; this change fixes them in the lock layer and wires the manager on top of it:Initand was never revalidated again: nothing re-created the record when it was later deleted, and no expiry was signalled. On backends where the creator identity outlives the session (oxia scopescreatedBySelfto the client identity), the adoption also left the record bound to the expired session. The adoption is now a version-compared write that re-binds the record to the live session and leaves the lock valid; aBadVersionfrom that write, typically our own concurrent re-acquire, routes back to a fresh read instead of expiring the lock.On top of that, the manager:
onRegistrationExpired(), the point where the consumer re-attempts the registration. Only the still-tracked handle owns the notification: the completion of a voluntary unregister's own release cannot drive the bookie into a re-registration of what was just torn down.Where the earlier revision of this PR escalated after three consecutive confirmations of a foreign holder, the escalation is now the lock layer's single, server-authoritative expiry observation — same terminal state (the consumer's re-registration fails and the bookie exits), one code path, no second writer. A same-value record of this bookie, including a stale copy left by an expired session, is re-adopted or re-created rather than waited out.
On zk the session events already flow today; on oxia this additionally requires the store to honor
enableSessionWatcher(companion PR #26712).Verifying this change
./gradlew :pulsar-metadata:test— 1006 tests, all green../gradlew :pulsar-metadata:checkstyleMain :pulsar-metadata:checkstyleTest :pulsar-metadata:spotlessCheck— clean.Tests
LockManagerTest(new): a lock that adopts an already present record is valid and re-creates the path when the record is deleted afterwards (the regression that pins the adoption fix), against every metadata store backend; a handle superseded by a newer acquire does not resurrect the path after the replacement is released; a revalidation whose read is still in flight when the lock is released does not re-create the path behind the release.PulsarRegistrationManagerRevalidationTest(rewritten, state-based assertions on the fake store's records, sessions and listener invocations — no log-text coupling): a swept registration is re-created on the live session after the re-establishment without notifying the listeners; a zombie record in the server-side grace window is re-bound to the live session both under oxia identity semantics (adoption write) and under zk ownership semantics (delete-and-recreate); a foreign different-value holder expires the registration through the listener; a voluntary unregister and a readonly conversion never notify; a failed unregister is retryable and the record survives it; the direct-record cleanup of an unreleasable lock completes before the unregister returns; close releases all registrations.PulsarRegistrationManagerTest(70 tests) is unchanged and green.