This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 3043cd20e231 CAMEL-25052: camel-file - FileLockClusterView stop must
end a running leadership check (#26932)
3043cd20e231 is described below
commit 3043cd20e231933d18f7e2dd3e7362d2de5ce22a
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 14:17:45 2026 +0530
CAMEL-25052: camel-file - FileLockClusterView stop must end a running
leadership check (#26932)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../file/cluster/FileLockClusterView.java | 276 +++++++++++++++-----
.../file/cluster/FileLockClusterViewStopTest.java | 284 +++++++++++++++++++++
2 files changed, 495 insertions(+), 65 deletions(-)
diff --git
a/components/camel-file/src/main/java/org/apache/camel/component/file/cluster/FileLockClusterView.java
b/components/camel-file/src/main/java/org/apache/camel/component/file/cluster/FileLockClusterView.java
index 4c24bb6ff2ac..cf9e2cb6ba8c 100644
---
a/components/camel-file/src/main/java/org/apache/camel/component/file/cluster/FileLockClusterView.java
+++
b/components/camel-file/src/main/java/org/apache/camel/component/file/cluster/FileLockClusterView.java
@@ -30,7 +30,6 @@ import java.util.Objects;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
-import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicReference;
@@ -54,13 +53,20 @@ public class FileLockClusterView extends
AbstractCamelClusterView {
private final Path leaderLockPath;
private final Path leaderDataPath;
private final AtomicReference<FileLockClusterLeaderInfo>
clusterLeaderInfoRef = new AtomicReference<>();
- private RandomAccessFile leaderLockFile;
- private RandomAccessFile leaderDataFile;
- private FileLock lock;
- private ScheduledFuture<?> task;
+ // Written under stateLock (published by acquireLock, taken over by doStop
or by the leadership-lost path of a
+ // current check). Volatile as the leadership check and the cluster data
tasks read them without the lock.
+ private volatile RandomAccessFile leaderLockFile;
+ private volatile RandomAccessFile leaderDataFile;
+ private volatile FileLock lock;
private int heartbeatTimeoutMultiplier;
private long acquireLockIntervalMilliseconds;
private FileLockClusterTaskExecutor clusterTaskExecutor;
+ // Guards generation and the hand-over of lock, leaderLockFile and
leaderDataFile between the leadership check and
+ // doStop. It is only held for short state changes, never during file I/O
or while listeners are notified.
+ private final ReentrantLock stateLock = new ReentrantLock();
+ // Incremented on every start and stop. A leadership check belongs to the
generation of the start that scheduled
+ // it, and it ends without taking the lock or rescheduling itself once
that generation is over.
+ private long generation;
FileLockClusterView(FileLockClusterService cluster, String namespace) {
super(cluster, namespace);
@@ -98,6 +104,7 @@ public class FileLockClusterView extends
AbstractCamelClusterView {
try {
contextStartLock.lock();
+ // Defensive only: doStop and the leadership-lost path always take
over and clear the lock and files
if (leaderLockFile != null) {
closeInternal();
fireLeadershipChangedEvent((CamelClusterMember) null);
@@ -124,76 +131,164 @@ public class FileLockClusterView extends
AbstractCamelClusterView {
heartbeatTimeoutMultiplier = service.getHeartbeatTimeoutMultiplier();
- scheduleTryLock(true);
+ long gen;
+ stateLock.lock();
+ try {
+ gen = ++generation;
+ } finally {
+ stateLock.unlock();
+ }
+ scheduleTryLock(true, gen);
localMember.setStatus(ClusterMemberStatus.STARTED);
}
@Override
protected void doStop() throws Exception {
- if (localMember.isLeader() && leaderDataFile != null) {
- clusterTaskExecutor.run(ThrowingHelper.wrapAsSupplier(new
ThrowingSupplier<Void, Throwable>() {
- @Override
- public Void get() throws Throwable {
- try {
- FileChannel channel = leaderDataFile.getChannel();
- channel.truncate(0);
- channel.force(true);
- } catch (Exception e) {
- // Log and ignore since we need to release the file
lock and do cleanup
- LOGGER.debug("Failed to truncate {} on {} stop",
leaderDataPath, getClass().getSimpleName(), e);
+ final boolean wasLeader;
+ final FileLock heldLock;
+ final RandomAccessFile heldLockFile;
+ final RandomAccessFile heldDataFile;
+ stateLock.lock();
+ try {
+ // A leadership check that is running now can no longer publish a
lock it acquires: it re-checks the
+ // generation under stateLock and releases the lock instead. So
the lock and files taken here are the
+ // only ones this view holds.
+ generation++;
+ wasLeader = localMember.isLeader();
+ localMember.setStatus(ClusterMemberStatus.STOPPED);
+ heldLock = lock;
+ heldLockFile = leaderLockFile;
+ heldDataFile = leaderDataFile;
+ lock = null;
+ leaderLockFile = null;
+ leaderDataFile = null;
+ } finally {
+ stateLock.unlock();
+ }
+
+ try {
+ if (wasLeader && heldDataFile != null) {
+ clusterTaskExecutor.run(ThrowingHelper.wrapAsSupplier(new
ThrowingSupplier<Void, Throwable>() {
+ @Override
+ public Void get() throws Throwable {
+ try {
+ FileChannel channel = heldDataFile.getChannel();
+ channel.truncate(0);
+ channel.force(true);
+ } catch (Exception e) {
+ // Log and ignore since we need to release the
file lock and do cleanup
+ LOGGER.debug("Failed to truncate {} on {} stop",
leaderDataPath, getClass().getSimpleName(), e);
+ }
+ return null;
}
- return null;
- }
- }));
+ }));
+ }
+ } finally {
+ // The fields no longer reference the lock and files, so they must
be released even if the truncate task
+ // timed out (for example on a hanging NFS mount), otherwise the
stopped view would keep the lock
+ releaseFileLock(heldLock);
+ closeFile(heldLockFile);
+ closeFile(heldDataFile);
+ clusterLeaderInfoRef.set(null);
}
+ }
- closeInternal();
- localMember.setStatus(ClusterMemberStatus.STOPPED);
- clusterLeaderInfoRef.set(null);
+ /**
+ * The lock, lock file and data file taken over from this view by {@link
#takeOver(long)}.
+ */
+ private record HeldLock(FileLock lock, RandomAccessFile lockFile,
RandomAccessFile dataFile) {
+ void release() {
+ releaseFileLock(lock);
+ closeFile(lockFile);
+ closeFile(dataFile);
+ }
}
- private void closeInternal() {
- if (task != null) {
- task.cancel(true);
+ /**
+ * If generation {@code gen} is still current, sets the member to FOLLOWER
and takes the lock and files over from
+ * the view, so that the caller releases them outside stateLock. Returns
null if the view has been stopped or
+ * started again since, in which case doStop (or the check of the newer
generation) owns them.
+ */
+ private HeldLock takeOver(long gen) {
+ stateLock.lock();
+ try {
+ if (gen != generation) {
+ return null;
+ }
+ HeldLock held = new HeldLock(lock, leaderLockFile, leaderDataFile);
+ lock = null;
+ leaderLockFile = null;
+ leaderDataFile = null;
+ localMember.setStatus(ClusterMemberStatus.FOLLOWER);
+ return held;
+ } finally {
+ stateLock.unlock();
}
+ }
+ private void closeInternal() {
releaseFileLock();
closeLockFiles();
}
private void closeLockFiles() {
- if (leaderLockFile != null) {
+ closeFile(leaderLockFile);
+ leaderLockFile = null;
+ closeFile(leaderDataFile);
+ leaderDataFile = null;
+ }
+
+ private static void closeFile(RandomAccessFile file) {
+ if (file != null) {
try {
- leaderLockFile.close();
+ file.close();
} catch (Exception ignore) {
LOGGER.warn("{}", ignore.getMessage(), ignore);
}
- leaderLockFile = null;
}
+ }
+
+ private void releaseFileLock() {
+ releaseFileLock(lock);
+ }
- if (leaderDataFile != null) {
+ private static void releaseFileLock(FileLock fileLock) {
+ if (fileLock != null) {
try {
- leaderDataFile.close();
+ fileLock.release();
} catch (Exception ignore) {
LOGGER.warn("{}", ignore.getMessage(), ignore);
}
- leaderDataFile = null;
}
}
- private void releaseFileLock() {
- if (lock != null) {
- try {
- lock.release();
- } catch (Exception ignore) {
- LOGGER.warn("{}", ignore.getMessage(), ignore);
+ private boolean isCurrentGeneration(long gen) {
+ stateLock.lock();
+ try {
+ return gen == generation;
+ } finally {
+ stateLock.unlock();
+ }
+ }
+
+ private boolean setFollower(long gen) {
+ stateLock.lock();
+ try {
+ if (gen != generation) {
+ return false;
}
+ localMember.setStatus(ClusterMemberStatus.FOLLOWER);
+ return true;
+ } finally {
+ stateLock.unlock();
}
}
- private void tryLock() {
- if (isStarting() || isStarted()) {
+ private void tryLock(long gen) {
+ // A check scheduled by an earlier start, or by a start that has been
stopped since, ends here without
+ // rescheduling itself, so that a stop ends the chain and a quick stop
and start does not add a second chain
+ if (isCurrentGeneration(gen) && (isStarting() || isStarted())) {
Exception reason = null;
try {
@@ -211,19 +306,23 @@ public class FileLockClusterView extends
AbstractCamelClusterView {
// Non-null lock at this point signifies leadership has been
lost or relinquished
if (lock != null) {
- LOGGER.info("Lock on file {} lost (lock={},
cluster-member-id={})", leaderLockPath, lock,
- localMember.getUuid());
- localMember.setStatus(ClusterMemberStatus.FOLLOWER);
- fireLeadershipChangedEvent((CamelClusterMember) null);
- clusterLeaderInfoRef.set(null);
- releaseFileLock();
- closeLockFiles();
- lock = null;
+ // Only if the view has not been stopped meanwhile: doStop
then owns the lock and files, and the
+ // member stays STOPPED
+ HeldLock held = takeOver(gen);
+ if (held != null) {
+ LOGGER.info("Lock on file {} lost (lock={},
cluster-member-id={})", leaderLockPath, held.lock(),
+ localMember.getUuid());
+ fireLeadershipChangedEvent((CamelClusterMember) null);
+ clusterLeaderInfoRef.set(null);
+ held.release();
+ }
return;
}
- // Must be follower to reach here
- localMember.setStatus(ClusterMemberStatus.FOLLOWER);
+ // Must be follower to reach here. A stopped view stays
STOPPED and its check ends here
+ if (!setFollower(gen)) {
+ return;
+ }
// Get & update cluster leader state
LOGGER.debug("Reading cluster leader state from {}",
leaderDataPath);
@@ -247,18 +346,9 @@ public class FileLockClusterView extends
AbstractCamelClusterView {
// Attempt to obtain cluster leadership
LOGGER.debug("Try to acquire a lock on {}
(cluster-member-id={})", leaderLockPath, localMember.getUuid());
- lock = null;
- leaderLockFile = createRandomAccessFile(leaderLockPath);
- leaderDataFile = createRandomAccessFile(leaderDataPath);
- if (leaderLockFile != null && leaderDataFile != null) {
- lock = leaderLockFile.getChannel().tryLock(0,
Math.max(1, leaderLockFile.getChannel().size()), false);
- }
-
- if (lockIsValid()) {
+ if (acquireLock(gen)) {
LOGGER.info("Lock on file {} acquired (lock={},
cluster-member-id={})", leaderLockPath, lock,
localMember.getUuid());
- localMember.setStatus(ClusterMemberStatus.LEADER);
- clusterLeaderInfoRef.set(null);
fireLeadershipChangedEvent(localMember);
writeClusterLeaderInfo(true);
} else {
@@ -275,11 +365,63 @@ public class FileLockClusterView extends
AbstractCamelClusterView {
if (lock == null) {
LOGGER.debug("Lock on file {} not acquired
(cluster-member-id={})", leaderLockPath, localMember.getUuid(),
reason);
- closeLockFiles();
}
- scheduleTryLock(false);
+ if (isCurrentGeneration(gen)) {
+ scheduleTryLock(false, gen);
+ }
+ }
+ }
+ }
+
+ /**
+ * Opens the lock and data files and tries to lock the lock file. The
files and the lock are only published to this
+ * view, and the member only becomes leader, if the view has not been
stopped since the check of generation
+ * {@code gen} started. Otherwise the lock is released again. The file I/O
runs without holding stateLock.
+ */
+ private boolean acquireLock(long gen) throws Exception {
+ if (!isCurrentGeneration(gen)) {
+ // stopped while the cluster data was read: do not take the lock,
even briefly
+ return false;
+ }
+ RandomAccessFile newLockFile = null;
+ RandomAccessFile newDataFile = null;
+ FileLock newLock = null;
+ boolean published = false;
+ try {
+ newLockFile = createRandomAccessFile(leaderLockPath);
+ newDataFile = createRandomAccessFile(leaderDataPath);
+ if (newLockFile != null && newDataFile != null) {
+ newLock = newLockFile.getChannel().tryLock(0, Math.max(1,
newLockFile.getChannel().size()), false);
+ }
+
+ if (lockIsValid(newLock)) {
+ stateLock.lock();
+ try {
+ if (gen == generation) {
+ lock = newLock;
+ leaderLockFile = newLockFile;
+ leaderDataFile = newDataFile;
+ localMember.setStatus(ClusterMemberStatus.LEADER);
+ clusterLeaderInfoRef.set(null);
+ published = true;
+ }
+ } finally {
+ stateLock.unlock();
+ }
+
+ if (!published) {
+ LOGGER.debug("Lock on file {} acquired after the view was
stopped, releasing it (cluster-member-id={})",
+ leaderLockPath, localMember.getUuid());
+ }
+ }
+ } finally {
+ if (!published) {
+ releaseFileLock(newLock);
+ closeFile(newLockFile);
+ closeFile(newDataFile);
}
}
+ return published;
}
void validateAcquireLockInterval(FileLockClusterLeaderInfo
clusterLeaderInfo) {
@@ -293,7 +435,7 @@ public class FileLockClusterView extends
AbstractCamelClusterView {
}
}
- void scheduleTryLock(boolean isFirstRun) {
+ void scheduleTryLock(boolean isFirstRun, long gen) {
long offset = System.currentTimeMillis() %
acquireLockIntervalMilliseconds;
long delay = acquireLockIntervalMilliseconds - offset;
if (delay <= 0) {
@@ -320,7 +462,7 @@ public class FileLockClusterView extends
AbstractCamelClusterView {
getClusterService().unwrap(FileLockClusterService.class)
.getExecutor()
- .schedule(this::tryLock, delay, TimeUnit.MILLISECONDS);
+ .schedule(() -> tryLock(gen), delay, TimeUnit.MILLISECONDS);
}
boolean isLeaderStale(FileLockClusterLeaderInfo clusterLeaderInfo,
FileLockClusterLeaderInfo previousClusterLeaderInfo) {
@@ -405,7 +547,11 @@ public class FileLockClusterView extends
AbstractCamelClusterView {
}
boolean lockIsValid() throws ExecutionException, TimeoutException {
- if (lock != null && lock.isValid()) {
+ return lockIsValid(lock);
+ }
+
+ private boolean lockIsValid(FileLock fileLock) throws ExecutionException,
TimeoutException {
+ if (fileLock != null && fileLock.isValid()) {
return clusterTaskExecutor.run(ThrowingHelper.wrapAsSupplier(new
ThrowingSupplier<Boolean, Throwable>() {
@Override
public Boolean get() throws Throwable {
diff --git
a/core/camel-core/src/test/java/org/apache/camel/component/file/cluster/FileLockClusterViewStopTest.java
b/core/camel-core/src/test/java/org/apache/camel/component/file/cluster/FileLockClusterViewStopTest.java
new file mode 100644
index 000000000000..6bf113a7114a
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/component/file/cluster/FileLockClusterViewStopTest.java
@@ -0,0 +1,284 @@
+/*
+ * 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.camel.component.file.cluster;
+
+import java.io.RandomAccessFile;
+import java.nio.channels.FileChannel;
+import java.nio.channels.FileLock;
+import java.nio.channels.OverlappingFileLockException;
+import java.nio.file.Path;
+import java.nio.file.StandardOpenOption;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.cluster.CamelClusterView;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A stopped {@link FileLockClusterView} must end its leadership check: a
check that is still running when the view
+ * stops must not take the lock or report leadership afterwards, and a quick
stop and start must not leave two checks
+ * running. The stop must release the lock even if the cluster data cannot be
written.
+ */
+public class FileLockClusterViewStopTest {
+
+ private static final String NAMESPACE = "ns";
+
+ @TempDir
+ Path root;
+
+ private final List<CamelContext> contexts = new ArrayList<>();
+
+ @AfterEach
+ public void stopContexts() {
+ contexts.forEach(CamelContext::stop);
+ }
+
+ @Test
+ public void testViewStoppedDuringLeadershipCheckDoesNotKeepTheLock()
throws Exception {
+ CamelContext contextA = startContext();
+ FileLockClusterService serviceA = newService("A");
+ contextA.addService(serviceA);
+ CamelClusterView viewA = serviceA.getView(NAMESPACE);
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
viewA.getLocalMember().isLeader()
+ &&
FileLockClusterUtils.readClusterLeaderInfo(root.resolve(NAMESPACE + ".dat")) !=
null);
+
+ // B follows A, and its next check after A has gone is held just
before it opens the lock file
+ HookedFileLockClusterService serviceB = new
HookedFileLockClusterService();
+ configure(serviceB, "B");
+ CamelContext contextB = startContext();
+ contextB.addService(serviceB);
+ CamelClusterView viewB = serviceB.getView(NAMESPACE);
+
+ contextA.stop();
+ assertTrue(serviceB.reached.await(10, TimeUnit.SECONDS), "B did not
start a leadership check");
+
+ // B's view is released by its last user (a master route or clustered
route stopping) during that check
+ serviceB.releaseView(viewB);
+ serviceB.release.countDown();
+ // the check runs on the single thread of B's executor: once this task
has run, the check is over
+ serviceB.getExecutor().submit(() -> {
+ }).get(30, TimeUnit.SECONDS);
+
+ // the check of the stopped view did not schedule another one
+ assertTrue(serviceB.scheduler.getQueue().isEmpty(), "the stopped view
still schedules leadership checks");
+ assertFalse(viewB.getLocalMember().isLeader(), "the stopped view
reports leadership");
+ assertFalse(serviceB.isLeader(NAMESPACE), "the stopped view reports
leadership");
+ assertTrue(isLockFree(root.resolve(NAMESPACE)), "the stopped view
holds the lock");
+
+ // so another member can take over
+ CamelContext contextC = startContext();
+ FileLockClusterService serviceC = newService("C");
+ contextC.addService(serviceC);
+ CamelClusterView viewC = serviceC.getView(NAMESPACE);
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
viewC.getLocalMember().isLeader());
+ }
+
+ @Test
+ public void testQuickRestartDoesNotAddLeadershipChecks() throws Exception {
+ HookedFileLockClusterService service = new
HookedFileLockClusterService();
+ configure(service, "A");
+ service.release.countDown();
+ CamelContext context = startContext();
+ context.addService(service);
+ CamelClusterView view = service.getView(NAMESPACE);
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
view.getLocalMember().isLeader());
+
+ // a master route or clustered route stopped and started again within
one acquireLockInterval
+ for (int i = 0; i < 3; i++) {
+ service.releaseView(view);
+ service.getView(NAMESPACE);
+ }
+
+ // every running check has exactly one run scheduled at a time. The
four starts (the first one and three
+ // restarts) queued four checks. The checks of the earlier starts end
at their next run without rescheduling,
+ // so after one interval only the check of the last start is left
+ ScheduledThreadPoolExecutor executor = service.scheduler;
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
executor.getQueue().size() == 1);
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
view.getLocalMember().isLeader());
+ await().during(2, TimeUnit.SECONDS).atMost(3,
TimeUnit.SECONDS).until(() -> executor.getQueue().size() <= 1);
+ }
+
+ @Test
+ public void testStopReleasesTheLockWhenTheDataFileCannotBeTruncated()
throws Exception {
+ HangingFileLockClusterService service = new
HangingFileLockClusterService();
+ configure(service, "A");
+ CamelContext context = startContext();
+ context.addService(service);
+ CamelClusterView view = service.getView(NAMESPACE);
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
view.getLocalMember().isLeader()
+ &&
FileLockClusterUtils.readClusterLeaderInfo(root.resolve(NAMESPACE + ".dat")) !=
null);
+
+ // the cluster data storage hangs (for example an NFS mount), so the
truncate on stop times out
+ try {
+ service.setClusterDataTaskMaxAttempts(1);
+ service.setClusterDataTaskTimeout(200, TimeUnit.MILLISECONDS);
+ service.hang();
+ assertThrows(RuntimeException.class, () ->
service.releaseView(view),
+ "the truncate of the data file did not time out");
+ } finally {
+ service.unhang();
+ }
+
+ assertFalse(view.getLocalMember().isLeader(), "the stopped view
reports leadership");
+ assertTrue(isLockFree(root.resolve(NAMESPACE)), "the stopped view
holds the lock");
+ }
+
+ private CamelContext startContext() {
+ CamelContext context = new DefaultCamelContext();
+ context.disableJMX();
+ contexts.add(context);
+ context.start();
+ return context;
+ }
+
+ private FileLockClusterService newService(String id) {
+ FileLockClusterService service = new FileLockClusterService();
+ configure(service, id);
+ return service;
+ }
+
+ private void configure(FileLockClusterService service, String id) {
+ service.setId("node-" + id);
+ service.setRoot(root.toString());
+ service.setAcquireLockDelay(100, TimeUnit.MILLISECONDS);
+ service.setAcquireLockInterval(500, TimeUnit.MILLISECONDS);
+ service.setClusterDataTaskTimeout(30, TimeUnit.SECONDS);
+ }
+
+ private static boolean isLockFree(Path lockFile) throws Exception {
+ try (FileChannel channel = FileChannel.open(lockFile,
StandardOpenOption.READ, StandardOpenOption.WRITE)) {
+ FileLock lock = channel.tryLock();
+ if (lock == null) {
+ return false;
+ }
+ lock.release();
+ return true;
+ } catch (OverlappingFileLockException e) {
+ // held by another channel of this JVM
+ return false;
+ }
+ }
+
+ /**
+ * Once {@link #hang()} is called, no leadership check runs any more, and
the cluster data tasks run on an executor
+ * whose only thread is blocked, so that they time out like the I/O on a
hanging network file system.
+ */
+ private static final class HangingFileLockClusterService extends
FileLockClusterService {
+ private final CountDownLatch unblock = new CountDownLatch(1);
+ private final ScheduledThreadPoolExecutor scheduler = new
ScheduledThreadPoolExecutor(1);
+ private final ExecutorService hangingExecutor =
Executors.newSingleThreadExecutor();
+ private volatile boolean hanging;
+
+ void hang() throws InterruptedException {
+ // wait until no leadership check is running, and hold the next
ones
+ CountDownLatch paused = new CountDownLatch(1);
+ scheduler.execute(() -> {
+ paused.countDown();
+ awaitUnblock();
+ });
+ assertTrue(paused.await(10, TimeUnit.SECONDS), "the leadership
checks were not paused");
+ hangingExecutor.execute(this::awaitUnblock);
+ hanging = true;
+ }
+
+ void unhang() {
+ hanging = false;
+ unblock.countDown();
+ hangingExecutor.shutdown();
+ }
+
+ private void awaitUnblock() {
+ try {
+ unblock.await(30, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+
+ @Override
+ ScheduledExecutorService getExecutor() {
+ return scheduler;
+ }
+
+ @Override
+ ExecutorService getClusterDataTaskExecutor() {
+ return hanging ? hangingExecutor :
super.getClusterDataTaskExecutor();
+ }
+
+ @Override
+ protected void doStop() throws Exception {
+ super.doStop();
+ scheduler.shutdownNow();
+ }
+ }
+
+ /**
+ * Holds the first attempt to open the lock file until {@link #release} is
counted down, and schedules the
+ * leadership checks on an executor the test can inspect.
+ */
+ private static final class HookedFileLockClusterService extends
FileLockClusterService {
+ final CountDownLatch reached = new CountDownLatch(1);
+ final CountDownLatch release = new CountDownLatch(1);
+ final ScheduledThreadPoolExecutor scheduler = new
ScheduledThreadPoolExecutor(1);
+
+ @Override
+ protected FileLockClusterView createView(String namespace) {
+ return new FileLockClusterView(this, namespace) {
+ @Override
+ RandomAccessFile createRandomAccessFile(Path path) throws
ExecutionException, TimeoutException {
+ if (path.getFileName().toString().equals(NAMESPACE) &&
reached.getCount() > 0) {
+ reached.countDown();
+ try {
+ release.await(30, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ return super.createRandomAccessFile(path);
+ }
+ };
+ }
+
+ @Override
+ ScheduledExecutorService getExecutor() {
+ return scheduler;
+ }
+
+ @Override
+ protected void doStop() throws Exception {
+ super.doStop();
+ scheduler.shutdownNow();
+ }
+ }
+}