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
