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 8843e8884508 CAMEL-25213: camel-support - FileStateRepository must not
lose the stored state when it is stopped or rewritten
8843e8884508 is described below
commit 8843e8884508c23b8991ff488cdd16524d788630
Author: allthingssecurity <[email protected]>
AuthorDate: Fri Oct 2 02:13:45 2026 +0530
CAMEL-25213: camel-support - FileStateRepository must not lose the stored
state when it is stopped or rewritten
FileStateRepository (used e.g. as camel-kafka offsetRepository) could lose
its stored state:
- doStop rewrote the file and cleared the map without cacheAndStoreLock, so
a
concurrent setState caused a ConcurrentModificationException and lost
keys.
- trunkStore truncated the file before rewriting it, so a failure or crash
during the rewrite left it empty or partial.
- appendToStore wrote a line in several writes, and loadStore failed on a
line without '=', so a torn last line prevented a restart.
doStop now holds the lock. trunkStore writes to <file>.tmp with the store's
POSIX permissions, syncs it, and moves it atomically over the real file
(resolving symlinks, so the link is kept). Lines are appended in a single
write, and loadStore skips an incomplete line with a WARN. No API or file
format change.
Closes #27162
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../processor/state/FileStateRepositoryTest.java | 264 +++++++++++++++++++++
.../processor/state/FileStateRepository.java | 73 ++++--
2 files changed, 323 insertions(+), 14 deletions(-)
diff --git
a/core/camel-core/src/test/java/org/apache/camel/support/processor/state/FileStateRepositoryTest.java
b/core/camel-core/src/test/java/org/apache/camel/support/processor/state/FileStateRepositoryTest.java
index 6dab400afb4c..896787ae36af 100644
---
a/core/camel-core/src/test/java/org/apache/camel/support/processor/state/FileStateRepositoryTest.java
+++
b/core/camel-core/src/test/java/org/apache/camel/support/processor/state/FileStateRepositoryTest.java
@@ -17,15 +17,36 @@
package org.apache.camel.support.processor.state;
import java.io.File;
+import java.io.IOException;
+import java.nio.file.FileSystems;
+import java.nio.file.Files;
+import java.nio.file.LinkOption;
+import java.nio.file.Path;
+import java.nio.file.attribute.PosixFilePermission;
+import java.nio.file.attribute.PosixFilePermissions;
+import java.util.AbstractSet;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
import org.apache.camel.TestSupport;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.condition.DisabledOnOs;
+import org.junit.jupiter.api.condition.OS;
import static
org.apache.camel.support.processor.state.FileStateRepository.fileStateRepository;
+import static org.awaitility.Awaitility.await;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.junit.jupiter.api.Assumptions.assumeTrue;
public class FileStateRepositoryTest extends TestSupport {
@@ -128,9 +149,252 @@ public class FileStateRepositoryTest extends TestSupport {
assertTrue(repositoryStore.length() < previousSize);
}
+ @Test
+ public void shouldSkipIncompleteLineWhenLoading() throws Exception {
+ // Given a store whose last line is incomplete (such as after a crash
while appending)
+ Files.writeString(repositoryStore.toPath(),
"key1=value1\nkey2=value2\nkey3");
+
+ // When starting a FileStateRepository with that store
+ FileStateRepository repository = createRepository();
+
+ // Then it starts and has the complete entries
+ assertEquals("value1", repository.getState("key1"));
+ assertEquals("value2", repository.getState("key2"));
+ assertNull(repository.getState("key3"));
+ }
+
+ @Test
+ public void shouldKeepStoreWhenRewritingFails() throws Exception {
+ // Given a FileStateRepository with some content, whose 1st level
cache fails while it is written to the store
+ FailingMap cache = new FailingMap();
+ FileStateRepository repository = fileStateRepository(repositoryStore,
cache);
+ repository.start();
+ repository.setState("key1", "value1");
+ repository.setState("key2", "value2");
+ repository.setState("key3", "value3");
+
+ // When stopping it (which rewrites the store) and the rewrite fails
half way
+ cache.fail = true;
+ assertThrows(IllegalStateException.class, repository::stop);
+
+ // Then the store still has the previous content
+ FileStateRepository newRepository = createRepository();
+ assertEquals("value1", newRepository.getState("key1"));
+ assertEquals("value2", newRepository.getState("key2"));
+ assertEquals("value3", newRepository.getState("key3"));
+ }
+
+ @Test
+ public void shouldNotLoseStateWhenUpdatedWhileStopping() throws Exception {
+ // Given a FileStateRepository with some content
+ CountDownLatch rewriting = new CountDownLatch(1);
+ CountDownLatch resume = new CountDownLatch(1);
+ PausingMap cache = new PausingMap(rewriting, resume);
+ FileStateRepository repository = fileStateRepository(repositoryStore,
cache);
+ repository.start();
+ for (int i = 0; i < 5; i++) {
+ repository.setState("key" + i, "value" + i);
+ }
+
+ // When the state is updated by another thread while the repository is
stopping (rewriting the store)
+ AtomicReference<Exception> stopFailure = new AtomicReference<>();
+ cache.pauseThread = new Thread(() -> {
+ try {
+ repository.stop();
+ } catch (Exception e) {
+ stopFailure.set(e);
+ }
+ });
+ cache.pauseThread.start();
+ assertTrue(rewriting.await(10, TimeUnit.SECONDS));
+
+ Thread updater = new Thread(() -> repository.setState("key5",
"value5"));
+ updater.start();
+ // the update must wait for the stop to complete
+ await().atMost(10, TimeUnit.SECONDS).until(() -> updater.getState() ==
Thread.State.WAITING
+ || updater.getState() == Thread.State.TERMINATED);
+ resume.countDown();
+ cache.pauseThread.join(10000);
+ updater.join(10000);
+
+ // Then stopping succeeded, and no state is lost
+ assertNull(stopFailure.get());
+ FileStateRepository newRepository = createRepository();
+ for (int i = 0; i < 6; i++) {
+ assertEquals("value" + i, newRepository.getState("key" + i));
+ }
+ }
+
+ @Test
+ @DisabledOnOs(OS.WINDOWS)
+ public void shouldKeepSymlinkedStoreWhenRewriting() throws Exception {
+ // Given a store which is a symbolic link to a file in another
directory
+ Path realStore = testDirectory("real",
true).resolve("file-state-repository.dat");
+ Files.writeString(realStore, "key1=value1\n");
+ Path link = repositoryStore.toPath();
+ try {
+ Files.createSymbolicLink(link, realStore);
+ } catch (UnsupportedOperationException | IOException e) {
+ assumeTrue(false, "Symbolic links are not supported: " +
e.getMessage());
+ }
+
+ // When updating the state and stopping the repository (which rewrites
the store)
+ FileStateRepository repository = createRepository();
+ repository.setState("key2", "value2");
+ repository.stop();
+
+ // Then the store is still the link, the real file has the state, and
no temporary file is left
+ assertTrue(Files.isSymbolicLink(link));
+ assertEquals(realStore, Files.readSymbolicLink(link));
+ assertTrue(Files.readString(realStore).contains("key1=value1\n"));
+ assertTrue(Files.readString(realStore).contains("key2=value2\n"));
+ assertFalse(Files.exists(Path.of(realStore + ".tmp")));
+ assertFalse(Files.exists(Path.of(link + ".tmp"),
LinkOption.NOFOLLOW_LINKS));
+ FileStateRepository newRepository = createRepository();
+ assertEquals("value1", newRepository.getState("key1"));
+ assertEquals("value2", newRepository.getState("key2"));
+ }
+
+ @Test
+ public void shouldKeepStorePermissionsWhenRewriting() throws Exception {
+
assumeTrue(FileSystems.getDefault().supportedFileAttributeViews().contains("posix"),
+ "POSIX file permissions are not supported");
+
+ // Given a store that only its owner can read
+ Path store = repositoryStore.toPath();
+ Files.writeString(store, "key1=value1\n");
+ Set<PosixFilePermission> permissions =
PosixFilePermissions.fromString("rw-------");
+ Files.setPosixFilePermissions(store, permissions);
+ CountDownLatch rewriting = new CountDownLatch(1);
+ CountDownLatch resume = new CountDownLatch(1);
+ PausingMap cache = new PausingMap(rewriting, resume);
+ FileStateRepository repository = fileStateRepository(repositoryStore,
cache);
+ repository.start();
+ repository.setState("key2", "value2");
+
+ // When stopping the repository (which rewrites the store), paused
while the state is written
+ AtomicReference<Exception> stopFailure = new AtomicReference<>();
+ cache.pauseThread = new Thread(() -> {
+ try {
+ repository.stop();
+ } catch (Exception e) {
+ stopFailure.set(e);
+ }
+ });
+ cache.pauseThread.start();
+ assertTrue(rewriting.await(10, TimeUnit.SECONDS));
+ Set<PosixFilePermission> whileWriting =
Files.getPosixFilePermissions(Path.of(store + ".tmp"));
+ resume.countDown();
+ cache.pauseThread.join(10000);
+
+ // Then the temporary file already had the permissions of the store
while the state was written to it,
+ // and the rewritten store keeps them
+ assertNull(stopFailure.get());
+ assertEquals(permissions, whileWriting);
+ assertEquals(permissions, Files.getPosixFilePermissions(store));
+ assertEquals("value2", createRepository().getState("key2"));
+ }
+
private FileStateRepository createRepository() {
FileStateRepository repository = fileStateRepository(repositoryStore);
repository.start();
return repository;
}
+
+ /**
+ * A cache which fails while its entries are iterated (when the store is
rewritten), if requested.
+ */
+ private static final class FailingMap extends HashMap<String, String> {
+ private volatile boolean fail;
+
+ @Override
+ public Set<Map.Entry<String, String>> entrySet() {
+ Set<Map.Entry<String, String>> entries = super.entrySet();
+ if (!fail) {
+ return entries;
+ }
+ return new AbstractSet<>() {
+ @Override
+ public Iterator<Map.Entry<String, String>> iterator() {
+ Iterator<Map.Entry<String, String>> it =
entries.iterator();
+ return new Iterator<>() {
+ private int count;
+
+ @Override
+ public boolean hasNext() {
+ return it.hasNext();
+ }
+
+ @Override
+ public Map.Entry<String, String> next() {
+ if (count++ == 1) {
+ throw new IllegalStateException("Simulated
failure while rewriting the store");
+ }
+ return it.next();
+ }
+ };
+ }
+
+ @Override
+ public int size() {
+ return entries.size();
+ }
+ };
+ }
+ }
+
+ /**
+ * A cache whose iteration (when the store is rewritten) by the given
thread pauses after the first entry.
+ */
+ private static final class PausingMap extends HashMap<String, String> {
+ private final CountDownLatch rewriting;
+ private final CountDownLatch resume;
+ private volatile Thread pauseThread;
+
+ private PausingMap(CountDownLatch rewriting, CountDownLatch resume) {
+ this.rewriting = rewriting;
+ this.resume = resume;
+ }
+
+ @Override
+ public Set<Map.Entry<String, String>> entrySet() {
+ Set<Map.Entry<String, String>> entries = super.entrySet();
+ if (Thread.currentThread() != pauseThread) {
+ return entries;
+ }
+ return new AbstractSet<>() {
+ @Override
+ public Iterator<Map.Entry<String, String>> iterator() {
+ Iterator<Map.Entry<String, String>> it =
entries.iterator();
+ return new Iterator<>() {
+ private int count;
+
+ @Override
+ public boolean hasNext() {
+ return it.hasNext();
+ }
+
+ @Override
+ public Map.Entry<String, String> next() {
+ Map.Entry<String, String> answer = it.next();
+ if (count++ == 0) {
+ rewriting.countDown();
+ try {
+ resume.await(10, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ return answer;
+ }
+ };
+ }
+
+ @Override
+ public int size() {
+ return entries.size();
+ }
+ };
+ }
+ }
}
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/processor/state/FileStateRepository.java
b/core/camel-support/src/main/java/org/apache/camel/support/processor/state/FileStateRepository.java
index b5016c428987..1492796a0c51 100644
---
a/core/camel-support/src/main/java/org/apache/camel/support/processor/state/FileStateRepository.java
+++
b/core/camel-support/src/main/java/org/apache/camel/support/processor/state/FileStateRepository.java
@@ -19,9 +19,17 @@ package org.apache.camel.support.processor.state;
import java.io.File;
import java.io.FileOutputStream;
import java.io.IOException;
+import java.nio.file.AtomicMoveNotSupportedException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.StandardCopyOption;
+import java.nio.file.attribute.PosixFileAttributeView;
+import java.nio.file.attribute.PosixFilePermission;
+import java.nio.file.attribute.PosixFilePermissions;
import java.util.HashMap;
import java.util.Map;
import java.util.Objects;
+import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
@@ -182,12 +190,9 @@ public class FileStateRepository extends ServiceSupport
implements StateReposito
if (!fileStore.exists()) {
FileUtil.createNewFile(fileStore);
}
- // append to store
+ // append to store (as a single write, so the line is not split in
several writes)
fos = new FileOutputStream(fileStore, true);
- fos.write(key.getBytes());
- fos.write(KEY_VALUE_DELIMITER.getBytes());
- fos.write(value.getBytes());
- fos.write(STORE_DELIMITER.getBytes());
+ fos.write((key + KEY_VALUE_DELIMITER + value +
STORE_DELIMITER).getBytes());
} catch (IOException e) {
throw RuntimeCamelException.wrapRuntimeCamelException(e);
} finally {
@@ -200,19 +205,48 @@ public class FileStateRepository extends ServiceSupport
implements StateReposito
*/
protected void trunkStore() {
LOG.info("Trunking state filestore: {}", fileStore);
+ // write the 1st level cache to a temporary file and then replace the
file store with it, so the file store
+ // is never left truncated or half written (such as if writing fails,
or the JVM crashes while writing)
+ Path target = fileStore.toPath();
+ File tmp = null;
+ boolean written = false;
FileOutputStream fos = null;
try {
- fos = new FileOutputStream(fileStore);
+ if (Files.exists(target)) {
+ // replace the real file of a symlinked store (so the link is
kept, and the temporary file is on the
+ // same file system as the store)
+ target = target.toRealPath();
+ }
+ tmp = new File(target + ".tmp");
+ if (Files.exists(target) && Files.getFileAttributeView(target,
PosixFileAttributeView.class) != null) {
+ // give the temporary file the permissions of the store before
the state is written to it, so the
+ // state is never readable with wider permissions than the
store
+ Set<PosixFilePermission> permissions =
Files.getPosixFilePermissions(target);
+ Files.deleteIfExists(tmp.toPath());
+ Files.createFile(tmp.toPath(),
PosixFilePermissions.asFileAttribute(permissions));
+ // the umask may have removed some of them
+ Files.setPosixFilePermissions(tmp.toPath(), permissions);
+ }
+ fos = new FileOutputStream(tmp);
for (Map.Entry<String, String> entry : cache.entrySet()) {
- fos.write(entry.getKey().getBytes());
- fos.write(KEY_VALUE_DELIMITER.getBytes());
- fos.write(entry.getValue().getBytes());
- fos.write(STORE_DELIMITER.getBytes());
+ fos.write((entry.getKey() + KEY_VALUE_DELIMITER +
entry.getValue() + STORE_DELIMITER).getBytes());
+ }
+ fos.getFD().sync();
+ fos.close();
+ fos = null;
+ try {
+ Files.move(tmp.toPath(), target,
StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE);
+ } catch (AtomicMoveNotSupportedException e) {
+ Files.move(tmp.toPath(), target,
StandardCopyOption.REPLACE_EXISTING);
}
+ written = true;
} catch (IOException e) {
throw RuntimeCamelException.wrapRuntimeCamelException(e);
} finally {
IOHelper.close(fos, "Trunking file state repository", LOG);
+ if (!written && tmp != null) {
+ FileUtil.deleteFile(tmp);
+ }
}
}
@@ -243,6 +277,11 @@ public class FileStateRepository extends ServiceSupport
implements StateReposito
while (scanner.hasNext()) {
String line = scanner.next();
int separatorIndex = line.indexOf(KEY_VALUE_DELIMITER);
+ if (separatorIndex < 0) {
+ // an incomplete line (such as the last line, if the JVM
crashed while appending to the store)
+ LOG.warn("Skipping invalid line in state filestore: {}",
fileStore);
+ continue;
+ }
String key = line.substring(0, separatorIndex);
String value = line.substring(separatorIndex +
KEY_VALUE_DELIMITER.length());
cache.put(key, value);
@@ -266,10 +305,16 @@ public class FileStateRepository extends ServiceSupport
implements StateReposito
@Override
protected void doStop() throws Exception {
- // reset will trunk and clear the cache
- trunkStore();
- cache.clear();
- init.set(false);
+ // must hold the lock, as the store may still be updated while stopping
+ cacheAndStoreLock.lock();
+ try {
+ // reset will trunk and clear the cache
+ trunkStore();
+ cache.clear();
+ init.set(false);
+ } finally {
+ cacheAndStoreLock.unlock();
+ }
}
public File getFileStore() {