This is an automated email from the ASF dual-hosted git repository.

spetz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git


The following commit(s) were added to refs/heads/master by this push:
     new 328b28974 fix(journal): reject recovery slot collisions (#4009)
328b28974 is described below

commit 328b289745e0e24bace1bb531fde4c0f9159d608
Author: aias00 <[email protected]>
AuthorDate: Tue Sep 1 14:56:53 2026 +0800

    fix(journal): reject recovery slot collisions (#4009)
---
 core/configs/src/server_config/metadata.rs |  4 +-
 core/journal/src/prepare_journal.rs        | 85 +++++++++++++++++++++++++++++-
 core/server/config.toml                    |  4 +-
 3 files changed, 89 insertions(+), 4 deletions(-)

diff --git a/core/configs/src/server_config/metadata.rs 
b/core/configs/src/server_config/metadata.rs
index 2ba8d821f..fbec9853f 100644
--- a/core/configs/src/server_config/metadata.rs
+++ b/core/configs/src/server_config/metadata.rs
@@ -100,7 +100,9 @@ pub struct MetadataConfig {
     /// Size of the metadata WAL's in-memory index, in slots (one
     /// committed-but-unsnapshotted op per slot). Headroom between forced
     /// checkpoints; more slots = rarer checkpoints, more memory, larger
-    /// per-checkpoint WAL rewrites.
+    /// per-checkpoint WAL rewrites. Reducing this for an existing data 
directory
+    /// is refused at boot when the live WAL suffix cannot fit without 
collisions;
+    /// restore the previous value to recover.
     pub journal_slots: usize,
 
     /// Slot count of the VSR client table: how many distinct clients
diff --git a/core/journal/src/prepare_journal.rs 
b/core/journal/src/prepare_journal.rs
index 278fb8e0e..b00f24b56 100644
--- a/core/journal/src/prepare_journal.rs
+++ b/core/journal/src/prepare_journal.rs
@@ -575,8 +575,25 @@ impl PrepareJournal {
 
             let slot = slot_for_op(header.op, slot_count);
 
-            // Note: Regarding duplicate op in WAL. We rewrite it with 
whichever
-            // is the latest entry.
+            // Match append's collision fence while rebuilding the index. 
Reopening
+            // with fewer configured slots can otherwise hide an unsnapshotted 
op
+            // even though its bytes remain in the WAL. A duplicate op 
deliberately
+            // keeps the latest entry, and a snapshotted op is safe to evict.
+            if let Some(existing) = headers[slot]
+                && existing.op != header.op
+                && existing.op > snapshot_op
+            {
+                return Err(JournalError::Io(io::Error::new(
+                    io::ErrorKind::InvalidData,
+                    format!(
+                        "journal slot collision while rebuilding the index: op 
{} and \
+                         unsnapshotted op {} map to slot {slot} with 
slot_count={slot_count} \
+                         (snapshot_op={snapshot_op}); restore the previous \
+                         metadata.journal_slots value or checkpoint before 
shrinking it",
+                        header.op, existing.op,
+                    ),
+                )));
+            }
             headers[slot] = Some(header);
             offsets[slot] = Some(pos);
 
@@ -2099,6 +2116,70 @@ mod tests {
         assert_eq!(journal.last_op(), Some(3));
     }
 
+    #[compio::test]
+    async fn 
reopen_with_fewer_slots_rejects_unsnapshotted_collision_without_changing_wal() {
+        let dir = tempdir().unwrap();
+        let path = dir.path().join("journal.wal");
+        let journal = PrepareJournal::open_with_slots(&path, 0, 
4).await.unwrap();
+        journal.append(make_prepare(1, 32)).await.unwrap();
+        journal.append(make_prepare(3, 32)).await.unwrap();
+        drop(journal);
+        let wal_before = std::fs::read(&path).unwrap();
+
+        let error = PrepareJournal::open_with_slots(&path, 0, 2)
+            .await
+            .expect_err("shrinking the index must not hide an unsnapshotted 
entry");
+
+        assert!(
+            error.to_string().contains("journal slot collision"),
+            "unexpected error: {error}"
+        );
+        assert_eq!(
+            std::fs::read(&path).unwrap(),
+            wal_before,
+            "a refused scan must leave the WAL intact"
+        );
+    }
+
+    #[compio::test]
+    async fn reopen_keeps_latest_duplicate_op() {
+        let dir = tempdir().unwrap();
+        let path = dir.path().join("journal.wal");
+        let journal = PrepareJournal::open_with_slots(&path, 0, 
4).await.unwrap();
+        journal.append(make_prepare(1, 32)).await.unwrap();
+        journal.set_snapshot_op(1);
+        journal.append(make_prepare(1, 64)).await.unwrap();
+        drop(journal);
+
+        let journal = PrepareJournal::open_with_slots(&path, 0, 
2).await.unwrap();
+        let header = *journal.header(1).expect("duplicate op must remain 
indexed");
+        assert_eq!(header.size as usize, HEADER_SIZE + 64);
+        assert_eq!(
+            journal
+                .entry_at(&header)
+                .await
+                .unwrap()
+                .unwrap()
+                .as_slice()
+                .len(),
+            HEADER_SIZE + 64
+        );
+    }
+
+    #[compio::test]
+    async fn reopen_with_fewer_slots_can_evict_snapshotted_entry() {
+        let dir = tempdir().unwrap();
+        let path = dir.path().join("journal.wal");
+        let journal = PrepareJournal::open_with_slots(&path, 0, 
4).await.unwrap();
+        journal.append(make_prepare(1, 32)).await.unwrap();
+        journal.append(make_prepare(3, 32)).await.unwrap();
+        drop(journal);
+
+        let journal = PrepareJournal::open_with_slots(&path, 1, 
2).await.unwrap();
+        assert!(journal.header(1).is_none());
+        assert_eq!(journal.header(3).map(|header| header.op), Some(3));
+    }
+
     const POISON_REASON: &str = "test: simulated post-rename failure";
 
     #[compio::test]
diff --git a/core/server/config.toml b/core/server/config.toml
index e70cfacee..79bdc46e8 100644
--- a/core/server/config.toml
+++ b/core/server/config.toml
@@ -960,7 +960,9 @@ prepare_queue_depth = 32
 # Size of the metadata WAL's in-memory index, in slots (one committed but
 # not-yet-snapshotted op per slot). Larger values buy more headroom
 # between forced checkpoints at the cost of memory and bigger WAL
-# rewrites per checkpoint.
+# rewrites per checkpoint. Reducing this for an existing data directory is
+# refused at boot if the live WAL suffix collides in the smaller index;
+# restore the previous value to recover.
 journal_slots = 1024
 
 # Slot count of the VSR client table: how many distinct clients (TCP/QUIC/WS

Reply via email to