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