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

mmodzelewski 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 ab257a039 refactor(consensus): single restore constructor + 
wedge/replay/ack fixes (#3909)
ab257a039 is described below

commit ab257a03996cc67d2ee2efad5501d8953d086bcd
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Wed Aug 19 06:45:54 2026 +0200

    refactor(consensus): single restore constructor + wedge/replay/ack fixes 
(#3909)
---
 core/configs/src/server_config/cluster.rs          |  43 ++++
 core/configs/src/server_config/defaults.rs         |   5 +
 core/consensus/src/fatal.rs                        |   4 +
 core/consensus/src/impls.rs                        | 110 ++++++++++
 core/consensus/src/plane_helpers.rs                |  17 +-
 .../tests/server/partition_view_durability_vsr.rs  |  22 +-
 core/metadata/src/impls/metadata.rs                |  12 +-
 core/metadata/src/impls/recovery.rs                |  38 +++-
 core/partitions/src/iggy_partition.rs              |   4 +-
 core/server/config.toml                            |   9 +
 core/server/src/bootstrap.rs                       | 243 ++++++++++++---------
 core/server/src/partition_helpers.rs               |  82 +++----
 core/shard/src/lib.rs                              |  64 +++++-
 13 files changed, 470 insertions(+), 183 deletions(-)

diff --git a/core/configs/src/server_config/cluster.rs 
b/core/configs/src/server_config/cluster.rs
index ad767c6b2..a1874a640 100644
--- a/core/configs/src/server_config/cluster.rs
+++ b/core/configs/src/server_config/cluster.rs
@@ -50,6 +50,10 @@ const MIN_HEARTBEAT_TO_COMMIT_BROADCAST_RATIO: u32 = 4;
 /// than escalating a progressing view change into a fresh cluster-wide 
election.
 const MIN_STATUS_TO_RETRANSMIT_RATIO: u32 = 4;
 
+/// Floor for a nonzero `superblock_wedged_fatal_timeout`: a shorter window
+/// would fail-stop the process on a transient disk hiccup instead of a wedge.
+const MIN_SUPERBLOCK_WEDGED_FATAL_TIMEOUT: Duration = Duration::from_secs(30);
+
 /// Default recovering-replica probe-attempt ceiling. Duplicated here rather
 /// than imported so `core/configs` keeps off a build-time edge onto
 /// `core/consensus` (mirroring [`super::partition`]); `core/server`'s
@@ -172,6 +176,16 @@ fn default_repair_chunk_max() -> usize {
     SERVER_CONFIG.cluster.repair_chunk_max as usize
 }
 
+/// serde fallback for configs written before the field existed; the value
+/// itself lives in `core/server/config.toml` like every other default.
+fn default_superblock_wedged_fatal_timeout() -> IggyDuration {
+    SERVER_CONFIG
+        .cluster
+        .superblock_wedged_fatal_timeout
+        .parse()
+        .unwrap()
+}
+
 #[serde_as]
 #[derive(Debug, Deserialize, Serialize, Clone, ConfigEnv)]
 #[serde(deny_unknown_fields)]
@@ -262,6 +276,17 @@ pub struct ClusterConfig {
     /// be > 0 and <= `MAX_REPAIR_CHUNK_MAX`.
     #[serde(default = "default_repair_chunk_max")]
     pub repair_chunk_max: usize,
+    /// How long the metadata superblock may stay unwritable before the replica
+    /// fail-stops. While wedged the replica is already fenced quorum-invisible
+    /// and peers elect around it; this converts the log-only limp into a
+    /// distinct exit status a supervisor can act on. Zero (and the `0` /
+    /// `disabled` / `unlimited` sentinels, which all parse to zero) disables
+    /// the fail-stop; nonzero values below
+    /// `MIN_SUPERBLOCK_WEDGED_FATAL_TIMEOUT` are rejected at boot.
+    #[serde(default = "default_superblock_wedged_fatal_timeout")]
+    #[serde_as(as = "DisplayFromStr")]
+    #[config_env(leaf)]
+    pub superblock_wedged_fatal_timeout: IggyDuration,
     /// Full roster of cluster members. Intended to be byte-identical across
     /// every node so operators ship one config. The running node's identity
     /// is supplied out-of-band via the `--replica-id` CLI flag, which
@@ -747,6 +772,22 @@ impl Validatable<ConfigurationError> for ClusterConfig {
             return Err(ConfigurationError::InvalidConfigurationValue);
         }
 
+        // Also ahead of the enabled gate: a solo replica persists the
+        // superblock too, so the fail-stop window applies regardless.
+        let superblock_fatal_window = 
self.superblock_wedged_fatal_timeout.get_duration();
+        if !superblock_fatal_window.is_zero()
+            && superblock_fatal_window < MIN_SUPERBLOCK_WEDGED_FATAL_TIMEOUT
+        {
+            eprintln!(
+                "Invalid cluster configuration: 
cluster.superblock_wedged_fatal_timeout '{}' must \
+                 be zero (disabled) or at least {}s (a shorter window 
fail-stops the process on a \
+                 transient disk hiccup)",
+                self.superblock_wedged_fatal_timeout,
+                MIN_SUPERBLOCK_WEDGED_FATAL_TIMEOUT.as_secs()
+            );
+            return Err(ConfigurationError::InvalidConfigurationValue);
+        }
+
         if !self.enabled {
             return Ok(());
         }
@@ -1402,6 +1443,7 @@ mod tests {
             view_probe_attempts_max: default_view_probe_attempts_max(),
             repair_retry_interval: default_repair_retry_interval(),
             repair_chunk_max: default_repair_chunk_max(),
+            superblock_wedged_fatal_timeout: 
default_superblock_wedged_fatal_timeout(),
             nodes: Vec::new(),
             auth: ClusterAuthConfig {
                 enabled: true,
@@ -1737,6 +1779,7 @@ mod cluster_validate_tests {
             view_probe_attempts_max: default_view_probe_attempts_max(),
             repair_retry_interval: default_repair_retry_interval(),
             repair_chunk_max: default_repair_chunk_max(),
+            superblock_wedged_fatal_timeout: 
default_superblock_wedged_fatal_timeout(),
             nodes,
             auth: ClusterAuthConfig::default(),
             tls: ClusterTlsConfig::default(),
diff --git a/core/configs/src/server_config/defaults.rs 
b/core/configs/src/server_config/defaults.rs
index 7238769b5..1705e861b 100644
--- a/core/configs/src/server_config/defaults.rs
+++ b/core/configs/src/server_config/defaults.rs
@@ -92,6 +92,11 @@ impl Default for ClusterConfig {
                 .view_change_status_timeout
                 .parse()
                 .unwrap(),
+            superblock_wedged_fatal_timeout: SERVER_CONFIG
+                .cluster
+                .superblock_wedged_fatal_timeout
+                .parse()
+                .unwrap(),
             request_start_view_retransmit_interval: SERVER_CONFIG
                 .cluster
                 .request_start_view_retransmit_interval
diff --git a/core/consensus/src/fatal.rs b/core/consensus/src/fatal.rs
index d08983411..341bda389 100644
--- a/core/consensus/src/fatal.rs
+++ b/core/consensus/src/fatal.rs
@@ -26,6 +26,10 @@ pub enum FatalReason {
     /// back. The durable log is intact up to the previous op, so recovery 
re-derives
     /// the frontier and restarting is the repair.
     UnreconcilableLogFrontier = 2,
+    /// The superblock stayed unwritable past the configured fail-stop window.
+    /// The replica was already fenced quorum-invisible, so exiting hands the
+    /// wedge to a supervisor instead of a log reader.
+    SuperblockWedged = 3,
 }
 
 impl FatalReason {
diff --git a/core/consensus/src/impls.rs b/core/consensus/src/impls.rs
index 50bb8523d..33838ee63 100644
--- a/core/consensus/src/impls.rs
+++ b/core/consensus/src/impls.rs
@@ -1113,6 +1113,49 @@ where
     clock: ConsensusClock,
 }
 
+/// Boot-time timer set for a consensus group, one struct so every plane's
+/// restore path applies the same values in the same place.
+#[derive(Debug, Clone, Copy)]
+pub struct ConsensusTimers {
+    pub normal_heartbeat_ticks: u64,
+    pub commit_message_ticks: u64,
+    pub prepare_ticks: u64,
+    pub view_change_retransmit_ticks: u64,
+    pub view_change_status_ticks: u64,
+    pub request_start_view_ticks: u64,
+    pub probe_attempts_max: u32,
+}
+
+/// How a restored replica joins its group.
+#[derive(Debug, Clone, Copy)]
+pub enum JoinMode {
+    /// Fresh group or solo replica: plain init; the group needs its view-0
+    /// primary to exist.
+    Init,
+    /// Prior life detected: join quorum-invisible and probe for the current
+    /// view (`RequestStartView`) instead of resuming a role the cluster may
+    /// have elected past.
+    ProbeAsBackup {
+        /// Also await a state-transfer offer before serving, replacing
+        /// snapshot-shaped state from the live primary.
+        await_state_transfer: bool,
+    },
+}
+
+/// Restored state handed to [`VsrConsensus::restored`].
+#[derive(Debug, Clone, Copy)]
+pub struct VsrRestore<'a> {
+    pub timers: &'a ConsensusTimers,
+    /// `(view, log_view)` read back from the group's durable superblock.
+    pub durable_view: Option<(u32, u32)>,
+    /// View inferred from the last journaled prepare, consulted only when no
+    /// durable record exists; `log_view` cannot be inferred and stays 0.
+    pub view_fallback: Option<u32>,
+    /// Non-zero boot incarnation; `None` keeps the default.
+    pub incarnation: Option<u128>,
+    pub join: JoinMode,
+}
+
 impl<B: MessageBus, P: Pipeline<Entry = PipelineEntry>> VsrConsensus<B, P> {
     /// # Panics
     /// - If `replica >= replica_count`.
@@ -1136,6 +1179,73 @@ impl<B: MessageBus, P: Pipeline<Entry = PipelineEntry>> 
VsrConsensus<B, P> {
         )
     }
 
+    /// Restore constructor: the one ordered boot path for every plane, so the
+    /// metadata and partition planes cannot diverge in restore order. Timers
+    /// first, then the durable view, then role selection - a probe must never
+    /// advertise a view older than the recorded one.
+    ///
+    /// # Panics
+    /// - If `replica >= replica_count`.
+    /// - If `replica_count < 1`.
+    pub fn restored(
+        cluster: u128,
+        replica: u8,
+        replica_count: u8,
+        group: u64,
+        message_bus: B,
+        pipeline: P,
+        restore: VsrRestore<'_>,
+    ) -> Self {
+        let mut consensus = Self::new(
+            cluster,
+            replica,
+            replica_count,
+            group,
+            message_bus,
+            pipeline,
+        );
+        let timers = restore.timers;
+        consensus.set_normal_heartbeat_ticks(timers.normal_heartbeat_ticks);
+        consensus.set_commit_message_ticks(timers.commit_message_ticks);
+        consensus.set_prepare_ticks(timers.prepare_ticks);
+        
consensus.set_view_change_retransmit_ticks(timers.view_change_retransmit_ticks);
+        
consensus.set_view_change_status_ticks(timers.view_change_status_ticks);
+        
consensus.set_request_start_view_ticks(timers.request_start_view_ticks);
+        consensus.set_probe_attempts_max(timers.probe_attempts_max);
+        if let Some(incarnation) = restore.incarnation {
+            consensus.set_incarnation(incarnation);
+        }
+        if let Some((view, log_view)) = restore.durable_view {
+            // The one line proving the durable record was READ BACK, not 
merely
+            // written: a replica that came back at view 0 is otherwise
+            // indistinguishable from one that resumed correctly until it 
votes.
+            tracing::info!(
+                group,
+                view,
+                log_view,
+                "restored group view from its superblock"
+            );
+            consensus.set_view(view);
+            consensus.set_log_view(log_view);
+            consensus.mark_superblock_durable(view, log_view);
+        } else if let Some(view) = restore.view_fallback {
+            consensus.set_view(view);
+        }
+        match restore.join {
+            JoinMode::Init => consensus.init(),
+            JoinMode::ProbeAsBackup {
+                await_state_transfer,
+            } => {
+                consensus.init_as_backup();
+                consensus.begin_view_probe();
+                if await_state_transfer {
+                    consensus.begin_state_transfer_await();
+                }
+            }
+        }
+        consensus
+    }
+
     /// [`Self::new`] with an explicit time source. Simulator and clock
     /// tests only; production wiring stays on the system-clock default.
     ///
diff --git a/core/consensus/src/plane_helpers.rs 
b/core/consensus/src/plane_helpers.rs
index 57c6d296a..4e4bcd7ef 100644
--- a/core/consensus/src/plane_helpers.rs
+++ b/core/consensus/src/plane_helpers.rs
@@ -771,7 +771,14 @@ pub fn panic_if_hash_chain_would_break_in_same_view(
     }
 }
 
-// TODO: Figure out how to make this check the journal if it contains the 
prepare.
+/// Ack a prepare back to its primary once the owning plane vouches for it.
+///
+/// `is_persisted` is the caller's journal-containment verdict for `header`:
+/// consensus is sans-io and cannot consult the journal itself, so the plane
+/// that owns the journal must vouch that this exact prepare is durable before
+/// the ack leaves. `false` withholds the ack; the primary's retransmit
+/// re-drives it once a later persist succeeds.
+///
 /// # Panics
 /// - If `header.command` is not `Command::Prepare`.
 /// - If `header.view > consensus.view()`.
@@ -779,7 +786,7 @@ pub fn panic_if_hash_chain_would_break_in_same_view(
 pub async fn send_prepare_ok<B, P>(
     consensus: &VsrConsensus<B, P>,
     header: &PrepareHeader,
-    is_persisted: Option<bool>,
+    is_persisted: bool,
 ) where
     B: MessageBus,
     P: Pipeline<Entry = PipelineEntry>,
@@ -794,7 +801,7 @@ pub async fn send_prepare_ok<B, P>(
         return;
     }
 
-    if is_persisted == Some(false) {
+    if !is_persisted {
         return;
     }
 
@@ -1227,7 +1234,7 @@ mod tests {
             ..Default::default()
         };
 
-        futures::executor::block_on(send_prepare_ok(&consensus, 
&prepare_header, Some(true)));
+        futures::executor::block_on(send_prepare_ok(&consensus, 
&prepare_header, true));
 
         let mut buf = Vec::new();
         consensus.drain_loopback_into(&mut buf);
@@ -1874,7 +1881,7 @@ mod tests {
             ..Default::default()
         };
 
-        futures::executor::block_on(send_prepare_ok(&consensus, 
&prepare_header, Some(true)));
+        futures::executor::block_on(send_prepare_ok(&consensus, 
&prepare_header, true));
 
         let mut buf = Vec::new();
         consensus.drain_loopback_into(&mut buf);
diff --git a/core/integration/tests/server/partition_view_durability_vsr.rs 
b/core/integration/tests/server/partition_view_durability_vsr.rs
index 9d266e447..1e8a9480a 100644
--- a/core/integration/tests/server/partition_view_durability_vsr.rs
+++ b/core/integration/tests/server/partition_view_durability_vsr.rs
@@ -52,9 +52,11 @@ const MESSAGES_COUNT: u32 = 10;
 /// then the restarted node's probe and repair), and CI runners are slow. 60s 
bounds
 /// the worst case without hanging the suite.
 const CONVERGE_TIMEOUT: Duration = Duration::from_secs(60);
-/// Boot line reporting the `(view, log_view)` a partition group restored from 
its
-/// superblock: the recovered replica's own account of what it read back.
-const RESTORED_VIEW_MARKER: &str = "restored partition view from its 
superblock";
+/// Boot line reporting the `(view, log_view)` a consensus group restored from
+/// its superblock: the recovered replica's own account of what it read back.
+/// One line per group since the unified restore constructor, so the parser
+/// filters the metadata group's line out by its `group` field.
+const RESTORED_VIEW_MARKER: &str = "restored group view from its superblock";
 const POLL_INTERVAL: Duration = Duration::from_millis(250);
 
 #[iggy_harness(cluster_nodes = 3, server(system.sharding.cpu_allocation = 
"0..1"))]
@@ -200,17 +202,23 @@ fn restored_partition_view(harness: &TestHarness, node: 
usize) -> Option<(u32, u
     }
     log.lines()
         .filter(|line| line.contains(RESTORED_VIEW_MARKER))
+        // The metadata group logs the same restore line; its view advances
+        // independently of the partition group under test, so folding it into
+        // the max would let a metadata view change satisfy a partition assert.
+        .filter(|line| {
+            field(line, " group=") != 
Some(iggy_binary_protocol::namespace::METADATA_GROUP)
+        })
         .filter_map(|line| {
             // Leading space so the `view` key cannot match inside `log_view`.
-            let view = field(line, " view=")?;
-            let log_view = field(line, " log_view=")?;
+            let view = u32::try_from(field(line, " view=")?).ok()?;
+            let log_view = u32::try_from(field(line, " log_view=")?).ok()?;
             Some((view, log_view))
         })
         .max()
 }
 
-/// Value of a space-prefixed `key=<u32>` tracing field.
-fn field(line: &str, key: &str) -> Option<u32> {
+/// Value of a space-prefixed `key=<u64>` tracing field.
+fn field(line: &str, key: &str) -> Option<u64> {
     let start = line.find(key)? + key.len();
     line[start..]
         .split(|character: char| !character.is_ascii_digit())
diff --git a/core/metadata/src/impls/metadata.rs 
b/core/metadata/src/impls/metadata.rs
index db84be4af..ee67e1097 100644
--- a/core/metadata/src/impls/metadata.rs
+++ b/core/metadata/src/impls/metadata.rs
@@ -3548,9 +3548,17 @@ where
         if !self.persist_superblock_if_needed(consensus).await {
             return;
         }
+        // Containment, not occupancy: after a refused append (slot collision,
+        // view-change race) the slot can hold a DIFFERENT prepare at this op,
+        // and acking it would vouch durability for bytes this replica never
+        // journaled. Checksum equality withholds the ack; the primary's
+        // retransmit re-drives it once the right prepare lands.
         let journal = self.journal.as_ref().unwrap();
-        let persisted = journal.handle().header(header.op as usize).is_some();
-        send_prepare_ok_common(consensus, header, Some(persisted)).await;
+        let persisted = journal
+            .handle()
+            .header(header.op as usize)
+            .is_some_and(|stored| stored.checksum == header.checksum);
+        send_prepare_ok_common(consensus, header, persisted).await;
     }
 }
 
diff --git a/core/metadata/src/impls/recovery.rs 
b/core/metadata/src/impls/recovery.rs
index f192d38f7..afdec53b9 100644
--- a/core/metadata/src/impls/recovery.rs
+++ b/core/metadata/src/impls/recovery.rs
@@ -358,6 +358,7 @@ pub async fn recover<M>(
     journal_slots: usize,
     clients_table_max: usize,
     seed_baseline: impl FnOnce(&M),
+    on_replayed_logout: impl Fn(&M, u128, iggy_common::IggyTimestamp),
 ) -> Result<RecoveredMetadata<M>, RecoveryError>
 where
     M: StateMachine<Input = Message<PrepareHeader>, Error = IggyError>
@@ -618,10 +619,15 @@ where
         }
         if header.operation == Operation::Logout {
             client_table.remove_client(header.client);
-            // TODO: the commit paths also run `remove_consumer_group_member`
-            // here; recovery has no `StreamsFrontend` bound, so replayed
-            // logouts leave stale group members (pre-existing, harmless for
-            // dead connections but a divergence from the live apply).
+            // Logout's only state-machine effect, mirrored from the live
+            // commit path: the caller drops the client from its consumer
+            // groups (`remove_consumer_group_member`) so replay and live
+            // apply converge on the same group membership.
+            on_replayed_logout(
+                &mux_stm,
+                header.client,
+                iggy_common::IggyTimestamp::from(header.timestamp),
+            );
             last_applied_op = Some(header.op);
             continue;
         }
@@ -956,6 +962,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -981,6 +988,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1016,6 +1024,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1073,6 +1082,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1132,6 +1142,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await;
 
@@ -1182,6 +1193,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1219,6 +1231,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1288,6 +1301,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1363,6 +1377,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1428,6 +1443,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1477,12 +1493,14 @@ mod tests {
             journal.storage_ref().fsync().await.unwrap();
         }
 
+        let replayed_logouts = std::cell::RefCell::new(Vec::new());
         let recovered = recover::<TestStm>(
             dir.path(),
             SOLO,
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, client, _| replayed_logouts.borrow_mut().push(client),
         )
         .await
         .unwrap();
@@ -1491,6 +1509,11 @@ mod tests {
             None,
             "logged-out session must not be resurrected"
         );
+        assert_eq!(
+            replayed_logouts.into_inner(),
+            vec![CLIENT],
+            "replay must hand the logged-out client to the group-removal hook"
+        );
         assert_eq!(recovered.last_applied_op, Some(2));
     }
 
@@ -1545,6 +1568,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await;
         assert!(matches!(
@@ -1577,6 +1601,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await;
         assert!(matches!(
@@ -1608,6 +1633,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1660,6 +1686,7 @@ mod tests {
                 journal::prepare_journal::DEFAULT_SLOT_COUNT,
                 CLIENTS_TABLE_MAX,
                 |_| {},
+                |_, _, _| {},
             )
             .await;
             match result {
@@ -1700,6 +1727,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         {
@@ -1764,6 +1792,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
@@ -1816,6 +1845,7 @@ mod tests {
             journal::prepare_journal::DEFAULT_SLOT_COUNT,
             CLIENTS_TABLE_MAX,
             |_| {},
+            |_, _, _| {},
         )
         .await
         .unwrap();
diff --git a/core/partitions/src/iggy_partition.rs 
b/core/partitions/src/iggy_partition.rs
index 76267af57..4a366076d 100644
--- a/core/partitions/src/iggy_partition.rs
+++ b/core/partitions/src/iggy_partition.rs
@@ -4717,7 +4717,9 @@ where
         // consumer-offset ops (via `apply_replicated_operation`) append
         // to that journal before `send_prepare_ok` fires, so every op
         // that reaches here is journal-backed and ACKs as durable.
-        send_prepare_ok_common(self.consensus(), header, Some(true)).await;
+        // (`header_by_op` is a linear scan, so re-proving that here would
+        // put O(journal) on every ack; the call-order invariant stands in.)
+        send_prepare_ok_common(self.consensus(), header, true).await;
     }
 }
 
diff --git a/core/server/config.toml b/core/server/config.toml
index d10bfa4e0..f420aeeb3 100644
--- a/core/server/config.toml
+++ b/core/server/config.toml
@@ -559,6 +559,15 @@ repair_retry_interval = "1s"
 # the queue and drops frames. Must be > 0 and <= 1024.
 repair_chunk_max = 128
 
+# How long the metadata superblock may stay unwritable before the replica
+# fail-stops (duration). A replica that cannot persist its view is already
+# fenced quorum-invisible and retries with capped backoff; past this window the
+# process exits with a distinct status so a supervisor restarts or replaces it
+# instead of an operator finding the wedge in logs. "0" disables the fail-stop
+# and leaves the replica fenced indefinitely. Nonzero values must be at least
+# 30s so a transient disk hiccup cannot kill the process.
+superblock_wedged_fatal_timeout = "2m"
+
 # Replica-to-replica authentication (PSK + BLAKE3 keyed-MAC handshake).
 [cluster.auth]
 # When true, every replica peer must complete the authenticated handshake or be
diff --git a/core/server/src/bootstrap.rs b/core/server/src/bootstrap.rs
index a78cc847d..01ef2d1fa 100644
--- a/core/server/src/bootstrap.rs
+++ b/core/server/src/bootstrap.rs
@@ -26,7 +26,7 @@ use crate::dispatch::{
 use crate::http;
 use crate::partition_helpers::{
     build_partition_fresh, configure_consumer_offsets, ensure_initial_segment,
-    open_partition_superblock, restore_partition_view,
+    open_partition_superblock,
 };
 use crate::segment_recovery::{RecoveredSegment, load_persisted_segments};
 use crate::server_error::{ServerError, ShardJoinFailure, ShardJoinFailureKind};
@@ -37,8 +37,8 @@ use configs::sharding::{
     INBOX_CAPACITY_MAX, SHUTDOWN_DRAIN_TIMEOUT_MAX, SHUTDOWN_POLL_INTERVAL_MAX,
 };
 use consensus::{
-    ClientTable, LocalPipeline, MetadataHandle, PartitionsHandle, 
PipelineEntry, Sequencer,
-    VsrConsensus,
+    ClientTable, ConsensusTimers, JoinMode, LocalPipeline, MetadataHandle, 
PartitionsHandle,
+    PipelineEntry, Sequencer, VsrConsensus, VsrRestore,
 };
 // `try_send` / `try_recv` resolve through these traits on `MAsyncTx` /
 // `MAsyncRx`; the metadata-handoff loops below depend on the
@@ -1016,6 +1016,11 @@ async fn shard_main(
                 |mux_stm| {
                     ensure_default_root_user(mux_stm);
                 },
+                |mux_stm, client, timestamp| {
+                    mux_stm
+                        .streams()
+                        .remove_consumer_group_member(client, timestamp);
+                },
             )
             .await
             .map_err(ServerError::MetadataRecovery)?;
@@ -2012,6 +2017,7 @@ async fn build_shard_for_thread(
     // Repair pacing is shared by both planes' repair loops, so it is a
     // per-shard tunable set once here rather than per consensus group.
     shard.set_repair_retry_ticks(repair_retry_ticks(config));
+    
shard.set_superblock_wedged_fatal_failures(superblock_wedged_fatal_failures(config));
     shard.set_served_segment_cache_bytes_max(
         config
             .partition
@@ -2093,6 +2099,28 @@ fn duration_to_ticks(interval: Duration) -> u64 {
     u64::try_from(ticks.max(1)).unwrap_or(u64::MAX)
 }
 
+/// `[cluster] superblock_wedged_fatal_timeout` as a consecutive-failure count.
+/// Retries pin at the backoff cap after warmup, so the window divided by
+/// [`journal::superblock::SUPERBLOCK_RETRY_BACKOFF_MAX_MICROS`] bounds how
+/// long a wedged replica may limp before it fail-stops. Zero stays zero
+/// (fail-stop disabled).
+pub(crate) fn superblock_wedged_fatal_failures(config: &ServerConfig) -> u64 {
+    superblock_window_to_failures(
+        config
+            .cluster
+            .superblock_wedged_fatal_timeout
+            .get_duration(),
+    )
+}
+
+fn superblock_window_to_failures(window: Duration) -> u64 {
+    if window.is_zero() {
+        return 0;
+    }
+    let cap_micros = 
u128::from(journal::superblock::SUPERBLOCK_RETRY_BACKOFF_MAX_MICROS);
+    u64::try_from((window.as_micros() / cap_micros).max(1)).unwrap_or(u64::MAX)
+}
+
 /// `[cluster] heartbeat_timeout` in consensus ticks. Every consensus group
 /// (metadata and per-partition planes alike) gets the same window: the failure
 /// it guards against - a primary that stopped heartbeating - is host-level, 
not
@@ -2181,6 +2209,20 @@ pub(crate) fn request_start_view_ticks(config: 
&ServerConfig) -> u64 {
     )
 }
 
+/// The full `[cluster]` timer set every consensus group boots with, built
+/// once so the planes cannot diverge in what they apply.
+pub(crate) fn consensus_timers(config: &ServerConfig) -> ConsensusTimers {
+    ConsensusTimers {
+        normal_heartbeat_ticks: cluster_heartbeat_ticks(config),
+        commit_message_ticks: commit_broadcast_ticks(config),
+        prepare_ticks: prepare_retransmit_ticks(config),
+        view_change_retransmit_ticks: view_change_retransmit_ticks(config),
+        view_change_status_ticks: view_change_status_ticks(config),
+        request_start_view_ticks: request_start_view_ticks(config),
+        probe_attempts_max: config.cluster.view_probe_attempts_max,
+    }
+}
+
 /// `[cluster] repair_retry_interval` in consensus ticks: how long a stalled
 /// journal-repair stream waits before re-requesting its window. Both planes'
 /// repair loops share it, so it is applied once per shard (not per consensus
@@ -2235,50 +2277,10 @@ fn restore_metadata_consensus(
     );
     let prepare_queue_depth = config.metadata.prepare_queue_depth;
 
-    let mut consensus = VsrConsensus::new(
-        topology.cluster_id,
-        topology.self_replica_id,
-        replica_count,
-        server_common::sharding::METADATA_GROUP,
-        bus,
-        // Request queue keeps the stock 2x ratio over the prepare queue
-        // (32 -> 64 at defaults): buffered requests are cheap relative to
-        // in-flight prepares and drain as prepares commit.
-        LocalPipeline::with_capacities(prepare_queue_depth, 
prepare_queue_depth * 2),
-    );
-    consensus.set_normal_heartbeat_ticks(cluster_heartbeat_ticks(config));
-    consensus.set_commit_message_ticks(commit_broadcast_ticks(config));
-    consensus.set_prepare_ticks(prepare_retransmit_ticks(config));
-    
consensus.set_view_change_retransmit_ticks(view_change_retransmit_ticks(config));
-    consensus.set_view_change_status_ticks(view_change_status_ticks(config));
-    consensus.set_request_start_view_ticks(request_start_view_ticks(config));
-    consensus.set_probe_attempts_max(config.cluster.view_probe_attempts_max);
-    // Fresh random incarnation each boot, so a StartView addressed to a 
previous
-    // incarnation still in flight is ignored (`handle_start_view` guard). `| 
1`
-    // guarantees the non-zero the guard treats as set. The deterministic 
simulator
-    // overrides this with a seed-derived value bumped per restart.
-    consensus.set_incarnation(rand::random::<u128>() | 1);
-
     let last_header = journal
         .last_op()
         .and_then(|op| usize::try_from(op).ok())
         .and_then(|op| journal.header(op).map(|header| *header));
-    // View and log_view come from the durable superblock when present. A 
present but
-    // unreadable superblock already refused boot in `recover()`, so reaching 
the
-    // `else` means it is genuinely absent: a fresh node, or one that took 
writes but
-    // never checkpointed or changed view. There, inferring the view from the 
last WAL
-    // prepare is safe, since the persist-before-send gate guarantees this 
replica
-    // never externalized a view beyond what a re-probe re-derives, and it 
re-probes
-    // as a backup below. log_view cannot be inferred and stays 0 until the 
next
-    // superblock write.
-    if let Some(state) = recovered_state {
-        consensus.set_view(state.view);
-        consensus.set_log_view(state.log_view);
-        consensus.mark_superblock_durable(state.view, state.log_view);
-    } else if let Some(header) = last_header {
-        consensus.set_view(header.view);
-    }
-
     // On a RESTART in a cluster, rejoin as a quorum-invisible backup and
     // probe for the current view (`RequestStartView`): the view's primary
     // answers with a `StartView`, the replica adopts it as a backup, and
@@ -2295,18 +2297,51 @@ fn restore_metadata_consensus(
     // and an empty journal; gating on the WAL alone would `init()` it into
     // `Status::Normal` as primary for a view the cluster may have moved past,
     // with `ceded_primaryship` false and no probe to correct it.
-    if replica_count > 1 && (restored_op > 0 || recovered_state.is_some()) {
-        consensus.init_as_backup();
-        consensus.begin_view_probe();
-        // Restart in a cluster: replace snapshot-shaped metadata state
-        // (snapshot + client table) from the live primary the probe finds,
-        // then journal-repair the tail. If the probe exhausts instead --
-        // full-cluster bootstrap, nobody live to fetch from -- the election
-        // fallback clears the stage and this local recovery stands.
-        consensus.begin_state_transfer_await();
+    //
+    // The rejoin also awaits a state transfer: snapshot-shaped metadata state
+    // (snapshot + client table) is replaced from the live primary the probe
+    // finds, then journal repair fills the tail. If the probe exhausts
+    // instead -- full-cluster bootstrap, nobody live to fetch from -- the
+    // election fallback clears the stage and this local recovery stands.
+    let join = if replica_count > 1 && (restored_op > 0 || 
recovered_state.is_some()) {
+        JoinMode::ProbeAsBackup {
+            await_state_transfer: true,
+        }
     } else {
-        consensus.init();
-    }
+        JoinMode::Init
+    };
+    let timers = consensus_timers(config);
+    let consensus = VsrConsensus::restored(
+        topology.cluster_id,
+        topology.self_replica_id,
+        replica_count,
+        server_common::sharding::METADATA_GROUP,
+        bus,
+        // Request queue keeps the stock 2x ratio over the prepare queue
+        // (32 -> 64 at defaults): buffered requests are cheap relative to
+        // in-flight prepares and drain as prepares commit.
+        LocalPipeline::with_capacities(prepare_queue_depth, 
prepare_queue_depth * 2),
+        VsrRestore {
+            timers: &timers,
+            // View and log_view come from the durable superblock when present.
+            // A present but unreadable superblock already refused boot in
+            // `recover()`, so no durable record means genuinely absent: a
+            // fresh node, or one that took writes but never checkpointed or
+            // changed view. There, inferring the view from the last WAL
+            // prepare is safe, since the persist-before-send gate guarantees
+            // this replica never externalized a view beyond what a re-probe
+            // re-derives, and it re-probes as a backup.
+            durable_view: recovered_state.map(|state| (state.view, 
state.log_view)),
+            view_fallback: last_header.map(|header| header.view),
+            // Fresh random incarnation each boot, so a StartView addressed to
+            // a previous incarnation still in flight is ignored
+            // (`handle_start_view` guard). `| 1` guarantees the non-zero the
+            // guard treats as set. The deterministic simulator overrides this
+            // with a seed-derived value bumped per restart.
+            incarnation: Some(rand::random::<u128>() | 1),
+            join,
+        },
+    );
     consensus.sequencer().set_sequence(restored_op);
     // A SOLO replica's durable journal head IS its commit point: quorum is
     // 1-of-1, so an entry commits the instant it is durable, and the acks
@@ -2388,39 +2423,6 @@ fn restore_metadata_consensus(
     consensus
 }
 
-#[allow(clippy::too_many_arguments)]
-/// Build the ticked-and-bounded consensus a loaded partition group joins
-/// with; the fresh-create path configures its own inside
-/// `build_partition_fresh`.
-fn loaded_partition_consensus(
-    config: &ServerConfig,
-    namespace: IggyNamespace,
-    cluster_id: u128,
-    self_replica_id: u8,
-    replica_count: u8,
-    bus: Rc<IggyMessageBus>,
-) -> VsrConsensus<Rc<IggyMessageBus>> {
-    // Request queue holds 2x the prepare depth (buffered requests drain as
-    // prepares commit); depth is the per-partition `[partition]` knob.
-    let prepare_queue_depth = config.partition.prepare_queue_depth;
-    let consensus = VsrConsensus::new(
-        cluster_id,
-        self_replica_id,
-        replica_count,
-        namespace.inner(),
-        bus,
-        LocalPipeline::with_capacities(prepare_queue_depth, 
prepare_queue_depth * 2),
-    );
-    consensus.set_normal_heartbeat_ticks(cluster_heartbeat_ticks(config));
-    consensus.set_commit_message_ticks(commit_broadcast_ticks(config));
-    consensus.set_prepare_ticks(prepare_retransmit_ticks(config));
-    
consensus.set_view_change_retransmit_ticks(view_change_retransmit_ticks(config));
-    consensus.set_view_change_status_ticks(view_change_status_ticks(config));
-    consensus.set_request_start_view_ticks(request_start_view_ticks(config));
-    consensus.set_probe_attempts_max(config.cluster.view_probe_attempts_max);
-    consensus
-}
-
 /// Recover this partition's persisted segment chain, stamping each segment
 /// with the topic's effective segment size (the per-topic value when the
 /// topic was created with one, else the shard-wide configured size).
@@ -2472,20 +2474,9 @@ async fn load_partition(
     let stream_id = namespace.stream_id();
     let topic_id = namespace.topic_id();
     let partition_id = namespace.partition_id();
-    let mut consensus = loaded_partition_consensus(
-        config,
-        namespace,
-        cluster_id,
-        self_replica_id,
-        replica_count,
-        bus,
-    );
-
     // (view, log_view) come from the group's durable superblock when present;
     // a present but unverifiable record already refused boot inside
-    // `open_partition_superblock`. Restored BEFORE choosing how to join, so
-    // the backup probe below never advertises a view older than the recorded
-    // one.
+    // `open_partition_superblock`.
     let partition_dir = config
         .system
         .get_partition_path(stream_id, topic_id, partition_id);
@@ -2498,9 +2489,6 @@ async fn load_partition(
         },
     )
     .await?;
-    if let Some(state) = recovered_state.as_ref() {
-        restore_partition_view(&mut consensus, state);
-    }
 
     // A recovered partition lost its journal state with the process: the
     // partition journal is in-memory and segments carry no op numbers, so
@@ -2512,12 +2500,34 @@ async fn load_partition(
     // at the serving peer's retention point. The probe re-broadcasts on its
     // timeout, so it needs no live mesh at boot. Single-replica groups
     // have no peer to ask and keep the plain init.
-    if replica_count > 1 {
-        consensus.init_as_backup();
-        consensus.begin_view_probe();
+    let join = if replica_count > 1 {
+        JoinMode::ProbeAsBackup {
+            await_state_transfer: false,
+        }
     } else {
-        consensus.init();
-    }
+        JoinMode::Init
+    };
+    // Request queue holds 2x the prepare depth (buffered requests drain as
+    // prepares commit); depth is the per-partition `[partition]` knob.
+    let prepare_queue_depth = config.partition.prepare_queue_depth;
+    let timers = consensus_timers(config);
+    let consensus = VsrConsensus::restored(
+        cluster_id,
+        self_replica_id,
+        replica_count,
+        namespace.inner(),
+        bus,
+        LocalPipeline::with_capacities(prepare_queue_depth, 
prepare_queue_depth * 2),
+        VsrRestore {
+            timers: &timers,
+            durable_view: recovered_state
+                .as_ref()
+                .map(|state| (state.view, state.log_view)),
+            view_fallback: None,
+            incarnation: None,
+            join,
+        },
+    );
 
     // No prepare-timestamp floor is restored here: the partition consensus
     // journal is non-durable today, so there is no persisted head to observe
@@ -3965,6 +3975,25 @@ const fn operation_triggers_partition_reconcile(op: 
Operation) -> bool {
 mod tests {
     use super::*;
 
+    #[test]
+    fn superblock_fatal_window_converts_to_capped_backoff_retries() {
+        assert_eq!(
+            superblock_window_to_failures(Duration::ZERO),
+            0,
+            "zero window must stay the disabled sentinel"
+        );
+        assert_eq!(
+            superblock_window_to_failures(Duration::from_mins(2)),
+            120,
+            "past warmup one retry rides each 1s backoff cap"
+        );
+        assert_eq!(
+            superblock_window_to_failures(Duration::from_micros(500)),
+            1,
+            "a sub-cap window still needs one failure to fire"
+        );
+    }
+
     #[test]
     fn fresh_cluster_bootstrap_requires_explicit_root_credentials() {
         assert!(matches!(
diff --git a/core/server/src/partition_helpers.rs 
b/core/server/src/partition_helpers.rs
index 8fb4a0e33..820d40f5f 100644
--- a/core/server/src/partition_helpers.rs
+++ b/core/server/src/partition_helpers.rs
@@ -28,7 +28,7 @@ use crate::offset_recovery::{load_consumer_group_offsets, 
load_consumer_offsets}
 use crate::server_error::ServerError;
 use compio::fs::create_dir_all;
 use configs::server::ServerConfig;
-use consensus::{LocalPipeline, VsrConsensus, VsrState};
+use consensus::{JoinMode, LocalPipeline, VsrConsensus, VsrRestore, VsrState};
 use iggy_common::{
     ConsumerGroupOffsets, ConsumerOffsets, IggyByteSize, IggyError, 
IggyTimestamp, PartitionStats,
     TopicRuntimeOptions,
@@ -44,7 +44,7 @@ use std::path::{Path, PathBuf};
 use std::rc::Rc;
 use std::sync::Arc;
 use std::sync::atomic::Ordering;
-use tracing::{error, info, warn};
+use tracing::{error, warn};
 
 /// Create the on-disk directory hierarchy for a partition.
 ///
@@ -498,29 +498,6 @@ pub(crate) async fn open_partition_superblock(
     Ok((Rc::new(superblock), recovered_state))
 }
 
-/// Restore `(view, log_view)` from a recovered superblock record and mark
-/// them durable (read back from disk, durable by definition). Runs BEFORE
-/// `init` / `init_as_backup` so the join path never advertises a view older
-/// than the recorded one.
-pub(crate) fn restore_partition_view(
-    consensus: &mut VsrConsensus<Rc<IggyMessageBus>>,
-    state: &VsrState,
-) {
-    // The one line proving the durable record was READ BACK, not merely 
written:
-    // the group's whole anti-regression guarantee rests on this call running, 
and
-    // a replica that came back at view 0 is otherwise indistinguishable from 
one
-    // that resumed correctly until it votes.
-    info!(
-        namespace_raw = consensus.group(),
-        view = state.view,
-        log_view = state.log_view,
-        "restored partition view from its superblock"
-    );
-    consensus.set_view(state.view);
-    consensus.set_log_view(state.log_view);
-    consensus.mark_superblock_durable(state.view, state.log_view);
-}
-
 /// Materialise a brand-new [`IggyPartition`] for a namespace that has no 
on-disk state yet.
 ///
 /// Counterpart to bootstrap's `load_partition`, which hydrates from
@@ -586,26 +563,6 @@ pub async fn build_partition_fresh(
             source
         })?;
 
-    // Request queue holds 2x the prepare depth (buffered requests drain as
-    // prepares commit); depth is the per-partition `[partition]` knob.
-    let prepare_queue_depth = config.partition.prepare_queue_depth;
-    let mut consensus = VsrConsensus::new(
-        cluster_id,
-        self_replica_id,
-        replica_count,
-        namespace.inner(),
-        bus,
-        LocalPipeline::with_capacities(prepare_queue_depth, 
prepare_queue_depth * 2),
-    );
-    
consensus.set_normal_heartbeat_ticks(crate::bootstrap::cluster_heartbeat_ticks(config));
-    
consensus.set_commit_message_ticks(crate::bootstrap::commit_broadcast_ticks(config));
-    
consensus.set_prepare_ticks(crate::bootstrap::prepare_retransmit_ticks(config));
-    consensus
-        
.set_view_change_retransmit_ticks(crate::bootstrap::view_change_retransmit_ticks(config));
-    
consensus.set_view_change_status_ticks(crate::bootstrap::view_change_status_ticks(config));
-    
consensus.set_request_start_view_ticks(crate::bootstrap::request_start_view_ticks(config));
-    consensus.set_probe_attempts_max(config.cluster.view_probe_attempts_max);
-
     // The hierarchy create above guarantees the directory exists; recover this
     // group's durable (view, log_view) before choosing how to join, so a
     // restart materialization resumes from the view it last recorded instead
@@ -622,9 +579,6 @@ pub async fn build_partition_fresh(
         },
     )
     .await?;
-    if let Some(state) = recovered_state.as_ref() {
-        restore_partition_view(&mut consensus, state);
-    }
 
     // A partition directory that already holds segment bytes is a RESTART
     // materialization, not a fresh create: this replica's group state died
@@ -635,12 +589,34 @@ pub async fn build_partition_fresh(
     // peer, byte-identical by the deterministic-roll/replicated-ciphertext
     // design. A truly fresh create keeps the plain init: every group needs
     // its view-0 primary to exist.
-    if restarted {
-        consensus.init_as_backup();
-        consensus.begin_view_probe();
+    let join = if restarted {
+        JoinMode::ProbeAsBackup {
+            await_state_transfer: false,
+        }
     } else {
-        consensus.init();
-    }
+        JoinMode::Init
+    };
+    // Request queue holds 2x the prepare depth (buffered requests drain as
+    // prepares commit); depth is the per-partition `[partition]` knob.
+    let prepare_queue_depth = config.partition.prepare_queue_depth;
+    let timers = crate::bootstrap::consensus_timers(config);
+    let consensus = VsrConsensus::restored(
+        cluster_id,
+        self_replica_id,
+        replica_count,
+        namespace.inner(),
+        bus,
+        LocalPipeline::with_capacities(prepare_queue_depth, 
prepare_queue_depth * 2),
+        VsrRestore {
+            timers: &timers,
+            durable_view: recovered_state
+                .as_ref()
+                .map(|state| (state.view, state.log_view)),
+            view_fallback: None,
+            incarnation: None,
+            join,
+        },
+    );
 
     let mut partition = IggyPartition::new(stats, consensus);
     partition.set_runtime_options(runtime_options);
diff --git a/core/shard/src/lib.rs b/core/shard/src/lib.rs
index 726bcb017..4e20691b0 100644
--- a/core/shard/src/lib.rs
+++ b/core/shard/src/lib.rs
@@ -29,10 +29,10 @@ pub use router::CONSENSUS_TICK_INTERVAL;
 use consensus::LocalPipeline;
 use consensus::{
     ChunkProgress, CommitOutcome, Consensus, ConsensusClock, DVC_HEADERS_MAX, 
DvcHeaderKind,
-    DvcSuffix, MergedLog, MetadataHandle, MuxPlane, PartitionsHandle, 
Pipeline, Plane, PlaneKind,
-    STATE_TRANSFER_MAX_DECODE_RETRIES, STATE_TRANSFER_MAX_STALL_RETRIES, 
Sequencer, Status,
-    VsrAction, VsrConsensus, build_deny_reply_from_request_header, dvc_blank, 
dvc_header_kind,
-    encode_prepare_headers, restamp_prepare_view, verify_prepare_integrity,
+    DvcSuffix, FatalReason, MergedLog, MetadataHandle, MuxPlane, 
PartitionsHandle, Pipeline, Plane,
+    PlaneKind, STATE_TRANSFER_MAX_DECODE_RETRIES, 
STATE_TRANSFER_MAX_STALL_RETRIES, Sequencer,
+    Status, VsrAction, VsrConsensus, build_deny_reply_from_request_header, 
dvc_blank,
+    dvc_header_kind, encode_prepare_headers, fatal, restamp_prepare_view, 
verify_prepare_integrity,
 };
 #[cfg(any(test, feature = "simulator"))]
 use crossfire::AsyncRxTrait;
@@ -1379,6 +1379,12 @@ where
     /// `[cluster] repair_retry_interval` at bootstrap.
     repair_retry_ticks: Cell<u32>,
 
+    /// Consecutive metadata superblock write failures tolerated before the
+    /// process fail-stops. Defaults to 0 (disabled) so the simulator and tests
+    /// keep a wedged-but-fenced replica alive; the server arms it from
+    /// `[cluster] superblock_wedged_fatal_timeout` at bootstrap.
+    superblock_wedged_fatal_failures: Cell<u64>,
+
     /// Live `[partition] transfer_served_cache_bytes_max`: the byte budget for
     /// segment payloads this shard keeps resident to serve chunk requests.
     /// Defaults to [`SERVED_SEGMENT_CACHE_BYTES_DEFAULT`]; the server
@@ -1525,6 +1531,7 @@ where
             partition_artifact_len_max: 
Cell::new(PARTITION_ARTIFACT_LEN_DEFAULT),
             repair_chunk_max: Cell::new(REPAIR_CHUNK_MAX),
             repair_retry_ticks: Cell::new(partitions::REPAIR_RETRY_TICKS),
+            superblock_wedged_fatal_failures: Cell::new(0),
             bus_max_message_size: Cell::new(DEFAULT_BUS_MAX_MESSAGE_SIZE),
             metadata_transfer_attempts: Cell::new(0),
             metadata_transfer_decode_failures: Cell::new(None),
@@ -1538,6 +1545,14 @@ where
         self.repair_retry_ticks.set(ticks);
     }
 
+    /// Arm the superblock fail-stop bound (consecutive write failures).
+    /// Called once per shard at bootstrap; the simulator and tests keep the
+    /// disabled default (0) so a wedged-but-fenced replica stays observable
+    /// in-process.
+    pub fn set_superblock_wedged_fatal_failures(&self, failures: u64) {
+        self.superblock_wedged_fatal_failures.set(failures);
+    }
+
     /// Override the serving-side resident payload budget from configuration.
     /// Called once per shard at bootstrap.
     pub fn set_served_segment_cache_bytes_max(&self, bytes: u64) {
@@ -1841,6 +1856,7 @@ where
             partition_artifact_len_max: 
Cell::new(PARTITION_ARTIFACT_LEN_DEFAULT),
             repair_chunk_max: Cell::new(REPAIR_CHUNK_MAX),
             repair_retry_ticks: Cell::new(partitions::REPAIR_RETRY_TICKS),
+            superblock_wedged_fatal_failures: Cell::new(0),
             bus_max_message_size: Cell::new(DEFAULT_BUS_MAX_MESSAGE_SIZE),
             metadata_transfer_attempts: Cell::new(0),
             metadata_transfer_decode_failures: Cell::new(None),
@@ -2393,6 +2409,12 @@ const fn parked_footprint(len: usize) -> usize {
     len.next_multiple_of(MESSAGE_ALIGN)
 }
 
+/// Whether consecutive superblock write failures crossed the fail-stop bound.
+/// `fatal_after == 0` disables the fail-stop.
+const fn superblock_wedged(failures: u64, fatal_after: u64) -> bool {
+    fatal_after != 0 && failures >= fatal_after
+}
+
 /// Reconciler passes a frame may survive before it is answered rather than 
held.
 ///
 /// Passes, not seconds, and deliberately not described in seconds: a pass 
fires
@@ -8228,6 +8250,20 @@ where
         if metadata.persist_superblock_if_needed(consensus).await {
             dispatch_vsr_actions(consensus, metadata.journal.as_ref(), 
&wire_actions).await;
         }
+        let superblock_failures = metadata.superblock_write_failures();
+        if superblock_wedged(
+            superblock_failures,
+            self.superblock_wedged_fatal_failures.get(),
+        ) {
+            fatal(
+                FatalReason::SuperblockWedged,
+                &format!(
+                    "metadata superblock persist failed {superblock_failures} 
consecutive times, \
+                     past the [cluster] superblock_wedged_fatal_timeout 
window; exiting so a \
+                     supervisor handles the wedge instead of the replica 
limping fenced"
+                ),
+            );
+        }
 
         // Repair a lost primary self-ack: `RetransmitPrepares` to self is a
         // no-op, so the timer-driven retransmit above cannot recover the
@@ -9769,3 +9805,23 @@ mod control_frame_tests {
         assert!(body.is_empty());
     }
 }
+
+#[cfg(test)]
+mod superblock_fail_stop_tests {
+    //! The bound must stay disabled at 0: the simulator asserts a wedged
+    //! replica survives fenced in-process, and only the server arms it.
+
+    use super::superblock_wedged;
+
+    #[test]
+    fn zero_bound_never_fires() {
+        assert!(!superblock_wedged(u64::MAX, 0));
+    }
+
+    #[test]
+    fn bound_fires_at_and_past_the_threshold() {
+        assert!(!superblock_wedged(119, 120));
+        assert!(superblock_wedged(120, 120));
+        assert!(superblock_wedged(121, 120));
+    }
+}

Reply via email to