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();
+        }
+    }
+}

Reply via email to