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() {

Reply via email to