Registrations are held as ephemeral resource locks, which the coordination layer
+ * revalidates and re-establishes on its own after a metadata store session loss: when the
+ * session is re-established, every tracked lock is revalidated, a lock whose record was swept
+ * is re-created on the live session, and a record still owned by an expired session is
+ * re-written and re-bound. This manager therefore carries no re-registration loop of its own:
+ * it maps the expiry of a registration lock — the point where the lock layer concludes the
+ * registration cannot be held — to the bookkeeper
+ * {@link RegistrationListener#onRegistrationExpired()} contract, so that the consumer can
+ * attempt a fresh registration.
+ *
+ *
The bookkeeper {@code ZKRegistrationManager} reports a registration expiry on every
+ * session loss and relies on the surrounding zk client stack to keep re-issuing its
+ * operations against the rebuilt session. The metadata store stack used here does not retry
+ * store operations across a session loss; the lock layer's revalidation is what converges the
+ * registration instead, and the expiry notification fires only when the registration was
+ * genuinely lost to another owner — 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. A
+ * genuine conflict still ends in the same terminal state: the consumer's re-registration
+ * fails and the bookie exits.
+ */
@CustomLog
public class PulsarRegistrationManager implements RegistrationManager {
+ private static final long MUTATION_EXECUTOR_SHUTDOWN_TIMEOUT_MS = 5000;
+
private final MetadataStoreExtended store;
private final CoordinationService coordinationService;
private final LockManager lockManager;
@@ -73,7 +105,27 @@ public class PulsarRegistrationManager implements RegistrationManager {
private final Map> bookieRegistration = new ConcurrentHashMap<>();
private final Map> bookieRegistrationReadOnly = new ConcurrentHashMap<>();
- private final List listeners = new ArrayList<>();
+ private final List listeners = new CopyOnWriteArrayList<>();
+
+ /**
+ * Single-threaded executor that serializes all the registration mutations and the
+ * session-loss revalidation loop. Public registration calls submit here and block until
+ * completion, because the bookkeeper state machine invokes them synchronously and expects
+ * failures to propagate to the calling thread.
+ */
+ private final ScheduledExecutorService mutationExecutor = Executors.newSingleThreadScheduledExecutor(
+ new DefaultThreadFactory("bookie-registration-mutation"));
+
+ /**
+ * The registration-expired listeners are notified from this dedicated executor instead of
+ * inline on the mutation executor or on a store event thread: the listeners are third-party
+ * code and must stay isolated, so that a slow or failing listener can neither hold up the
+ * mutation executor nor break the notification of the other listeners.
+ */
+ private final ExecutorService listenerExecutor = Executors.newSingleThreadExecutor(
+ new DefaultThreadFactory("bookie-registration-listener"));
+
+ private final AtomicBoolean closed = new AtomicBoolean(false);
PulsarRegistrationManager(MetadataStoreExtended store, String ledgersRootPath, AbstractConfiguration> conf) {
this.store = store;
@@ -88,25 +140,38 @@ public class PulsarRegistrationManager implements RegistrationManager {
@Override
public void close() {
- for (ResourceLock rwBookie : bookieRegistration.values()) {
- try {
- rwBookie.release().get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
- } catch (ExecutionException | TimeoutException ignore) {
- log.error().attr("lock", rwBookie).exception(ignore.getCause()).log("Cannot release correctly");
- } catch (InterruptedException ignore) {
- log.error().attr("lock", rwBookie).exception(ignore).log("Cannot release correctly");
- Thread.currentThread().interrupt();
- }
+ if (!closed.compareAndSet(false, true)) {
+ return;
}
- for (ResourceLock roBookie : bookieRegistrationReadOnly.values()) {
- try {
- roBookie.release().get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
- } catch (ExecutionException | TimeoutException ignore) {
- log.error().attr("lock", roBookie).exception(ignore.getCause()).log("Cannot release correctly");
- } catch (InterruptedException ignore) {
- log.error().attr("lock", roBookie).exception(ignore).log("Cannot release correctly");
- 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();
+ listenerExecutor.shutdown();
+ try {
+ mutationExecutor.awaitTermination(MUTATION_EXECUTOR_SHUTDOWN_TIMEOUT_MS, MILLISECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+
+ for (Map> registrations :
+ List.of(bookieRegistration, bookieRegistrationReadOnly)) {
+ for (ResourceLock lock : registrations.values()) {
+ try {
+ lock.release().get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
+ } catch (ExecutionException | TimeoutException e) {
+ log.error().attr("lock", lock).exception(e.getCause()).log("Cannot release correctly");
+ try {
+ removeOwnRegistrationRecord(lock.getPath());
+ } catch (BookieException cleanupFailure) {
+ log.warn().attr("lock", lock).exception(cleanupFailure)
+ .log("Cannot remove the registration record directly");
+ }
+ discardUnreleasableLock(lock);
+ } catch (InterruptedException ignore) {
+ log.error().attr("lock", lock).exception(ignore).log("Cannot release correctly");
+ Thread.currentThread().interrupt();
+ }
}
}
try {
@@ -132,34 +197,29 @@ public String getClusterInstanceId() throws BookieException {
@Override
public void registerBookie(BookieId bookieId, boolean readOnly, BookieServiceInfo bookieServiceInfo)
throws BookieException {
- String regPath = bookieRegistrationPath + "/" + bookieId;
- String regPathReadOnly = bookieReadonlyRegistrationPath + "/" + bookieId;
log.info().attr("bookieId", bookieId).attr("readOnly", readOnly).attr("info", bookieServiceInfo)
.log("RegisterBookie");
+ runBlockingMutation(() -> doRegisterBookie(bookieId, readOnly, bookieServiceInfo));
+ }
+ /**
+ * Internal registration path, to be run only on the mutation executor. It is only ever
+ * called through the public {@link #registerBookie(BookieId, boolean, BookieServiceInfo)}
+ * adapter, which blocks on the same single-thread executor.
+ */
+ private void doRegisterBookie(BookieId bookieId, boolean readOnly, BookieServiceInfo bookieServiceInfo)
+ throws BookieException {
try {
if (readOnly) {
- ResourceLock rwRegistration = bookieRegistration.remove(bookieId);
- if (rwRegistration != null) {
- log.info().attr("bookieId", bookieId)
- .log("Bookie was already registered as writable, unregistering");
- rwRegistration.release().get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
- }
+ unregisterTrackedRegistration(bookieId, false);
bookieRegistrationReadOnly.put(bookieId,
- lockManager.acquireLock(regPathReadOnly, bookieServiceInfo)
- .get(BLOCKING_CALL_TIMEOUT, MILLISECONDS));
+ acquireRegistrationLock(bookieId, true, bookieServiceInfo));
} else {
- ResourceLock roRegistration = bookieRegistrationReadOnly.remove(bookieId);
- if (roRegistration != null) {
- log.info().attr("bookieId", bookieId)
- .log("Bookie was already registered as read-only, unregistering");
- roRegistration.release().get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
- }
+ unregisterTrackedRegistration(bookieId, true);
bookieRegistration.put(bookieId,
- lockManager.acquireLock(regPath, bookieServiceInfo)
- .get(BLOCKING_CALL_TIMEOUT, MILLISECONDS));
+ acquireRegistrationLock(bookieId, false, bookieServiceInfo));
}
} catch (ExecutionException | TimeoutException ee) {
log.error().exception(ee).log("Exception registering ephemeral node for Bookie");
@@ -179,22 +239,85 @@ public void registerBookie(BookieId bookieId, boolean readOnly, BookieServiceInf
@Override
public void unregisterBookie(BookieId bookieId, boolean readOnly) throws BookieException {
+ runBlockingMutation(() -> doUnregisterBookie(bookieId, readOnly));
+ }
+
+ private void doUnregisterBookie(BookieId bookieId, boolean readOnly) throws BookieException {
+ unregisterTrackedRegistration(bookieId, readOnly);
+ }
+
+ /**
+ * Releases the tracked registration of the bookie. The handle stays tracked until its
+ * release succeeded, so that a failed unregister, for example on a transient store
+ * failure, can be retried: a retry that finds no handle would silently succeed without
+ * deleting anything. Must run on the mutation executor.
+ */
+ private void unregisterTrackedRegistration(BookieId bookieId, boolean readOnly) throws BookieException {
+ ResourceLock registration = registrationMap(readOnly).get(bookieId);
+ if (registration == null) {
+ return;
+ }
+ releaseRegistrationLock(bookieId, readOnly, registration);
+ registrationMap(readOnly).remove(bookieId, registration);
+ }
+
+ /**
+ * Acquires the registration lock and watches its expiry: the expiry of a tracked
+ * registration, the point where the lock layer concludes the registration cannot be held,
+ * is what drives the bookkeeper registration-expired contract.
+ */
+ private ResourceLock acquireRegistrationLock(
+ BookieId bookieId, boolean readOnly, BookieServiceInfo bookieServiceInfo)
+ throws InterruptedException, ExecutionException, TimeoutException {
+ ResourceLock lock = lockManager
+ .acquireLock(registrationPath(bookieId, readOnly), bookieServiceInfo)
+ .get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
+ lock.getLockExpiredFuture().thenRun(() -> scheduleOnMutationExecutor(
+ () -> handleRegistrationExpired(bookieId, readOnly, lock)));
+ return lock;
+ }
+
+ /**
+ * The lock of a tracked registration expired. Only the still-tracked handle owns the
+ * notification: a voluntary unregister removes the handle only after its release
+ * succeeded, and the completion of that release must not drive the bookie into a
+ * re-registration of what was just torn down.
+ */
+ private void handleRegistrationExpired(BookieId bookieId, boolean readOnly,
+ ResourceLock lock) {
+ if (closed.get()) {
+ return;
+ }
+ if (registrationMap(readOnly).remove(bookieId, lock)) {
+ log.warn().attr("bookieId", bookieId)
+ .log("The bookie registration expired, notifying the listeners");
+ notifyRegistrationExpired();
+ }
+ }
+
+ /**
+ * Releases a registration lock, tolerating a BadVersion caused by a revalidation of the
+ * same lock still in flight (its version expectation was reset): in that case the record
+ * is removed directly, when it belongs to this store identity, and the unreleasable
+ * handle is discarded, because the registration it represented is torn down anyway.
+ */
+ private void releaseRegistrationLock(BookieId bookieId, boolean readOnly,
+ ResourceLock registration) throws BookieException {
try {
- if (readOnly) {
- ResourceLock roRegistration = bookieRegistrationReadOnly.get(bookieId);
- if (roRegistration != null) {
- roRegistration.release().get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
- }
- } else {
- ResourceLock rwRegistration = bookieRegistration.get(bookieId);
- if (rwRegistration != null) {
- rwRegistration.release().get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
- }
- }
+ registration.release().get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new BookieException.MetadataStoreException(ie);
} catch (ExecutionException | TimeoutException e) {
+ if (e instanceof ExecutionException
+ && e.getCause() instanceof MetadataStoreException.BadVersionException) {
+ log.warn().attr("bookieId", bookieId)
+ .log("Cannot release the registration lock through its handle,"
+ + " removing the registration record directly");
+ removeOwnRegistrationRecord(registrationPath(bookieId, readOnly));
+ discardUnreleasableLock(registration);
+ return;
+ }
throw new BookieException.MetadataStoreException(e);
}
}
@@ -366,7 +489,7 @@ public boolean format() throws Exception {
instanceId.getBytes(StandardCharsets.UTF_8), Optional.of(-1L))
.get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
- log.info("Successfully formatted BookKeeper metadata");
+ log.info().log("Successfully formatted BookKeeper metadata");
return true;
}
@@ -406,6 +529,137 @@ public boolean nukeExistingCluster() throws Exception {
@Override
public void addRegistrationListener(RegistrationListener listener) {
- // Not implemented. Does not seem to map into MetadataStoreExtended.
+ listeners.add(listener);
+ }
+
+ /**
+ * Runs a registration mutation on the mutation executor, blocking the calling thread until
+ * the mutation completes and propagating its failure, to preserve the synchronous contract
+ * of the bookkeeper registration calls.
+ */
+ private void runBlockingMutation(RegistrationMutation mutation) throws BookieException {
+ CompletableFuture future = new CompletableFuture<>();
+ // The budget covers the worst-case internal blocking sequence and the time the mutation
+ // may spend queued behind other work on the single-thread executor.
+ long deadline = System.currentTimeMillis() + 2 * BLOCKING_CALL_TIMEOUT + 5_000;
+ try {
+ mutationExecutor.execute(() -> {
+ try {
+ mutation.run();
+ future.complete(null);
+ } catch (Throwable t) {
+ future.completeExceptionally(t);
+ }
+ });
+ } catch (RejectedExecutionException e) {
+ throw new BookieException.MetadataStoreException(e);
+ }
+
+ try {
+ future.get(deadline - System.currentTimeMillis(), MILLISECONDS);
+ } catch (ExecutionException e) {
+ if (e.getCause() instanceof BookieException) {
+ throw (BookieException) e.getCause();
+ }
+ throw new BookieException.MetadataStoreException(e.getCause());
+ } catch (InterruptedException ie) {
+ Thread.currentThread().interrupt();
+ throw new BookieException.MetadataStoreException(ie);
+ } catch (TimeoutException te) {
+ throw new BookieException.MetadataStoreException(te);
+ }
+ }
+
+ private interface RegistrationMutation {
+
+ void run() throws BookieException;
+ }
+
+ private void scheduleOnMutationExecutor(Runnable task) {
+ try {
+ mutationExecutor.execute(task);
+ } catch (RejectedExecutionException ignore) {
+ // The executor rejects only after close. Run the continuation inline: after close
+ // there is no mutation left to race with, and the continuations only do terminal
+ // bookkeeping, like the cleanup of a stale in-flight registration write. A task
+ // accepted between close()'s closed-CAS and the shutdownNow() that follows is
+ // dropped instead; what it could have cleaned up is a leftover bound to the dying
+ // session, which the store's own shutdown clears.
+ if (closed.get()) {
+ task.run();
+ }
+ }
+ }
+
+ private void notifyRegistrationExpired() {
+ try {
+ listenerExecutor.execute(() -> {
+ for (RegistrationListener listener : listeners) {
+ try {
+ listener.onRegistrationExpired();
+ } catch (Throwable t) {
+ log.error().exception(t).log("Failed to notify the registration listener");
+ }
+ }
+ });
+ } catch (RejectedExecutionException ignore) {
+ // The registration manager was closed
+ }
+ }
+ /**
+ * Removes the registration record directly, completing inside the current mutation: a
+ * cleanup that only issues the get/delete and returns could race a subsequent registration
+ * that adopted the record — an adoption does not change the version — and delete the
+ * fresh record with the version this cleanup read. Only a record created by this store
+ * identity is removed.
+ */
+ private void removeOwnRegistrationRecord(String path) throws BookieException {
+ Optional result;
+ try {
+ result = store.get(path).get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
+ } catch (ExecutionException | TimeoutException e) {
+ throw new BookieException.MetadataStoreException(e);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new BookieException.MetadataStoreException(e);
+ }
+ if (result.isEmpty() || !result.get().getStat().isCreatedBySelf()) {
+ return;
+ }
+ try {
+ store.delete(path, Optional.of(result.get().getStat().getVersion()))
+ .get(BLOCKING_CALL_TIMEOUT, MILLISECONDS);
+ } catch (ExecutionException e) {
+ if (!(e.getCause() instanceof MetadataStoreException.NotFoundException)
+ && !(e.getCause() instanceof MetadataStoreException.BadVersionException)) {
+ throw new BookieException.MetadataStoreException(e);
+ }
+ // The record is already gone, or was rewritten in the meantime by a writer we
+ // cannot attribute to this registration: nothing of it remains to clean up.
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new BookieException.MetadataStoreException(e);
+ } catch (TimeoutException e) {
+ throw new BookieException.MetadataStoreException(e);
+ }
+ }
+
+ /**
+ * Marks a lock that could not be released, for example because it was lost to another
+ * owner, as expired, so that it is discarded by the lock manager instead of failing its
+ * own close path.
+ */
+ private void discardUnreleasableLock(ResourceLock lock) {
+ lock.getLockExpiredFuture().complete(null);
+ }
+
+ private Map> registrationMap(boolean readOnly) {
+ return readOnly ? bookieRegistrationReadOnly : bookieRegistration;
+ }
+
+ private String registrationPath(BookieId bookieId, boolean readOnly) {
+ return readOnly
+ ? bookieReadonlyRegistrationPath + "/" + bookieId
+ : bookieRegistrationPath + "/" + bookieId;
}
}
diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java
index f7562b60fd11a..d9e9dcf1fd032 100644
--- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java
+++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java
@@ -89,7 +89,13 @@ public CompletableFuture> acquireLock(String path, T value) {
lock.acquire(value).thenRun(() -> {
synchronized (LockManagerImpl.this) {
if (state == State.Ready) {
- locks.put(path, lock);
+ ResourceLockImpl replaced = locks.put(path, lock);
+ if (replaced != null && replaced != lock) {
+ // 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();
+ }
lock.getLockExpiredFuture().thenRun(() -> {
log.info().attr("path", path).log("Released resource lock");
synchronized (LockManagerImpl.this) {
diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java
index 60a847c472430..2275aec3f46d9 100644
--- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java
+++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java
@@ -107,30 +107,60 @@ public synchronized CompletableFuture release() {
revalidateTask.cancel(true);
}
- CompletableFuture result = new CompletableFuture<>();
+ // Serialize the deletion with an in-flight revalidation: without this, a
+ // revalidation whose read or write lands after the deletion would re-create the
+ // path and resurrect the resource.
+ return sequencer.sequential(() -> {
+ synchronized (ResourceLockImpl.this) {
+ if (state == State.Released) {
+ return CompletableFuture.completedFuture(null);
+ }
- store.delete(path, Optional.of(version))
- .thenRun(() -> {
- synchronized (ResourceLockImpl.this) {
- state = State.Released;
- }
- expiredFuture.complete(null);
- result.complete(null);
- }).exceptionally(ex -> {
- if (ex.getCause() instanceof MetadataStoreException.NotFoundException) {
- // The lock is not there on release. We can anyway proceed
- synchronized (ResourceLockImpl.this) {
- state = State.Released;
- }
- expiredFuture.complete(null);
- result.complete(null);
- } else {
- result.completeExceptionally(ex);
- }
- return null;
- });
+ long versionAtRelease = version;
+ CompletableFuture result = new CompletableFuture<>();
- return result;
+ store.delete(path, Optional.of(versionAtRelease))
+ .thenRun(() -> {
+ synchronized (ResourceLockImpl.this) {
+ state = State.Released;
+ }
+ expiredFuture.complete(null);
+ result.complete(null);
+ }).exceptionally(ex -> {
+ if (ex.getCause() instanceof MetadataStoreException.NotFoundException) {
+ // The lock is not there on release. We can anyway proceed
+ synchronized (ResourceLockImpl.this) {
+ state = State.Released;
+ }
+ expiredFuture.complete(null);
+ result.complete(null);
+ } else {
+ result.completeExceptionally(ex);
+ }
+ return null;
+ });
+ return result;
+ }
+ });
+ }
+
+ /**
+ * Marks the lock as released without touching the store, and cancels its pending
+ * revalidation. Used when the lock manager replaces this handle with a newer one for the
+ * same path: a revalidation of the superseded handle that is still scheduled or in flight
+ * would otherwise re-create the path and resurrect a torn-down resource. Does not undo a
+ * store write that is already in flight.
+ */
+ synchronized void retire() {
+ if (state == State.Released) {
+ return;
+ }
+ state = State.Released;
+ if (revalidateTask != null) {
+ revalidateTask.cancel(true);
+ revalidateTask = null;
+ }
+ expiredFuture.complete(null);
}
@Override
@@ -258,15 +288,33 @@ synchronized CompletableFuture silentRevalidateOnce() {
}
private synchronized CompletableFuture revalidate(T newValue) {
+ return revalidate(newValue, true);
+ }
+
+ /**
+ * One revalidation attempt. A {@code BadVersion} from any of its version-compared writes
+ * means a concurrent revalidation or acquire moved the record between the read and the
+ * write, not a conflict: with {@code retryOnBadVersion} the decision is made again from a
+ * fresh read, once, instead of failing or expiring the lock. A genuine foreign holder
+ * surfaces as a {@code LockBusy}, which is not retried.
+ */
+ private synchronized CompletableFuture revalidate(T newValue, boolean retryOnBadVersion) {
// Since the distributed lock has been expired, we don't need to revalidate it.
if (state != State.Valid && state != State.Init) {
return CompletableFuture.failedFuture(
new IllegalStateException("Lock was not in valid state: " + state));
}
- log.debug().attr("newValue", newValue).attr("version", version).log("doRevalidate");
return store.get(path)
.thenCompose(optGetResult -> {
if (!optGetResult.isPresent()) {
+ synchronized (ResourceLockImpl.this) {
+ if (state == State.Releasing || state == State.Released) {
+ // The lock was released while the read was in flight: the
+ // path is gone because the release deleted it, and it must
+ // not be resurrected by this revalidation.
+ return CompletableFuture.completedFuture(null);
+ }
+ }
// The lock just disappeared, try to acquire it again
// Reset the expectation on the version
setVersion(-1L);
@@ -295,11 +343,34 @@ private synchronized CompletableFuture revalidate(T newValue) {
// logical "owners" of the lock.
if (res.getStat().isCreatedBySelf()) {
- // If the new lock belongs to the same session, there's no
- // need to recreate it.
- version = res.getStat().getVersion();
- value = newValue;
- return CompletableFuture.completedFuture(null);
+ // The record was created by this client, but on stores where
+ // the creator identity outlives the session (the oxia backend
+ // scopes it to the client identity) it may still be bound to
+ // an expired session. Re-writing it with the expected version
+ // re-binds the record to the live session, and leaves the
+ // lock valid, also when it was adopted during an acquire.
+ byte[] payload;
+ try {
+ payload = serde.serialize(path, newValue);
+ } catch (Throwable t) {
+ return FutureUtils.exception(t);
+ }
+ return store.put(path, payload, Optional.of(res.getStat().getVersion()),
+ EnumSet.of(CreateOption.Ephemeral))
+ .thenAccept(putStat -> {
+ synchronized (ResourceLockImpl.this) {
+ // Record the fresh version even when the lock
+ // was released while the write was in flight:
+ // the write landed on the store, and the
+ // release deletes the record with this
+ // expectation.
+ version = putStat.getVersion();
+ value = newValue;
+ if (state != State.Releasing && state != State.Released) {
+ state = State.Valid;
+ }
+ }
+ });
} else {
// The lock needs to get recreated since it belong to an earlier
// session which maybe expiring soon
@@ -323,15 +394,21 @@ private synchronized CompletableFuture revalidate(T newValue) {
return FutureUtils.exception(
new LockBusyException("Resource at " + path + " is already locked"));
}
-
return store.delete(path, Optional.of(res.getStat().getVersion()))
.thenRun(() ->
- // Reset the expectation that the key is not there anymore
- setVersion(-1L)
+ // Reset the expectation that the key is not there anymore
+ setVersion(-1L)
)
.thenCompose(__ -> acquireWithNoRevalidation(newValue))
.thenRun(() -> log.info().attr("path", path).log("Successfully re-acquired lock"));
}
+ })
+ .exceptionallyCompose(ex -> {
+ Throwable cause = FutureUtil.unwrapCompletionException(ex);
+ if (retryOnBadVersion && cause instanceof BadVersionException) {
+ return revalidate(newValue, false);
+ }
+ return CompletableFuture.failedFuture(cause);
});
}
diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java
index ebd60bad5507d..5968246fd1051 100644
--- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java
+++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java
@@ -383,4 +383,72 @@ public void lockDeletedAndReacquired(String provider, Supplier urlSuppli
assertFalse(lock.getLockExpiredFuture().isDone());
});
}
+
+ /**
+ * A lock that adopts an already present record must come out valid: when the record is
+ * later deleted, the revalidation of the adopted handle re-creates it. Adoption used to
+ * leave the handle in the Init state, where revalidation is a no-op, so nothing would
+ * re-create the record and no expiry would be signalled.
+ */
+ @Test(dataProvider = "impl")
+ public void adoptedLockIsValidAndRecreatesThePathWhenDeleted(
+ String provider, Supplier urlSupplier) throws Exception {
+ @Cleanup
+ MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(),
+ MetadataStoreConfig.builder().fsyncEnable(false).build());
+
+ @Cleanup
+ CoordinationService coordinationService = new CoordinationServiceImpl(store);
+
+ @Cleanup
+ LockManager lockManager = coordinationService.getLockManager(String.class);
+
+ String key = newKey();
+ lockManager.acquireLock(key, "lock").join();
+
+ // The second acquire adopts the record the first one created, instead of failing.
+ ResourceLock adopted = lockManager.acquireLock(key, "lock").join();
+ assertFalse(adopted.getLockExpiredFuture().isDone());
+
+ store.delete(key, Optional.empty()).join();
+
+ Awaitility.await().untilAsserted(() -> {
+ Optional val = store.get(key).join();
+ assertTrue(val.isPresent());
+ assertFalse(adopted.getLockExpiredFuture().isDone());
+ });
+ }
+
+ /**
+ * A handle that the lock manager replaced must not resurrect the path: after the
+ * replacement is released, a revalidation of the superseded handle that was still
+ * scheduled or in flight must not re-create the record.
+ */
+ @Test(dataProvider = "impl")
+ public void supersededHandleDoesNotResurrectThePath(
+ String provider, Supplier urlSupplier) throws Exception {
+ @Cleanup
+ MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(),
+ MetadataStoreConfig.builder().fsyncEnable(false).build());
+
+ @Cleanup
+ CoordinationService coordinationService = new CoordinationServiceImpl(store);
+
+ @Cleanup
+ LockManager lockManager = coordinationService.getLockManager(String.class);
+
+ String key = newKey();
+ ResourceLock first = lockManager.acquireLock(key, "lock").join();
+ ResourceLock replacement = lockManager.acquireLock(key, "lock").join();
+
+ // The replacement retired the superseded handle.
+ Awaitility.await().until(first.getLockExpiredFuture()::isDone);
+
+ replacement.release().join();
+
+ Awaitility.await()
+ .during(1, TimeUnit.SECONDS)
+ .atMost(2, TimeUnit.SECONDS)
+ .untilAsserted(() -> assertFalse(store.get(key).join().isPresent()));
+ }
}
diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/FakeOxiaSessionMetadataStore.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/FakeOxiaSessionMetadataStore.java
new file mode 100644
index 0000000000000..4c3f31ab7c030
--- /dev/null
+++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/bookkeeper/FakeOxiaSessionMetadataStore.java
@@ -0,0 +1,322 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pulsar.metadata.bookkeeper;
+
+import io.opentelemetry.api.OpenTelemetry;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.NavigableMap;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.Set;
+import java.util.TreeMap;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.Executors;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.atomic.AtomicInteger;
+import lombok.Data;
+import org.apache.pulsar.metadata.api.GetResult;
+import org.apache.pulsar.metadata.api.MetadataStoreException;
+import org.apache.pulsar.metadata.api.Notification;
+import org.apache.pulsar.metadata.api.NotificationType;
+import org.apache.pulsar.metadata.api.Option;
+import org.apache.pulsar.metadata.api.OptionsHelper;
+import org.apache.pulsar.metadata.api.Stat;
+import org.apache.pulsar.metadata.api.extended.SessionEvent;
+import org.apache.pulsar.metadata.impl.AbstractMetadataStore;
+
+/**
+ * Minimal in-memory metadata store that models the metadata store semantics the session-loss
+ * driven bookie re-registration depends on:
+ *
+ *
the record identity is a store-level UUID, independent of the session: a record
+ * created by an earlier session of the same store instance still reports
+ * {@code createdBySelf == true} (cross-session adopt);
+ *
ephemeral puts lazily create a session and bind the record to it; a session loss
+ * purges (or, in a server-side grace window, keeps) the records of that session;
+ *
expected-version checked puts and deletes with BadVersion/NotFound semantics;
+ *
the first session established after an expiry is reported as
+ * {@link SessionEvent#SessionReestablished}, mirroring the session watcher mapping.
+ *
+ */
+class FakeOxiaSessionMetadataStore extends AbstractMetadataStore {
+
+ /** Session id used for foreign records: never matches a session of this store. */
+ private static final long FOREIGN_SESSION_ID = Long.MIN_VALUE;
+
+ @Data
+ static class Record {
+ final long version;
+ final byte[] data;
+ final long createdTimestamp;
+ final long modifiedTimestamp;
+ final boolean ephemeral;
+ final String creatorIdentity;
+ final long ownerSessionId;
+ }
+
+ private final String identity;
+ private final NavigableMap records = new TreeMap<>();
+
+ // guarded by this
+ private long nextSessionId = 1;
+ private Long currentSessionId;
+ private boolean seenExpired;
+
+ private final ScheduledExecutorService sessionEventDispatcher = Executors.newSingleThreadScheduledExecutor(
+ r -> new Thread(r, "fake-oxia-session-events"));
+
+ private volatile boolean unavailable = false;
+ /**
+ * When set, models the zk ownership semantics: an update (setData) neither transfers
+ * the ephemeral ownership nor changes the node type, and {@code createdBySelf} is
+ * session-scoped ({@code ephemeralOwner == current session}) instead of identity-scoped.
+ */
+ private volatile boolean zkOwnershipSemantics = false;
+ private final AtomicInteger deletesToFailWithBadVersion = new AtomicInteger();
+
+ FakeOxiaSessionMetadataStore(String identity) {
+ super("fake-oxia-session-store", OpenTelemetry.noop(), null, 1);
+ this.identity = identity;
+ }
+
+ // ---------------------------------------------------------------- store operations
+
+ @Override
+ protected synchronized CompletableFuture> storeGet(String path, Set