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 b98d3a5a66f3d88debbfe360776536e3213ada03
Author: Krishna Vishal <[email protected]>
AuthorDate: Fri Aug 14 18:58:34 2026 +0530

    feat(simulator): carry partition logs across a replica restart
    
    A restart drops the whole partition and rebuilds it, which costs a real
    server nothing: its messages are in segment files and boot recovers the
    offset counter from them. The simulator's partitions are in-memory, so a
    rebuilt one came back empty, its commit_offset regressed to zero, and the
    monotonicity invariant reported a discarded log as a consensus
    regression. That confined restart injection to the metadata plane, whose
    WAL the harness already retained.
    
    The harness now holds each group's log across the rebuild and hands it
    back on re-materialisation, the same role it already plays for the
    metadata WAL and the superblocks. `adopt_retained_log` lives in the
    partitions crate beside `restore_offset_frontier` rather than in the
    harness, because the two have to agree on maxing against what is already
    proved, and the one place built to catch violations of that rule is the
    wrong place to keep a second copy of it.
    
    Relaxing the invariant instead would have modelled total data loss rather
    than a restart, so it stays strict and its doc no longer claims the
    partition plane is exempt.
    
    The regression test asserts both the counter and a poll that returns
    messages: a carried-over counter over an empty log would satisfy the
    counter alone while having lost everything behind it.
    
    Partition-plane crash/restart runs now drain and converge across seeds
    where they previously tripped the invariant (commit_offset 79 -> 0). The
    metadata-plane wedge reported earlier is untouched by this and is now
    better characterised: it reproduces on every plane containing metadata
    ops at every seed tried, so it is systematic rather than seed luck.
---
 core/partitions/src/iggy_partition.rs     |  45 ++++++++
 core/partitions/src/lib.rs                |  13 +++
 core/shard/Cargo.toml                     |   6 +-
 core/shard/src/lib.rs                     |  18 ++++
 core/simulator/src/lib.rs                 | 166 +++++++++++++++++++++++++++++-
 core/simulator/src/workload/auditor.rs    |   1 -
 core/simulator/src/workload/invariants.rs |  21 ++--
 7 files changed, 250 insertions(+), 20 deletions(-)

diff --git a/core/partitions/src/iggy_partition.rs 
b/core/partitions/src/iggy_partition.rs
index 8dc976169..16e07230d 100644
--- a/core/partitions/src/iggy_partition.rs
+++ b/core/partitions/src/iggy_partition.rs
@@ -752,6 +752,51 @@ where
     /// simulator share one implementation. A copy in the harness was a copy of
     /// the max rule that had lost the max, in the one place built to catch
     /// violations of it.
+    /// Adopt a log carried over from a previous incarnation of this partition,
+    /// standing in for what segment recovery reads off disk at boot.
+    ///
+    /// A restart drops the whole partition and rebuilds it, which on a real
+    /// server loses nothing: the messages are in segment files and boot 
recovers
+    /// the offset counter from them. The simulator's partitions are 
in-memory, so
+    /// without this the rebuilt partition comes back empty and its
+    /// `commit_offset` regresses to zero -- reported as a consensus regression
+    /// when it is really the harness having thrown the data away.
+    ///
+    /// `durable_offset` and `write_offset` are what the caller recovered, 
exactly
+    /// as `segment_recovery` derives them from the segments it read. Applied 
as a
+    /// MAX against whatever the superblock frontier already proved, never as 
an
+    /// overwrite, for the same reason [`Self::restore_offset_frontier`] 
maxes: a
+    /// recovered value that is behind the frontier must not lower it.
+    ///
+    /// Lives here rather than in the harness so the rule is stated once, 
beside
+    /// the frontier restore it has to agree with.
+    #[cfg(any(test, feature = "simulator"))]
+    pub fn adopt_retained_log(
+        &mut self,
+        log: SegmentedLog<PartitionJournal<PartitionJournalMemStorage>, 
PartitionJournalMemStorage>,
+        durable_offset: u64,
+        write_offset: u64,
+    ) {
+        self.log = log;
+        // Empty carry-over: the previous incarnation never took a write, so 
there
+        // is no offset space to restore and claiming one would make the next
+        // prepare mint from a base no peer agrees on.
+        if write_offset == 0 && durable_offset == 0 && 
!self.should_increment_offset {
+            return;
+        }
+        let durable = durable_offset.max(self.offset.load(Ordering::Acquire));
+        let dirty = write_offset
+            .max(durable)
+            .max(self.dirty_offset.load(Ordering::Relaxed));
+        self.offset.store(durable, Ordering::Release);
+        self.dirty_offset.store(dirty, Ordering::Relaxed);
+        self.should_increment_offset = true;
+        // Everything carried over is already persisted as far as this replica 
is
+        // concerned, so the flush and commit paths must not re-persist or 
re-count
+        // it -- the same contract boot gives a partition recovered from 
segments.
+        self.recovered_durable_offset = Some(durable);
+    }
+
     pub fn restore_offset_frontier(&mut self, recovered: 
Option<&consensus::VsrState>) {
         let Some(frontier) = recovered
             .map(|state| state.offset_frontier)
diff --git a/core/partitions/src/lib.rs b/core/partitions/src/lib.rs
index 2b62402cf..1a773e747 100644
--- a/core/partitions/src/lib.rs
+++ b/core/partitions/src/lib.rs
@@ -38,6 +38,19 @@ pub use iggy_index_reader::IggyIndexReader;
 pub use iggy_index_writer::IggyIndexWriter;
 pub use iggy_partition::{IggyPartition, PurgeError};
 pub use iggy_partitions::IggyPartitions;
+
+/// A partition's message log, named so a caller can carry one across a 
rebuild.
+///
+/// Exists for the simulator, which has no segment files to recover from and so
+/// must hold the log itself for a restarted replica to come back with its data
+/// (see [`IggyPartition::adopt_retained_log`]). The generic parameters are the
+/// only pair `IggyPartition::log` is ever instantiated with, so this names the
+/// concrete type rather than widening anything.
+#[cfg(any(test, feature = "simulator"))]
+pub type RetainedPartitionLog = log::SegmentedLog<
+    journal::PartitionJournal<journal::PartitionJournalMemStorage>,
+    journal::PartitionJournalMemStorage,
+>;
 pub use journal::{EVICTED_RING_BYTES_MAX, EVICTED_RING_CAPACITY};
 pub use messages_writer::MessagesWriter;
 pub use offset_storage::delete_persisted_offset;
diff --git a/core/shard/Cargo.toml b/core/shard/Cargo.toml
index 874fb5353..826f747ce 100644
--- a/core/shard/Cargo.toml
+++ b/core/shard/Cargo.toml
@@ -28,7 +28,11 @@ publish = false
 # off the pump task. A `-p iggy-server` build excludes it; `cargo build
 # --workspace` unifies features and compiles it into the shared `shard`
 # unit (simulator requests it). Benign: no production caller.
-simulator = []
+# Forwards to `partitions/simulator` because `init_partition` now names
+# `partitions::RetainedPartitionLog` and calls `adopt_retained_log`, both gated
+# there. Without the forward this feature only compiles when something else in
+# the build happens to enable the partitions one.
+simulator = ["partitions/simulator"]
 
 [dependencies]
 compio = { workspace = true }
diff --git a/core/shard/src/lib.rs b/core/shard/src/lib.rs
index bae836f8b..607713758 100644
--- a/core/shard/src/lib.rs
+++ b/core/shard/src/lib.rs
@@ -3424,12 +3424,21 @@ where
     /// `recovered_state` is that store's last record, read by the caller (the
     /// store's read is async and this is not), mirroring how `new_shard` 
takes the
     /// metadata plane's.
+    ///
+    /// `retained` is the log a previous incarnation of this group left behind,
+    /// standing in for the segments a real boot recovers from. `None` builds 
an
+    /// empty partition, which is right for a first materialisation and wrong 
for a
+    /// restart: a rebuilt partition with no data reports `commit_offset` 0 and
+    /// looks like a regression rather than a harness that discarded the log. 
The
+    /// recovered offsets ride along because the caller plays the storage layer
+    /// here, exactly as it does for the metadata WAL.
     #[cfg(any(test, feature = "simulator"))]
     pub fn init_partition(
         &self,
         namespace: IggyNamespace,
         superblock: Option<Rc<SB>>,
         recovered_state: Option<consensus::VsrState>,
+        retained: Option<(partitions::RetainedPartitionLog, u64, u64)>,
     ) where
         B: MessageBus + Clone,
     {
@@ -3466,6 +3475,15 @@ where
         if let Some(superblock) = superblock {
             partition.set_superblock(superblock, recovered_state.as_ref());
         }
+        // Retained log before the frontier restore, so the restore maxes 
against
+        // the offsets the log actually proved rather than the zeroes of an 
empty
+        // one. Both are max rules, so the order only decides which value each 
sees
+        // first, never the outcome; doing it in this order keeps the frontier
+        // restore's own precondition (`should_increment_offset` already set 
by a
+        // recovered offset space) meaningful.
+        if let Some((log, durable_offset, write_offset)) = retained {
+            partition.adopt_retained_log(log, durable_offset, write_offset);
+        }
         // The SAME call the boot paths make, not a copy of it: this restore is
         // a max against what the segments already proved, and a harness 
running
         // a divergent copy of that rule cannot catch a violation of it. 
Without
diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs
index 7c38b1668..7d0ce2413 100644
--- a/core/simulator/src/lib.rs
+++ b/core/simulator/src/lib.rs
@@ -38,7 +38,9 @@ use message_bus::installer::conn_info::{ClientConnMeta, 
ClientTransportKind};
 use metadata::impls::metadata::StreamsFrontend;
 use network::Network;
 use packet::{PacketSimulatorOptions, ProcessId};
-use partitions::{Partition, PartitionOffsets, PollFragments, PollingArgs, 
PollingConsumer};
+use partitions::{
+    Partition, PartitionOffsets, PollFragments, PollingArgs, PollingConsumer, 
RetainedPartitionLog,
+};
 use rand::RngExt;
 use rand_xoshiro::Xoshiro256Plus;
 use rand_xoshiro::rand_core::SeedableRng;
@@ -106,6 +108,16 @@ pub struct SimReplica {
     /// leaves the gate, its write-failure fence, and view recovery all
     /// unexercised.
     pub partition_superblocks: RefCell<HashMap<IggyNamespace, 
Rc<SimSuperblock>>>,
+    /// One retained message log per partition group, with the offsets 
recovered
+    /// from it. Harness-owned for the same reason as `metadata_journal`: a
+    /// restart drops and rebuilds the partition, and a real server loses 
nothing
+    /// there because its messages are in segment files that boot recovers the
+    /// offset counter from. The simulator has no segment files, so without 
this
+    /// the rebuilt partition comes back empty, `commit_offset` regresses to 
zero,
+    /// and the monotonicity invariant reports a discarded log as a consensus
+    /// 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)>>,
     /// 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).
@@ -347,6 +359,7 @@ impl Simulator {
                 metadata_journal,
                 metadata_incarnation,
                 partition_superblocks: RefCell::new(HashMap::new()),
+                partition_logs: RefCell::new(HashMap::new()),
                 _stop_txs: stop_txs,
                 pump_tasks,
             });
@@ -771,6 +784,12 @@ impl Simulator {
         // as a rebooted server partition reads the record in its directory.
         let partition_superblocks =
             std::mem::take(&mut 
*self.replicas[idx].partition_superblocks.borrow_mut());
+        // Take each live partition's log while its shard is still standing. 
This
+        // is the harness playing the storage layer: a real server's messages 
are
+        // in segment files and its boot recovers the offset counter from 
them, so
+        // a partition rebuilt with nothing would model total data loss rather 
than
+        // a restart. Read before the rebuild because the rebuild drops the 
shards.
+        let partition_logs = self.retain_partition_logs(idx, 
&partition_superblocks);
 
         // Recover the durable VSR state from the retained superblock before 
the
         // rebuild, as production reads it in restore_metadata_consensus.
@@ -834,6 +853,7 @@ impl Simulator {
             metadata_journal,
             metadata_incarnation,
             partition_superblocks: RefCell::new(partition_superblocks),
+            partition_logs: RefCell::new(partition_logs),
             _stop_txs: stop_txs,
             pump_tasks,
         };
@@ -863,6 +883,42 @@ impl Simulator {
         self.crashed.remove(&replica_index);
     }
 
+    /// Take the message log out of every partition this replica has
+    /// materialised, together with the offsets recovered from it.
+    ///
+    /// Called while the outgoing shards are still alive, so this is the last
+    /// point the data can be read. `std::mem::take` leaves the doomed 
partition
+    /// with an empty log, which nothing observes: the shard it belongs to is
+    /// dropped moments later.
+    ///
+    /// Keyed off `partition_superblocks` rather than the live partitions map
+    /// because that is already the record of which groups this replica has
+    /// materialised, and it is what the restart re-materialises from.
+    fn retain_partition_logs(
+        &self,
+        replica_idx: usize,
+        materialised: &HashMap<IggyNamespace, Rc<SimSuperblock>>,
+    ) -> HashMap<IggyNamespace, (RetainedPartitionLog, u64, u64)> {
+        let replica = &self.replicas[replica_idx];
+        let mut retained = HashMap::with_capacity(materialised.len());
+        for &namespace in materialised.keys() {
+            let partitions = 
replica.partition_shard(namespace).plane.partitions();
+            let Some(partition) = partitions.get_mut_by_ns(&namespace) else {
+                continue;
+            };
+            let offsets = partition.offsets();
+            retained.insert(
+                namespace,
+                (
+                    std::mem::take(&mut partition.log),
+                    offsets.commit_offset,
+                    offsets.write_offset,
+                ),
+            );
+        }
+        retained
+    }
+
     /// Advance consensus timeouts on every live replica without a full
     /// step cycle: fires the pumps' virtual tick timers and runs the
     /// executor to quiescence.
@@ -988,7 +1044,17 @@ fn materialise_partition(replica: &SimReplica, namespace: 
IggyNamespace) {
     let recovered_state = superblock
         .read_latest_sync()
         .and_then(|bytes| VsrState::try_from(bytes.as_slice()).ok());
-    replica.shards[usize::from(owner)].init_partition(namespace, 
Some(superblock), recovered_state);
+    // Hand back the log this group left behind, if it has been materialised
+    // before on this replica. Removed rather than cloned: the rebuilt 
partition
+    // becomes its sole owner, and a second materialisation without a restart 
in
+    // between would otherwise resurrect a log the live partition has moved 
past.
+    let retained = replica.partition_logs.borrow_mut().remove(&namespace);
+    replica.shards[usize::from(owner)].init_partition(
+        namespace,
+        Some(superblock),
+        recovered_state,
+        retained,
+    );
     // Commit the namespace before stamping the rows: a partition the metadata
     // plane never heard of is a shape production cannot produce, and the shard
     // refuses to serve client traffic whose routing-row epoch it cannot match
@@ -2984,7 +3050,7 @@ mod tests {
             executor.run_until_stalled(POLL_BUDGET); // borrow acquired; task 
parks
             let grow = Rc::clone(&sim.replicas[0].shards[0]);
             executor.spawn(async move {
-                grow.init_partition(ns_grow, None, None);
+                grow.init_partition(ns_grow, None, None, None);
             });
             executor.run_until_stalled(POLL_BUDGET); // grow while the borrow 
is live
         }))
@@ -3019,7 +3085,7 @@ mod tests {
         executor.run_until_stalled(POLL_BUDGET);
         let grow = Rc::clone(&sim.replicas[0].shards[0]);
         executor.spawn(async move {
-            grow.init_partition(ns_grow, None, None);
+            grow.init_partition(ns_grow, None, None, None);
         });
         executor.run_until_stalled(POLL_BUDGET);
 
@@ -3135,6 +3201,98 @@ 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 restarted replica comes back with the partition data it had, so its
+    /// `commit_offset` does not regress.
+    ///
+    /// The simulator has no segment files, so the log has to be carried across
+    /// the rebuild by hand (`retain_partition_logs` into
+    /// `IggyPartition::adopt_retained_log`). Without that the rebuilt 
partition
+    /// is empty and reports `commit_offset` 0, which models total data loss
+    /// rather than a restart and trips the monotonicity invariant on a harness
+    /// artefact. Asserted on a backup, which is where the driver's crash
+    /// injection puts a restart.
+    #[test]
+    fn restarted_replica_keeps_its_partition_offsets() {
+        use bytes::Bytes;
+
+        
server_common::MemoryPool::init_pool(&server_common::MemoryPoolConfigOther {
+            enabled: false,
+            size: iggy_common::IggyByteSize::from(0u64),
+            bucket_capacity: 1,
+        });
+
+        let replica_count: u8 = 3;
+        let client_id: u128 = 1;
+        let network_opts = packet::PacketSimulatorOptions {
+            node_count: replica_count,
+            client_count: 1,
+            seed: 0x5EED_0079,
+            ..packet::PacketSimulatorOptions::default()
+        };
+        let mut sim = Simulator::new(
+            usize::from(replica_count),
+            std::iter::once(client_id),
+            network_opts,
+        );
+        let client = SimClient::new(client_id);
+        let ns = IggyNamespace::new(1, 1, 0);
+        sim.init_partition(ns);
+        sim.register_client_with_primary(&client);
+
+        // Enough committed writes that an empty rebuild is unmistakable.
+        for sequence in 0..8u32 {
+            let msg = client.send_messages(ns, 
&[Bytes::from(format!("retained-{sequence}"))]);
+            sim.submit_request(client_id, 0, msg.into_generic());
+            for _ in 0..40 {
+                sim.step();
+            }
+        }
+
+        let backup = 1usize;
+        let before = sim
+            .offsets(backup, ns)
+            .expect("backup hosts the namespace")
+            .commit_offset;
+        assert!(
+            before > 0,
+            "backup must have committed some writes before the restart, else a 
\
+             rebuilt-empty partition would be indistinguishable from this 
state"
+        );
+
+        sim.replica_crash(u8::try_from(backup).expect("replica index fits 
u8"));
+        for _ in 0..50 {
+            sim.tick();
+        }
+        sim.replica_restart(u8::try_from(backup).expect("replica index fits 
u8"));
+
+        let after = sim
+            .offsets(backup, ns)
+            .expect("namespace re-materialised on restart")
+            .commit_offset;
+        assert_eq!(
+            after, before,
+            "restarted backup lost its partition offsets: the retained log did 
not \
+             carry across the rebuild"
+        );
+
+        // The data itself, not just the counter: a carried-over counter with 
an
+        // empty log would satisfy the assert above and still have lost every
+        // message.
+        let fragments = sim
+            .poll_messages(
+                backup,
+                ns,
+                PollingConsumer::Consumer(0, 0),
+                &PollingArgs::new(iggy_common::PollingStrategy::offset(0), 16, 
false),
+            )
+            .expect("poll against the re-materialised partition");
+        assert!(
+            !fragments.is_empty(),
+            "restarted backup served no messages: the offset counter came 
across \
+             but the log behind it did not"
+        );
+    }
+
     /// A lost `PrepareOk` does not wedge the metadata plane: once the acks 
flow
     /// again the primary reaches its commit quorum without client involvement.
     ///
diff --git a/core/simulator/src/workload/auditor.rs 
b/core/simulator/src/workload/auditor.rs
index ca51a6016..411a83af7 100644
--- a/core/simulator/src/workload/auditor.rs
+++ b/core/simulator/src/workload/auditor.rs
@@ -183,7 +183,6 @@ impl ServerAuditor {
         &self.stats
     }
 
-    #[must_use]
     /// The action of an outstanding request, if one is recorded for `key`.
     /// Diagnostic only: names what a stalled run is waiting on, which the bare
     /// `(client, request)` pair cannot.
diff --git a/core/simulator/src/workload/invariants.rs 
b/core/simulator/src/workload/invariants.rs
index a20b99281..d7f1b8a8d 100644
--- a/core/simulator/src/workload/invariants.rs
+++ b/core/simulator/src/workload/invariants.rs
@@ -51,20 +51,13 @@ impl Invariants {
     /// Globally:
     /// - total in-flight requests stay within the per-client queue ceiling.
     ///
-    /// Crashed replicas are skipped and their last-seen marks retained, which 
is
-    /// correct for `view` (the superblock carries it across a restart) but NOT
-    /// for partition `commit_offset`, because the simulator does not retain
-    /// partition data across a restart the way it retains the metadata WAL and
-    /// the superblocks. `IggyPartitions` is dropped and rebuilt empty, so a
-    /// restarted replica reports `commit_offset` 0 and this trips, reporting a
-    /// harness limitation as a consensus regression.
-    ///
-    /// So restart injection (`WorkloadOptions::restart_per_tick_ratio`) is
-    /// metadata-plane only today. Driving it with partition traffic needs the
-    /// partition journal held by the harness, per-namespace, exactly as
-    /// `SimReplica::metadata_journal` already is; a real server's segments 
are on
-    /// disk and do survive, so retaining them is the faithful model, and the
-    /// alternative of relaxing this check would instead model total data loss.
+    /// Crashed replicas are skipped and their last-seen marks retained, which
+    /// holds across a restart for both quantities: the superblock carries 
`view`,
+    /// and the harness carries each partition's log
+    /// (`Simulator::retain_partition_logs`) so a rebuilt partition recovers 
the
+    /// offsets it had rather than reporting zero. The latter is the harness
+    /// standing in for the segment files a real server recovers from; without 
it
+    /// this check trips on a discarded log and calls it a consensus 
regression.
     ///
     /// # Panics
     /// On any regression or in-flight overflow. The workload seed is in the

Reply via email to