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));
+ }
+}