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 cc66878d5f62 CAMEL-25061: camel-file - FileLockClusterService must be 
able to take the leadership again after a restart (#26944)
cc66878d5f62 is described below

commit cc66878d5f629cb82ea6c7e71194c5ba00510f68
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 14:09:34 2026 +0530

    CAMEL-25061: camel-file - FileLockClusterService must be able to take the 
leadership again after a restart (#26944)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../file/cluster/FileLockClusterService.java       |   3 +
 .../cluster/FileLockClusterServiceRestartTest.java | 109 +++++++++++++++++++++
 2 files changed, 112 insertions(+)

diff --git 
a/components/camel-file/src/main/java/org/apache/camel/component/file/cluster/FileLockClusterService.java
 
b/components/camel-file/src/main/java/org/apache/camel/component/file/cluster/FileLockClusterService.java
index d7f3f5289859..8333b9f5b841 100644
--- 
a/components/camel-file/src/main/java/org/apache/camel/component/file/cluster/FileLockClusterService.java
+++ 
b/components/camel-file/src/main/java/org/apache/camel/component/file/cluster/FileLockClusterService.java
@@ -264,6 +264,9 @@ public class FileLockClusterService extends 
AbstractCamelClusterService<FileLock
             } else {
                 clusterDataTaskExecutor.shutdown();
             }
+
+            // a new one is created when the service is started again
+            clusterDataTaskExecutor = null;
         }
     }
 
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/component/file/cluster/FileLockClusterServiceRestartTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/component/file/cluster/FileLockClusterServiceRestartTest.java
new file mode 100644
index 000000000000..c6adf0896940
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/component/file/cluster/FileLockClusterServiceRestartTest.java
@@ -0,0 +1,109 @@
+/*
+ * 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.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.concurrent.TimeUnit;
+
+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.BeforeEach;
+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.assertTrue;
+
+/**
+ * A {@link FileLockClusterService} that is stopped and started again in the 
same JVM, for example with the stop and
+ * start JMX operations of the cluster service, must be able to take the 
leadership again.
+ */
+public class FileLockClusterServiceRestartTest {
+
+    private static final String NAMESPACE = "ns";
+
+    @TempDir
+    Path root;
+
+    private CamelContext context;
+    private FileLockClusterService service;
+
+    @BeforeEach
+    public void setUp() throws Exception {
+        service = new FileLockClusterService();
+        service.setId("node-A");
+        service.setRoot(root.toString());
+        service.setAcquireLockDelay(100, TimeUnit.MILLISECONDS);
+        service.setAcquireLockInterval(200, TimeUnit.MILLISECONDS);
+
+        context = new DefaultCamelContext();
+        context.disableJMX();
+        context.addService(service);
+        context.start();
+    }
+
+    @AfterEach
+    public void tearDown() {
+        context.stop();
+    }
+
+    @Test
+    public void testLeaderAgainAfterServiceRestart() throws Exception {
+        CamelClusterView view = service.getView(NAMESPACE);
+        awaitLeader(view);
+
+        service.stop();
+        assertFalse(view.getLocalMember().isLeader());
+        assertTrue(isLockFree(root.resolve(NAMESPACE)), "the stopped service 
holds the lock");
+
+        // the view is started again with the service, and reads and writes 
the cluster data again
+        service.start();
+        awaitLeader(view);
+
+        // and it still works after another restart
+        service.stop();
+        service.start();
+        awaitLeader(view);
+    }
+
+    private void awaitLeader(CamelClusterView view) {
+        // the leader writes its heartbeat to the data file through the 
cluster data task executor
+        await().atMost(10, TimeUnit.SECONDS).until(() -> 
view.getLocalMember().isLeader()
+                && 
FileLockClusterUtils.readClusterLeaderInfo(root.resolve(NAMESPACE + ".dat")) != 
null);
+    }
+
+    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;
+        }
+    }
+}

Reply via email to