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

krishvishal pushed a commit to branch sim-workload-faults
in repository https://gitbox.apache.org/repos/asf/iggy.git

commit 007e6537d6e66bac332751da9ce68be559cff438
Author: Krishna Vishal <[email protected]>
AuthorDate: Sat Aug 15 11:12:23 2026 +0530

    feat(simulator): checkpoint for real and serve a chunked state transfer
    
    Three things each independently prevented a checkpoint, and all three had
    to go before one could happen. Metadata was built without a data
    directory, so there was no `SnapshotCoordinator` and `checkpoint_if_needed`
    returned immediately. The simulated journal reported no capacity limit, so
    `should_checkpoint` read `None` and never fired even where a coordinator
    existed. And `drain` was inherited from the trait default, which drains
    nothing, so a checkpoint that did run persisted a snapshot and then left
    the entire WAL in place.
    
    That last one is why arming the coordinator alone changed nothing
    observable: with every op still retained, a rejoining peer found
    everything it asked for, ordinary journal repair always sufficed, and no
    `RangeEvicted` was ever sent. The symptom looked like the transfer being
    unreachable when the cause was a WAL that never shrank.
    
    With all three fixed a lagging replica now completes a real transfer:
    RangeEvicted, RequestStateTransfer, StateTransferTarget, RequestStateChunk
    and StateChunk each cross the wire exactly once, followed by the tail
    repair that finishes the rejoin. Stamping a watermark by hand reached the
    handshake and then re-armed 335 times without converging; this converges.
    Those are the last two of the five frames no run had ever delivered.
    
    Checkpoints are opt-in through a separate constructor rather than the
    default. The coordinator persists through `std::fs`, and a harness whose
    defining property is that it touches nothing outside memory should not
    start writing files because a scenario forgot to say otherwise. The writes
    are synchronous and never reach the executor, so replay stays
    deterministic, and the caller owns the directory. The data directory is
    retained across a restart for the same reason the WAL and superblocks
    already are: a rebuilt replica has to read back the snapshot its previous
    incarnation wrote.
---
 core/simulator/Cargo.toml     |   3 +
 core/simulator/src/deps.rs    | 100 ++++++++++++++++++++++++
 core/simulator/src/lib.rs     | 173 ++++++++++++++++++++++++++++++++++++++++++
 core/simulator/src/replica.rs |   8 +-
 4 files changed, 283 insertions(+), 1 deletion(-)

diff --git a/core/simulator/Cargo.toml b/core/simulator/Cargo.toml
index adc3582a0..8c82f68da 100644
--- a/core/simulator/Cargo.toml
+++ b/core/simulator/Cargo.toml
@@ -50,6 +50,9 @@ strum = { workspace = true }
 tracing = { workspace = true }
 tracing-subscriber = { workspace = true }
 
+[dev-dependencies]
+tempfile = { workspace = true }
+
 [lints.clippy]
 enum_glob_use = "deny"
 pedantic = "deny"
diff --git a/core/simulator/src/deps.rs b/core/simulator/src/deps.rs
index 95a16c019..499bbce02 100644
--- a/core/simulator/src/deps.rs
+++ b/core/simulator/src/deps.rs
@@ -27,6 +27,7 @@ use metadata::stm::user::Users;
 use server_common::{Message, iobuf::Owned};
 use std::cell::{Cell, RefCell, UnsafeCell};
 use std::collections::HashMap;
+use std::ops::RangeInclusive;
 
 /// Fixed synthetic epoch for [`SimClock`]: 2026-01-01T00:00:00Z in micros.
 ///
@@ -107,6 +108,15 @@ pub struct SimJournal<S: Storage> {
     /// Snapshot watermark; see the `Journal::snapshot_op` impl for why a real
     /// value here is what makes `RangeEvicted` reachable at all.
     snapshot_op: Cell<u64>,
+    /// Slots this journal pretends to have, or `None` for unbounded.
+    ///
+    /// A real WAL is a fixed ring, and running low on slots is what forces a
+    /// checkpoint (`SnapshotCoordinator::should_checkpoint` gates on
+    /// `remaining_capacity`). An unbounded journal answers `None` there and so
+    /// never triggers one, which is why the simulator had never checkpointed 
even
+    /// where a coordinator existed. Defaults to unbounded so existing 
scenarios
+    /// keep their behaviour; a test that wants checkpoints sets a small count.
+    slot_count: Cell<Option<usize>>,
     /// Debug-only single-accessor tripwire. `entry` / `append` hold a
     /// [`JournalAccessGuard`] across their whole body, including the storage
     /// `.await`, so if a suspending storage tier ever let a second task touch
@@ -125,6 +135,7 @@ impl<S: Storage + Default> Default for SimJournal<S> {
             write_offset: Cell::new(0),
             last_op: Cell::new(None),
             snapshot_op: Cell::new(0),
+            slot_count: Cell::new(None),
             #[cfg(debug_assertions)]
             accessing: Cell::new(false),
         }
@@ -184,6 +195,26 @@ impl<S: Storage<Buffer = Vec<u8>>> Journal<S> for 
SimJournal<S> {
         self.last_op.get()
     }
 
+    /// Slots left before a checkpoint is forced, mirroring
+    /// `PrepareJournal::remaining_capacity`: the ring holds `slot_count` 
entries
+    /// and everything at or below the snapshot watermark is reclaimable, so 
what
+    /// is occupied is `last_op - snapshot_op`.
+    ///
+    /// `None` while unbounded, which is what the trait default gave before and
+    /// what `should_checkpoint` reads as "never checkpoint".
+    fn remaining_capacity(&self) -> Option<usize> {
+        let slot_count = self.slot_count.get()?;
+        let Some(last) = self.last_op.get() else {
+            return Some(slot_count);
+        };
+        let snapshot = self.snapshot_op.get();
+        if last <= snapshot {
+            return Some(slot_count);
+        }
+        let used = usize::try_from(last - snapshot).unwrap_or(usize::MAX);
+        Some(slot_count.saturating_sub(used))
+    }
+
     /// Drop the suffix, so a simulated backup whose entries disagree with a 
started
     /// view reconciles the way a real one does. Mirrors
     /// `PrepareJournal::truncate_from`, whose watermark stays put; here it 
never moves.
@@ -296,6 +327,66 @@ impl<S: Storage<Buffer = Vec<u8>>> Journal<S> for 
SimJournal<S> {
         let headers = unsafe { &*self.headers.get() };
         headers.get(&(idx as u64))
     }
+
+    /// Reclaim the prefix a checkpoint superseded, advancing the snapshot
+    /// watermark to the end of the drained range.
+    ///
+    /// Required, not inherited: the trait's default drains nothing, so a
+    /// simulated checkpoint persisted a snapshot and then left the whole WAL 
in
+    /// place. A peer's repair then found every op it asked for and journal 
repair
+    /// always sufficed, which is why arming the coordinator alone still never
+    /// produced a `RangeEvicted` or a state transfer.
+    ///
+    /// The watermark moves last, mirroring `PrepareJournal::drain`: advancing 
it
+    /// before the entries are gone would make live entries look evictable.
+    async fn drain(&self, ops: RangeInclusive<u64>) -> 
std::io::Result<Vec<Self::Entry>> {
+        #[cfg(debug_assertions)]
+        let _guard = JournalAccessGuard::new(&self.accessing);
+        let end_op = *ops.end();
+        let doomed: Vec<u64> = {
+            let headers = unsafe { &*self.headers.get() };
+            let mut doomed: Vec<u64> = headers
+                .keys()
+                .copied()
+                .filter(|op| ops.contains(op))
+                .collect();
+            // Sorted: the trait promises the removed entries in op order, and
+            // hash order would also make a replay of this drain diverge.
+            doomed.sort_unstable();
+            doomed
+        };
+
+        let mut drained = Vec::with_capacity(doomed.len());
+        for op in doomed {
+            // Read before removing, through `Storage` rather than the
+            // `MemStorage`-only sync path, so this stays generic like the 
rest of
+            // the impl. The borrow spans the read for the same reason 
`entry`'s
+            // does, and is sound on the same grounds: `MemStorage` never
+            // suspends, and the guard above trips if that ever changes.
+            let located = {
+                let headers = unsafe { &*self.headers.get() };
+                let offsets = unsafe { &*self.offsets.get() };
+                headers
+                    .get(&op)
+                    .and_then(|header| offsets.get(&op).map(|offset| 
(header.size, *offset)))
+            };
+            if let Some((size, offset)) = located
+                && let Ok(buffer) = self.storage.read_at(offset, vec![0; size 
as usize]).await
+                && let Ok(message) = 
Message::try_from(Owned::<4096>::copy_from_slice(&buffer))
+            {
+                drained.push(message);
+            }
+            let headers = unsafe { &mut *self.headers.get() };
+            let offsets = unsafe { &mut *self.offsets.get() };
+            headers.remove(&op);
+            offsets.remove(&op);
+        }
+
+        if end_op > self.snapshot_op.get() {
+            self.snapshot_op.set(end_op);
+        }
+        Ok(drained)
+    }
 }
 
 impl JournalHandle for SimJournal<MemStorage> {
@@ -315,6 +406,15 @@ impl SimJournal<MemStorage> {
         self.last_op.get()
     }
 
+    /// Bound this journal to `slots`, so running low on them forces a 
checkpoint.
+    ///
+    /// Unbounded by default (see the `slot_count` field): a journal that never
+    /// runs out never triggers `should_checkpoint`, so no checkpoint ever
+    /// happens and nothing produces the snapshot a state transfer serves.
+    pub fn set_slot_count(&self, slots: usize) {
+        self.slot_count.set(Some(slots));
+    }
+
     /// Forget one op, leaving a hole exactly where a lost prepare would.
     ///
     /// Tests only. The alternative is choreographing `Prepare`, `Commit` and
diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs
index 84829852b..27135d4b0 100644
--- a/core/simulator/src/lib.rs
+++ b/core/simulator/src/lib.rs
@@ -118,6 +118,10 @@ pub struct SimReplica {
     /// regression. Populated on the way into a restart (see
     /// [`Simulator::replica_restart`]) and consumed by 
`materialise_partition`.
     partition_logs: RefCell<HashMap<IggyNamespace, (RetainedPartitionLog, u64, 
u64)>>,
+    /// This replica's data directory when checkpoints are enabled, `None`
+    /// otherwise. Retained for the same reason as the WAL and the 
superblocks: a
+    /// restart must read back the snapshot its previous incarnation wrote.
+    data_dir: Option<std::path::PathBuf>,
     /// Keeps each pump's stop channel alive; dropping one would end that
     /// pump gracefully, which is reserved for future shutdown/restart
     /// tests (crash uses `DetExecutor::abort` instead).
@@ -252,6 +256,51 @@ impl Simulator {
         )
     }
 
+    /// [`Simulator::new`] with checkpoints enabled, each replica rooted at
+    /// `<data_dir_root>/replica-N`.
+    ///
+    /// A data directory is what arms the metadata `SnapshotCoordinator`; 
without
+    /// one `checkpoint_if_needed` returns immediately and no replica ever
+    /// checkpoints, so nothing produces the snapshot a state transfer serves.
+    ///
+    /// Opt-in and separate from the other constructors because the coordinator
+    /// persists through `std::fs`, and a harness whose defining property is 
that
+    /// it touches nothing outside memory should not start writing files 
because a
+    /// scenario forgot to say otherwise. The writes are synchronous and never
+    /// touch the executor, so replay stays deterministic; the caller owns the
+    /// directory's lifetime (a `tempfile::TempDir` in tests).
+    ///
+    /// Pair with [`Simulator::set_metadata_journal_slots`]: a checkpoint is 
forced
+    /// by the journal running low on slots, and the simulated journal is 
unbounded
+    /// until told otherwise.
+    ///
+    /// # Panics
+    /// Panics on duplicate `client_id`s, `shards_per_replica == 0`, or if the
+    /// per-replica directories cannot be created.
+    pub fn with_checkpoints(
+        replica_count: usize,
+        clients: impl Iterator<Item = u128>,
+        network_options: PacketSimulatorOptions,
+        data_dir_root: std::path::PathBuf,
+    ) -> Self {
+        Self::build_inner(
+            replica_count,
+            1,
+            clients,
+            network_options,
+            false,
+            Some(data_dir_root),
+        )
+    }
+
+    /// Bound every replica's metadata journal to `slots`, so filling it 
forces a
+    /// checkpoint. See [`deps::SimJournal::set_slot_count`].
+    pub fn set_metadata_journal_slots(&self, slots: usize) {
+        for replica in &self.replicas {
+            replica.metadata_journal.set_slot_count(slots);
+        }
+    }
+
     #[allow(clippy::cast_possible_truncation, clippy::too_many_lines)]
     fn build(
         replica_count: usize,
@@ -259,6 +308,25 @@ impl Simulator {
         clients: impl Iterator<Item = u128>,
         network_options: PacketSimulatorOptions,
         shell: bool,
+    ) -> Self {
+        Self::build_inner(
+            replica_count,
+            shards_per_replica,
+            clients,
+            network_options,
+            shell,
+            None,
+        )
+    }
+
+    #[allow(clippy::cast_possible_truncation, clippy::too_many_lines)]
+    fn build_inner(
+        replica_count: usize,
+        shards_per_replica: u16,
+        clients: impl Iterator<Item = u128>,
+        network_options: PacketSimulatorOptions,
+        shell: bool,
+        data_dir_root: Option<std::path::PathBuf>,
     ) -> Self {
         assert!(
             shards_per_replica >= 1,
@@ -312,6 +380,17 @@ impl Simulator {
             // restart increments in `replica_restart` never collide across 
replicas.
             // Non-zero.
             let metadata_incarnation = 1 + (u128::from(id) << 64);
+            // One directory per replica when checkpoints are enabled, created 
up
+            // front because the snapshot coordinator writes into
+            // `<dir>/metadata/` and does not create it. Retained on 
`SimReplica`
+            // so a restart reads back the snapshot its previous incarnation 
wrote,
+            // exactly as the retained WAL and superblock already do.
+            let replica_data_dir = data_dir_root.as_ref().map(|root| {
+                let dir = root.join(format!("replica-{i}"));
+                
std::fs::create_dir_all(dir.join(metadata::impls::METADATA_DIR))
+                    .expect("simulator data directory is creatable");
+                dir
+            });
 
             // One crossfire mesh per replica; every shard gets a clone of
             // the canonical senders vec and exclusively takes its inbox.
@@ -354,6 +433,7 @@ impl Simulator {
                     shard_journal,
                     None, // fresh boot: no recovered VSR state
                     metadata_incarnation,
+                    (shard_idx == 0).then(|| 
replica_data_dir.clone()).flatten(),
                 );
                 if shard_idx == 0 {
                     metadata_bundle = Some(
@@ -381,6 +461,7 @@ impl Simulator {
                 metadata_incarnation,
                 partition_superblocks: RefCell::new(HashMap::new()),
                 partition_logs: RefCell::new(HashMap::new()),
+                data_dir: replica_data_dir,
                 _stop_txs: stop_txs,
                 pump_tasks,
             });
@@ -880,6 +961,10 @@ impl Simulator {
             .read_latest_sync()
             .and_then(|bytes| VsrState::try_from(bytes.as_slice()).ok());
 
+        // Carried across the rebuild like the WAL and superblocks: the rebuilt
+        // replica has to find the snapshot its previous incarnation persisted.
+        let replica_data_dir = self.replicas[idx].data_dir.clone();
+
         let consensus_clock = 
ConsensusClock::new(Rc::new(SimClock::new(self.executor.timer())));
         let outbox = Rc::clone(&self.outboxes[idx]);
         let (senders, mut inboxes) =
@@ -914,6 +999,7 @@ impl Simulator {
                 shard_journal,
                 recovered_state,
                 metadata_incarnation,
+                (shard_idx == 0).then(|| replica_data_dir.clone()).flatten(),
             );
             if shard_idx == 0 {
                 metadata_bundle =
@@ -937,6 +1023,7 @@ impl Simulator {
             metadata_incarnation,
             partition_superblocks: RefCell::new(partition_superblocks),
             partition_logs: RefCell::new(partition_logs),
+            data_dir: replica_data_dir,
             _stop_txs: stop_txs,
             pump_tasks,
         };
@@ -3314,6 +3401,92 @@ mod tests {
     /// pipeline to `PIPELINE_PREPARE_QUEUE_MAX`; a request on `ns_b`
     /// still commits while `ns_a` is wedged (no quorum without backup
     /// acks); lifting the block drains `ns_a` completely.
+    /// A replica that checkpoints serves a real state transfer: the rejoining
+    /// peer fetches the snapshot in chunks rather than stalling at the 
handshake.
+    ///
+    /// The companion to
+    /// [`repair_below_the_snapshot_floor_escalates_to_state_transfer`], which
+    /// stamps a watermark without producing snapshot bytes and so reaches only
+    /// `StateTransferTarget`. Here the coordinator is armed with a data 
directory
+    /// and the journal is bounded, so the cluster checkpoints for real and the
+    /// transfer has artifacts to serve.
+    ///
+    /// Two things had to be true for a checkpoint to happen at all, and 
neither
+    /// was: metadata was built without a data directory, so there was no
+    /// `SnapshotCoordinator`; and the simulated journal reported no capacity
+    /// limit, so `should_checkpoint` never fired even where a coordinator 
existed.
+    #[test]
+    fn checkpointing_cluster_serves_a_chunked_state_transfer() {
+        use iggy_binary_protocol::Command2;
+
+        
server_common::MemoryPool::init_pool(&server_common::MemoryPoolConfigOther {
+            enabled: false,
+            size: iggy_common::IggyByteSize::from(0u64),
+            bucket_capacity: 1,
+        });
+
+        let root = tempfile::tempdir().expect("temp dir for the simulator's 
snapshots");
+        let replica_count: u8 = 3;
+        let client_id: u128 = 1;
+        let network_opts = packet::PacketSimulatorOptions {
+            node_count: replica_count,
+            client_count: 1,
+            seed: 0xC4E0_0001,
+            ..packet::PacketSimulatorOptions::default()
+        };
+        let mut sim = Simulator::with_checkpoints(
+            usize::from(replica_count),
+            std::iter::once(client_id),
+            network_opts,
+            root.path().to_path_buf(),
+        );
+        // Small enough that the ops below cross the margin; the coordinator 
forces
+        // a checkpoint once free slots fall to its margin (64 by default).
+        sim.set_metadata_journal_slots(80);
+
+        let client = SimClient::new(client_id);
+        sim.register_client_with_primary(&client);
+
+        let lagging = 2u8;
+        sim.replica_crash(lagging);
+
+        // Commit past the checkpoint margin while the lagging replica is 
down, so
+        // the survivors checkpoint and compact the prefix it is missing.
+        for sequence in 0..40u32 {
+            let msg = 
client.create_stream(&format!("wl-checkpoint-{sequence}"));
+            sim.submit_request(client_id, 0, msg.into_generic());
+            for _ in 0..40 {
+                sim.step();
+            }
+        }
+
+        let snapshot = root
+            .path()
+            .join("replica-0")
+            .join("metadata")
+            .join("snapshot.bin");
+        assert!(
+            snapshot.exists(),
+            "the primary never checkpointed, so there is no snapshot to 
transfer: \
+             raise the op count or lower the journal slot count"
+        );
+
+        sim.replica_restart(lagging);
+        for _ in 0..8_000 {
+            sim.step();
+        }
+
+        assert!(
+            sim.network.delivered_any(Command2::RequestStateChunk),
+            "the rejoining replica never asked for a chunk, so the transfer \
+             stalled at the handshake exactly as it does without a checkpoint"
+        );
+        assert!(
+            sim.network.delivered_any(Command2::StateChunk),
+            "no chunk was served: the peer offered a transfer it could not 
fulfil"
+        );
+    }
+
     /// A backup whose gap sits below the serving peer's snapshot floor is told
     /// `RangeEvicted` and converts its repair into a state transfer.
     ///
diff --git a/core/simulator/src/replica.rs b/core/simulator/src/replica.rs
index 2764b96da..473083d59 100644
--- a/core/simulator/src/replica.rs
+++ b/core/simulator/src/replica.rs
@@ -125,6 +125,7 @@ pub fn new_shard(
     metadata_journal: Option<Rc<SimJournal<MemStorage>>>,
     recovered_state: Option<VsrState>,
     incarnation: u128,
+    data_dir: Option<std::path::PathBuf>,
 ) -> (Rc<Replica>, Option<SimMetadataBundle>) {
     // Metadata is single-writer, mirroring the server bootstrap. Shard 0 owns
     // the only writable STM; every peer shard rebuilds a reader-mode mirror 
from
@@ -243,13 +244,18 @@ pub fn new_shard(
     });
     let metadata_snapshot = (shard_idx == 0).then(SimSnapshot::default);
 
+    // A data directory arms the `SnapshotCoordinator`, without which
+    // `checkpoint_if_needed` returns immediately and the replica never
+    // checkpoints. `None` keeps that (the default for scenarios that do not 
care);
+    // `Some` is opt-in, because the coordinator persists through `std::fs` 
and a
+    // harness that writes files is not what most specs want.
     let metadata = IggyMetadata::new(
         metadata_consensus,
         metadata_journal,
         metadata_snapshot,
         superblock,
         mux,
-        None,
+        data_dir,
     );
 
     // Reconstruct shard 0's committed metadata from the retained WAL, the sim 
analog

Reply via email to