Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,13 @@ public CompletableFuture<ResourceLock<T>> acquireLock(String path, T value) {
lock.acquire(value).thenRun(() -> {
synchronized (LockManagerImpl.this) {
if (state == State.Ready) {
locks.put(path, lock);
ResourceLockImpl<T> 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();

@void-ptr974 void-ptr974 Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Potential deadlock on the ZooKeeper path. For example:

  1. Thread A enters synchronized old.release(), sends store.delete(), and is descheduled before attaching its .thenRun completion. It still holds the old lock monitor.
  2. ZooKeeper completes the delete. Another caller acquires the same path; its create callback holds the LockManagerImpl monitor and blocks here in replaced.retire() waiting for A.
  3. A resumes. Because the delete future is already complete, .thenRun runs inline on A. Its expiredFuture.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.

}
lock.getLockExpiredFuture().thenRun(() -> {
log.info().attr("path", path).log("Released resource lock");
synchronized (LockManagerImpl.this) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -107,30 +107,60 @@ public synchronized CompletableFuture<Void> release() {
revalidateTask.cancel(true);
}

CompletableFuture<Void> 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<Void> 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
Expand Down Expand Up @@ -258,15 +288,33 @@ synchronized CompletableFuture<Void> silentRevalidateOnce() {
}

private synchronized CompletableFuture<Void> 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<Void> 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);
Expand Down Expand Up @@ -295,11 +343,34 @@ private synchronized CompletableFuture<Void> 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
Expand All @@ -323,15 +394,21 @@ private synchronized CompletableFuture<Void> 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);
});
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -383,4 +383,72 @@ public void lockDeletedAndReacquired(String provider, Supplier<String> 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<String> urlSupplier) throws Exception {
@Cleanup
MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(),
MetadataStoreConfig.builder().fsyncEnable(false).build());

@Cleanup
CoordinationService coordinationService = new CoordinationServiceImpl(store);

@Cleanup
LockManager<String> 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<String> adopted = lockManager.acquireLock(key, "lock").join();
assertFalse(adopted.getLockExpiredFuture().isDone());

store.delete(key, Optional.empty()).join();

Awaitility.await().untilAsserted(() -> {
Optional<GetResult> 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<String> urlSupplier) throws Exception {
@Cleanup
MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(),
MetadataStoreConfig.builder().fsyncEnable(false).build());

@Cleanup
CoordinationService coordinationService = new CoordinationServiceImpl(store);

@Cleanup
LockManager<String> lockManager = coordinationService.getLockManager(String.class);

String key = newKey();
ResourceLock<String> first = lockManager.acquireLock(key, "lock").join();
ResourceLock<String> 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()));
}
}
Loading
Loading