From 3e7ad6501368062c79f1b743a500ef97c55f508f Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Fri, 9 Oct 2026 04:58:41 +0800 Subject: [PATCH 1/3] feat(config): expose the startup readiness bounds under [tuning.startup] The metadata stall bound and the data-group recovery bound were constants in the bootstrap code, so an operator with a long WAL tail could not raise them. Both are newtypes now, read from [tuning.startup], which also stops a call site from transposing two adjacent Duration arguments. The shipped defaults change with them: the metadata stall bound moves from 30 s to 300 s and the data-group recovery bound from 60 s to 600 s. Those are the values production runs, and the backlog they cover takes minutes to replay, not seconds. The server config path already rejected such a key through `deny_unknown_fields`, but that message lists every valid field instead of the path that is read. A guard names the replacement path, so a misplaced key says where to put it. Both bounds are range-checked where the config is loaded, and the message names the key, the value and the range. Zero fails every boot on its first poll, and a bound near u64::MAX stops being a bound at all. The wait for the first authorization lease keeps its own constant: it is a hard deadline on one grant, with no progress reset, so it is not the metadata stall bound under another name. --- docs/getting-started.md | 22 +++ nodedb-test-support/src/single_node.rs | 19 ++- nodedb-types/src/config/tuning/config.rs | 19 +++ nodedb-types/src/config/tuning/mod.rs | 5 + nodedb-types/src/config/tuning/startup.rs | 143 ++++++++++++++++++ nodedb/src/bootstrap/cluster_ready.rs | 71 +++++++-- nodedb/src/bootstrap/data_group_recovery.rs | 22 ++- nodedb/src/config/server/config.rs | 115 ++++++++++++++ nodedb/src/config/server/domain.rs | 40 +++++ nodedb/src/control/cluster/test_one_node.rs | 10 +- .../src/control/security/auth_lease/status.rs | 6 + nodedb/src/main.rs | 2 + 12 files changed, 446 insertions(+), 28 deletions(-) create mode 100644 nodedb-types/src/config/tuning/startup.rs diff --git a/docs/getting-started.md b/docs/getting-started.md index 2ef35d74d..15f73f40a 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -589,6 +589,28 @@ a crashed node's descriptor leases block DDL until they expire. All three must be positive. `scope_expiry_interval_secs` has a floor of `10`. Below that the sweep costs more than the resolution it buys. +**Startup bounds:** + +Boot waits for the metadata raft group to apply its first entry, and for the +locally hosted data raft groups to replay their retained logs. Both waits are +bounded, and both are configurable. The defaults are sized for a cold restart of +a data dir that ran for weeks. The backlog takes minutes to apply, not seconds. + +| Config field | Default | +| ----------------------------------------------- | -------- | +| `tuning.startup.raft_ready_timeout_ms` | `300000` | +| `tuning.startup.data_group_recovery_timeout_ms` | `600000` | + +The metadata bound resets on every applied-index advance, so a large replay +finishes. Only a group that applies nothing fails. The data-group recovery bound +is a hard deadline. It does not reset, and a group still recovering when it +expires fails the boot. Both accept `1` to `86400000` ms (one day). A value +outside that range is rejected at load with the key, the value and the accepted +range named. A bound written under `[server]` is rejected with the path that is +read named instead. A key under `[tuning.startup]` that is not in the table is +rejected too. Two other boot waits keep fixed bounds and do not read this +section: Data Plane WAL replay (300 s) and the first authorization lease (30 s). + **Observability settings:** | Config field | Environment variable | Default | diff --git a/nodedb-test-support/src/single_node.rs b/nodedb-test-support/src/single_node.rs index 0934489a7..04534459b 100644 --- a/nodedb-test-support/src/single_node.rs +++ b/nodedb-test-support/src/single_node.rs @@ -117,9 +117,22 @@ pub async fn start( gateway_enable_gate: sequencer.register_gate(StartupPhase::GatewayEnable, "gateway"), }; // The harness cores replay their WAL before they report ready, so no - // replay receiver is owed here. - nodedb::bootstrap::cluster_ready::await_cluster_ready(shared, ready_rx, Vec::new(), gates) - .await?; + // replay receiver is owed here. The harness runs without a config file and + // pins short bounds, so a stuck group fails the test in seconds. The + // shipped defaults wait minutes for a replay backlog this harness never has. + let startup = nodedb_types::config::tuning::StartupTuning { + raft_ready_timeout_ms: 30_000, + data_group_recovery_timeout_ms: 60_000, + }; + nodedb::bootstrap::cluster_ready::await_cluster_ready( + shared, + ready_rx, + Vec::new(), + gates, + startup.raft_ready_timeout(), + startup.data_group_recovery_timeout(), + ) + .await?; Ok(raft) } diff --git a/nodedb-types/src/config/tuning/config.rs b/nodedb-types/src/config/tuning/config.rs index aac570e7f..a33741ca9 100644 --- a/nodedb-types/src/config/tuning/config.rs +++ b/nodedb-types/src/config/tuning/config.rs @@ -12,6 +12,7 @@ use super::memory::MemoryTuning; use super::network::{BridgeTuning, ClusterTransportTuning, NetworkTuning, WalTuning}; use super::scheduler::SchedulerTuning; use super::shutdown::ShutdownTuning; +use super::startup::StartupTuning; /// Top-level tuning configuration. /// @@ -51,6 +52,8 @@ pub struct TuningConfig { pub bitemporal: BitemporalTuning, #[serde(default)] pub maintenance: MaintenanceTuning, + #[serde(default)] + pub startup: StartupTuning, } impl TuningConfig { @@ -177,4 +180,20 @@ doc_cache_entries = 8192 assert_eq!(cfg.memory.overflow_max_bytes, 2 * 1024 * 1024 * 1024); assert_eq!(cfg.memory.doc_cache_entries, 8192); } + + #[test] + fn startup_bounds_are_read_from_the_aggregate() { + let cfg: TuningConfig = toml::from_str("").expect("deserialize"); + assert_eq!(cfg.startup.raft_ready_timeout_ms, 300_000); + assert_eq!(cfg.startup.data_group_recovery_timeout_ms, 600_000); + + let toml_str = r#" +[startup] +raft_ready_timeout_ms = 1234 +data_group_recovery_timeout_ms = 5678 +"#; + let cfg: TuningConfig = toml::from_str(toml_str).expect("deserialize"); + assert_eq!(cfg.startup.raft_ready_timeout_ms, 1234); + assert_eq!(cfg.startup.data_group_recovery_timeout_ms, 5678); + } } diff --git a/nodedb-types/src/config/tuning/mod.rs b/nodedb-types/src/config/tuning/mod.rs index d521e851f..8a9914dea 100644 --- a/nodedb-types/src/config/tuning/mod.rs +++ b/nodedb-types/src/config/tuning/mod.rs @@ -9,6 +9,7 @@ mod memory; mod network; mod scheduler; mod shutdown; +mod startup; pub use bitemporal::BitemporalTuning; pub use config::TuningConfig; @@ -23,3 +24,7 @@ pub use memory::MemoryTuning; pub use network::{BridgeTuning, ClusterTransportTuning, NetworkTuning, WalTuning}; pub use scheduler::SchedulerTuning; pub use shutdown::ShutdownTuning; +pub use startup::{ + DEFAULT_DATA_GROUP_RECOVERY_TIMEOUT_MS, DEFAULT_RAFT_READY_TIMEOUT_MS, + DataGroupRecoveryTimeout, RaftReadyTimeout, StartupTuning, +}; diff --git a/nodedb-types/src/config/tuning/startup.rs b/nodedb-types/src/config/tuning/startup.rs new file mode 100644 index 000000000..1969eeb75 --- /dev/null +++ b/nodedb-types/src/config/tuning/startup.rs @@ -0,0 +1,143 @@ +// SPDX-License-Identifier: Apache-2.0 + +//! Startup tuning: boot-time bounds applied by the readiness gates. + +use std::time::Duration; + +use serde::{Deserialize, Serialize}; + +/// Default bound for the metadata-group readiness stall. +/// +/// A node with a large metadata apply backlog (tens of thousands of entries +/// from a burst of cross-shard writes) needs minutes of replay before the +/// metadata group applies its first entry. A tighter bound turns a slow boot +/// into a restart loop. +pub const DEFAULT_RAFT_READY_TIMEOUT_MS: u64 = 300_000; + +/// Default bound for the local data raft groups' replay. +/// +/// Sized for a cold restart on a data dir that ran for weeks. The backlog of +/// committed entries takes minutes to apply, not seconds. +pub const DEFAULT_DATA_GROUP_RECOVERY_TIMEOUT_MS: u64 = 600_000; + +/// Metadata-group readiness stall bound, in the type system. +/// +/// Two adjacent `Duration` arguments transpose without a compile error. This +/// newtype pins the bound where it is chosen and handed to a gate; a helper +/// that takes plain `Duration` values converts at its own boundary. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct RaftReadyTimeout(pub Duration); + +/// Local data-group replay bound, in the type system. See [`RaftReadyTimeout`]. +/// +/// A separate type, not a shared one: the two bounds measure different waits +/// and are passed by different call sites, so making them interchangeable buys +/// nothing and costs a class of silent mix-ups. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct DataGroupRecoveryTimeout(pub Duration); + +impl From for RaftReadyTimeout { + fn from(value: Duration) -> Self { + Self(value) + } +} + +impl From for DataGroupRecoveryTimeout { + fn from(value: Duration) -> Self { + Self(value) + } +} + +fn default_raft_ready_timeout_ms() -> u64 { + DEFAULT_RAFT_READY_TIMEOUT_MS +} + +fn default_data_group_recovery_timeout_ms() -> u64 { + DEFAULT_DATA_GROUP_RECOVERY_TIMEOUT_MS +} + +/// Boot-time bounds for the readiness gates in `bootstrap::cluster_ready`. +/// +/// Both are range-checked where the config is loaded: below 1 ms the wait is +/// over before the group it waits for can answer, and above a day it stops +/// bounding anything while the instant it is added to is still finite. +/// +/// A key this section does not define is refused. A misspelt bound, such as +/// one without the `_ms` suffix, otherwise leaves the default in force and +/// fails a long recovery without a message that names the key. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct StartupTuning { + /// How long the metadata raft group can go without applying an entry + /// before the readiness gate fails startup. Every applied-index advance + /// resets the clock, so a large replay finishes. Only a stuck group fails. + /// Default: 300_000 (5 minutes). + #[serde(default = "default_raft_ready_timeout_ms")] + pub raft_ready_timeout_ms: u64, + + /// How long the locally hosted data raft groups have to replay their + /// retained logs before startup fails. Default: 600_000 (10 minutes). + #[serde(default = "default_data_group_recovery_timeout_ms")] + pub data_group_recovery_timeout_ms: u64, +} + +impl Default for StartupTuning { + fn default() -> Self { + Self { + raft_ready_timeout_ms: default_raft_ready_timeout_ms(), + data_group_recovery_timeout_ms: default_data_group_recovery_timeout_ms(), + } + } +} + +impl StartupTuning { + /// Metadata-group readiness stall bound. + pub fn raft_ready_timeout(&self) -> RaftReadyTimeout { + RaftReadyTimeout(Duration::from_millis(self.raft_ready_timeout_ms)) + } + + /// Data-group recovery bound. + pub fn data_group_recovery_timeout(&self) -> DataGroupRecoveryTimeout { + DataGroupRecoveryTimeout(Duration::from_millis(self.data_group_recovery_timeout_ms)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn defaults_are_the_backlogged_recovery_bounds() { + let cfg = StartupTuning::default(); + assert_eq!(cfg.raft_ready_timeout_ms, 300_000); + assert_eq!(cfg.data_group_recovery_timeout_ms, 600_000); + assert_eq!( + cfg.raft_ready_timeout(), + RaftReadyTimeout(Duration::from_secs(300)) + ); + assert_eq!( + cfg.data_group_recovery_timeout(), + DataGroupRecoveryTimeout(Duration::from_secs(600)) + ); + } + + #[test] + fn one_bound_set_leaves_the_other_at_its_default() { + let cfg: StartupTuning = + toml::from_str("raft_ready_timeout_ms = 1234").expect("deserialize"); + assert_eq!(cfg.raft_ready_timeout_ms, 1234); + assert_eq!(cfg.data_group_recovery_timeout_ms, 600_000); + } + + #[test] + fn an_unknown_key_is_refused() { + let error = toml::from_str::("data_group_recovery_timeout = 1800000") + .expect_err("a key without the _ms suffix is not a bound"); + let message = error.to_string(); + assert!(message.contains("data_group_recovery_timeout"), "{message}"); + assert!( + message.contains("data_group_recovery_timeout_ms"), + "the message must list the keys that exist: {message}" + ); + } +} diff --git a/nodedb/src/bootstrap/cluster_ready.rs b/nodedb/src/bootstrap/cluster_ready.rs index f35bfabf0..18ab1b314 100644 --- a/nodedb/src/bootstrap/cluster_ready.rs +++ b/nodedb/src/bootstrap/cluster_ready.rs @@ -7,6 +7,8 @@ use std::time::{Duration, Instant}; use tracing::info; +use nodedb_types::config::tuning::{DataGroupRecoveryTimeout, RaftReadyTimeout}; + use crate::bootstrap::schema_rehydrate::rehydrate_schema_registry; use crate::control::startup::ReadyGate; use crate::control::state::SharedState; @@ -30,6 +32,8 @@ pub async fn await_cluster_ready( mut raft_ready_rx: tokio::sync::watch::Receiver, data_plane_replay_done: Vec>>, gates: ClusterReadyGates, + raft_ready_timeout: RaftReadyTimeout, + data_group_recovery_timeout: DataGroupRecoveryTimeout, ) -> anyhow::Result<()> { let ClusterReadyGates { raft_gate, @@ -50,7 +54,7 @@ pub async fn await_cluster_ready( shared, &mut raft_ready_rx, &raft_gate, - RAFT_READY_STALL_TIMEOUT, + raft_ready_timeout, RAFT_READY_POLL_INTERVAL, ) .await?; @@ -156,7 +160,12 @@ pub async fn await_cluster_ready( // elections, so without this wait the gateway can open while a data // group's engines are still empty and an acknowledged write reads back as // if it never happened. Fail closed, like the replay wait above. - if let Err(e) = crate::bootstrap::data_group_recovery::await_data_group_recovery(shared).await { + if let Err(e) = crate::bootstrap::data_group_recovery::await_data_group_recovery( + shared, + data_group_recovery_timeout, + ) + .await + { data_groups_gate.fail(format!("data raft group recovery failed: {e}")); return Err(e); } @@ -272,7 +281,7 @@ pub async fn await_cluster_ready( Ok(_) => { crate::control::security::auth_lease::await_planning_admitted( shared, - RAFT_READY_STALL_TIMEOUT, + PLANNING_ADMISSION_TIMEOUT, ) .await } @@ -289,24 +298,28 @@ pub async fn await_cluster_ready( Ok(()) } -/// How long the metadata group may make NO replay progress before the boot -/// fails. Reset on every applied-index advance, so a large replay never trips -/// it — only a genuinely stuck group does. -const RAFT_READY_STALL_TIMEOUT: Duration = Duration::from_secs(30); - /// How often the stall check samples the applied index while waiting. const RAFT_READY_POLL_INTERVAL: Duration = Duration::from_secs(1); +/// How long boot waits for the first authorization lease before refusing to +/// open the gateway. +/// +/// Not the metadata group's stall bound, and not configurable with it: this is +/// a hard deadline on one grant, and it does not reset. The lease loop runs from +/// Raft start, so by the time boot reaches it every input of a grant is live; +/// 30 seconds is generous for a grant that is going to arrive at all. +const PLANNING_ADMISSION_TIMEOUT: Duration = Duration::from_secs(30); + /// Wait until the metadata raft group applies its first entry, and fail /// `raft_gate` when it does not. See [`raft_ready_progress`]. async fn wait_for_raft_ready( shared: &Arc, ready_rx: &mut tokio::sync::watch::Receiver, raft_gate: &ReadyGate, - stall_timeout: Duration, + stall_timeout: RaftReadyTimeout, poll_interval: Duration, ) -> anyhow::Result<()> { - match raft_ready_progress(shared, ready_rx, stall_timeout, poll_interval).await { + match raft_ready_progress(shared, ready_rx, stall_timeout.0, poll_interval).await { Ok(()) => { info!("metadata raft group ready — opening client listeners"); Ok(()) @@ -323,14 +336,20 @@ async fn wait_for_raft_ready( /// /// The same wait boot runs before it opens a listener. An in-process host /// that registers no startup gate calls this before it serves requests. +/// +/// `stall_timeout` is the caller's bound. Boot passes the configured +/// `[tuning.startup] raft_ready_timeout_ms` into [`await_cluster_ready`]; a +/// host with no config decides for itself rather than inheriting a value it +/// never read. pub async fn await_raft_ready( shared: &Arc, mut ready_rx: tokio::sync::watch::Receiver, + stall_timeout: RaftReadyTimeout, ) -> crate::Result<()> { raft_ready_progress( shared, &mut ready_rx, - RAFT_READY_STALL_TIMEOUT, + stall_timeout.0, RAFT_READY_POLL_INTERVAL, ) .await @@ -378,8 +397,7 @@ async fn raft_ready_progress( return Err(crate::Error::Internal { detail: format!( "metadata group applied no entry for {stall_timeout:?} \ - (applied_index stuck at {current}) — it failed to apply its \ - first entry" + (applied_index stuck at {current})" ), }); } @@ -433,7 +451,7 @@ mod tests { &shared, &mut ready_rx, &raft_gate, - Duration::from_millis(100), + RaftReadyTimeout(Duration::from_millis(100)), Duration::from_millis(10), ) .await; @@ -456,7 +474,7 @@ mod tests { &shared, &mut ready_rx, &raft_gate, - Duration::from_millis(100), + RaftReadyTimeout(Duration::from_millis(100)), Duration::from_millis(10), ) .await @@ -466,4 +484,27 @@ mod tests { "the failure must name the stall, not a generic timeout: {err}" ); } + + /// `await_raft_ready` reports the bound its caller passes, not a default. + /// + /// This pins the argument into the metadata gate only. It does not cover + /// the `await_cluster_ready` call in `main.rs`, where the two bounds arrive + /// as adjacent arguments of different types. + #[tokio::test(flavor = "multi_thread")] + async fn the_bound_a_caller_passes_is_the_bound_the_failure_reports() { + let shared = test_shared_state(); + let (_ready_tx, ready_rx) = tokio::sync::watch::channel(false); + + let err = await_raft_ready( + &shared, + ready_rx, + RaftReadyTimeout(Duration::from_millis(50)), + ) + .await + .expect_err("a group that applies nothing must fail"); + assert!( + err.to_string().contains("50ms"), + "the failure must name the caller's bound, not the shipped default: {err}" + ); + } } diff --git a/nodedb/src/bootstrap/data_group_recovery.rs b/nodedb/src/bootstrap/data_group_recovery.rs index 950a0e87e..430bd2d4d 100644 --- a/nodedb/src/bootstrap/data_group_recovery.rs +++ b/nodedb/src/bootstrap/data_group_recovery.rs @@ -23,6 +23,7 @@ use std::sync::Arc; use std::time::{Duration, Instant}; +use nodedb_types::config::tuning::DataGroupRecoveryTimeout; use tracing::{debug, info}; use crate::control::state::SharedState; @@ -31,11 +32,6 @@ use crate::control::state::SharedState; /// seconds, so a coarse poll costs nothing and avoids a busy loop. const POLL_INTERVAL: Duration = Duration::from_millis(50); -/// Upper bound on the whole wait. Generous relative to a randomized election -/// timeout plus replay of a retained log, but finite: a group that cannot elect -/// or cannot apply is a failure, not a reason to hang forever. -pub const DATA_GROUP_RECOVERY_TIMEOUT: Duration = Duration::from_secs(60); - /// True when `group_id` names a data group whose log carries user writes that /// must be replayed into the Data Plane before queries are served. /// @@ -170,13 +166,23 @@ fn pending_groups(statuses: Vec) -> Vec) -> anyhow::Result<()> { +pub async fn await_data_group_recovery( + shared: &Arc, + timeout: DataGroupRecoveryTimeout, +) -> anyhow::Result<()> { + // Typed at the call site, plain `Duration` inside, so the timeout message + // keeps naming the window the caller chose. + let timeout = timeout.0; let Some(status_fn) = shared.raft_status_fn.get() else { anyhow::bail!("data group recovery: no Raft status source: start_raft has not run"); }; let status_fn = Arc::clone(status_fn); - let deadline = Instant::now() + DATA_GROUP_RECOVERY_TIMEOUT; + let deadline = Instant::now() + timeout; loop { let pending = pending_groups(status_fn()); @@ -192,7 +198,7 @@ pub async fn await_data_group_recovery(shared: &Arc) -> anyhow::Res .collect::>() .join("; "); return Err(anyhow::anyhow!( - "data raft group recovery timeout after {DATA_GROUP_RECOVERY_TIMEOUT:?}: {detail}" + "data raft group recovery timeout after {timeout:?}: {detail}" )); } diff --git a/nodedb/src/config/server/config.rs b/nodedb/src/config/server/config.rs index 2656ba073..0e8e9897a 100644 --- a/nodedb/src/config/server/config.rs +++ b/nodedb/src/config/server/config.rs @@ -139,6 +139,30 @@ pub struct ServerConfig { pub scheduler: SchedulerConfig, } +/// Rejects a startup bound written at `[server]`, where nothing reads it, and +/// names the path that is read instead. `deny_unknown_fields` catches the key +/// anyway, but its message lists every valid field instead of the replacement. +fn reject_startup_bounds_at_the_server_path(content: &str) -> crate::Result<()> { + let Ok(doc) = toml::from_str::(content) else { + // A malformed document fails in the real parse below; nothing to + // inspect here. + return Ok(()); + }; + let Some(server) = doc.get("server").and_then(|v| v.as_table()) else { + return Ok(()); + }; + for key in ["raft_ready_timeout_ms", "data_group_recovery_timeout_ms"] { + if server.contains_key(key) { + return Err(crate::Error::Config { + detail: format!( + "`[server] {key}` is not read; set `[tuning.startup] {key}` instead" + ), + }); + } + } + Ok(()) +} + impl ServerConfig { /// Load configuration from a TOML file, falling back to defaults. pub fn from_file(path: &std::path::Path) -> crate::Result { @@ -150,6 +174,7 @@ impl ServerConfig { // that field — this is textual substitution before parsing, not a // second, competing expansion. let content = super::env_expand::expand_env(path, &content)?; + reject_startup_bounds_at_the_server_path(&content)?; let parsed: Self = toml::from_str(&content).map_err(|e| crate::Error::Config { detail: format!("invalid TOML config: {e}"), })?; @@ -477,4 +502,94 @@ mod tests { "unexpected error: {err}" ); } + + /// A startup bound written at `[server]` is rejected with the path that is + /// read named, not serde's full valid-field list. + #[test] + fn from_file_rejects_a_startup_bound_at_the_server_path() { + let path = write_temp_config( + "nodedb-server-path-startup-bound.toml", + "[server]\nraft_ready_timeout_ms = 300000\n", + ); + let err = ServerConfig::from_file(&path).unwrap_err(); + std::fs::remove_file(&path).ok(); + let msg = err.to_string(); + assert!(msg.contains("tuning.startup"), "{msg}"); + assert!(msg.contains("raft_ready_timeout_ms"), "{msg}"); + } + + /// The data-group bound is rejected at the same path. + #[test] + fn from_file_rejects_a_data_group_bound_at_the_server_path() { + let path = write_temp_config( + "nodedb-server-path-recovery-bound.toml", + "[server]\ndata_group_recovery_timeout_ms = 600000\n", + ); + let err = ServerConfig::from_file(&path).unwrap_err(); + std::fs::remove_file(&path).ok(); + let msg = err.to_string(); + assert!(msg.contains("tuning.startup"), "{msg}"); + assert!(msg.contains("data_group_recovery_timeout_ms"), "{msg}"); + } + + /// Both boot bounds live under `[tuning.startup]`; the other bound keeps + /// its default when only one is set. + #[test] + fn from_file_reads_the_startup_bounds_from_tuning() { + let path = write_temp_config( + "nodedb-startup-tuning.toml", + "[tuning.startup]\ndata_group_recovery_timeout_ms = 1234000\n", + ); + let cfg = ServerConfig::from_file(&path).expect("load config"); + std::fs::remove_file(&path).ok(); + assert_eq!(cfg.tuning.startup.data_group_recovery_timeout_ms, 1_234_000); + assert_eq!(cfg.tuning.startup.raft_ready_timeout_ms, 300_000); + } + + /// Neither bound is accepted outside the range a boot can wait for: zero + /// fails every boot on its first poll, `86400001` is past the stated ceiling, + /// and `u64::MAX` is not a bound at all. The message names the key, the value + /// and the range, so the operator does not have to read the source. + #[test] + fn from_file_rejects_a_startup_bound_out_of_range() { + for field in ["raft_ready_timeout_ms", "data_group_recovery_timeout_ms"] { + for value in ["0", "86400001", "18446744073709551615"] { + let path = write_temp_config( + "nodedb-startup-bound-out-of-range.toml", + &format!("[tuning.startup]\n{field} = {value}\n"), + ); + let err = ServerConfig::from_file(&path).unwrap_err(); + std::fs::remove_file(&path).ok(); + let msg = err.to_string(); + assert!( + msg.contains(field), + "the message must name the key for {field} = {value}: {msg}" + ); + assert!( + msg.contains(&format!( + "invalid value '{value}' for tuning.startup.{field}" + )), + "the message must name the value and the key for {field} = {value}: {msg}" + ); + assert!( + msg.contains("86400000"), + "the message must name the accepted range: {msg}" + ); + } + } + } + + /// Both ends of the accepted range load, and the value arrives intact. + #[test] + fn from_file_accepts_the_ends_of_the_startup_bound_range() { + for value in [1u64, 86_400_000] { + let path = write_temp_config( + "nodedb-startup-bound-in-range.toml", + &format!("[tuning.startup]\nraft_ready_timeout_ms = {value}\n"), + ); + let cfg = ServerConfig::from_file(&path).expect("a bound in range loads"); + std::fs::remove_file(&path).ok(); + assert_eq!(cfg.tuning.startup.raft_ready_timeout_ms, value); + } + } } diff --git a/nodedb/src/config/server/domain.rs b/nodedb/src/config/server/domain.rs index bb2cd2950..5ecf37dac 100644 --- a/nodedb/src/config/server/domain.rs +++ b/nodedb/src/config/server/domain.rs @@ -18,6 +18,20 @@ pub(super) const MIN_WAL_WRITE_BUFFER_BYTES: usize = 64 * 1024; /// shorter sweep costs more than the resolution it buys. pub(super) const MIN_SCOPE_EXPIRY_SECS: u64 = 10; +/// Smallest boot bound a caller can set, in milliseconds. +/// +/// A zero bound fails every boot on its first poll: the wait is over before the +/// group it waits for can answer. +pub(super) const MIN_STARTUP_BOUND_MS: u64 = 1; + +/// Largest boot bound, in milliseconds: one day. +/// +/// A bound past this stops being a bound: a stuck group and a healthy one are +/// indistinguishable for the whole window. On Linux a bound near `u64::MAX` +/// still computes, so the wait becomes effectively unbounded, and on a platform +/// whose `Instant` is a `u64` nanosecond counter the addition overflows. +pub(super) const MAX_STARTUP_BOUND_MS: u64 = 86_400_000; + /// Rejects an endpoint that carries no `http://` or `https://` host. pub(super) fn otlp_endpoint_has_host(raw: &str) -> bool { raw.strip_prefix("http://") @@ -38,6 +52,22 @@ pub(super) fn positive_u64(value: u64, field: &str) -> crate::Result<()> { Ok(()) } +/// A boot bound, in milliseconds, inside the range a boot can wait for. +/// +/// Zero fails every boot on its first poll: the wait is over before the group it +/// waits for can answer. A bound near `u64::MAX` is not a bound at all, and the +/// server refuses it here instead of letting a stuck boot look healthy. +fn startup_bound(value: u64, field: &str) -> crate::Result<()> { + if !(MIN_STARTUP_BOUND_MS..=MAX_STARTUP_BOUND_MS).contains(&value) { + return Err(reject( + field, + value, + &format!("a bound from {MIN_STARTUP_BOUND_MS} to {MAX_STARTUP_BOUND_MS} ms (one day)"), + )); + } + Ok(()) +} + fn positive_usize(value: usize, field: &str) -> crate::Result<()> { if value == 0 { return Err(reject(field, value, "a positive integer")); @@ -72,6 +102,16 @@ pub(super) fn validate_domain(config: &ServerConfig) -> crate::Result<()> { super::pitr::validate_pitr(config)?; super::backup::validate_backup(config)?; + let s = &config.tuning.startup; + startup_bound( + s.raft_ready_timeout_ms, + "tuning.startup.raft_ready_timeout_ms", + )?; + startup_bound( + s.data_group_recovery_timeout_ms, + "tuning.startup.data_group_recovery_timeout_ms", + )?; + let ts = &config.tuning.timeseries; positive_usize( ts.memtable_budget_bytes, diff --git a/nodedb/src/control/cluster/test_one_node.rs b/nodedb/src/control/cluster/test_one_node.rs index 7e97b2ee7..003bf77ff 100644 --- a/nodedb/src/control/cluster/test_one_node.rs +++ b/nodedb/src/control/cluster/test_one_node.rs @@ -14,13 +14,19 @@ use std::sync::Arc; use std::time::Duration; -use nodedb_types::config::tuning::ClusterTransportTuning; +use nodedb_types::config::tuning::{ClusterTransportTuning, RaftReadyTimeout}; use crate::bridge::dispatch::{CoreChannelDataSide, Dispatcher}; use crate::control::security::credential::CredentialStore; use crate::control::state::SharedState; use crate::wal::WalManager; +/// The metadata-group stall bound for this harness: fail fast instead of +/// inheriting the backlog-tolerant production default. The harness boots a +/// bare single-voter cluster, so a stall here is a test bug, not a slow +/// replay, and a short bound reports it sooner. +const TEST_RAFT_READY_TIMEOUT: RaftReadyTimeout = RaftReadyTimeout(Duration::from_secs(30)); + /// How long each shutdown step can take. const SHUTDOWN_STEP: Duration = Duration::from_secs(5); @@ -104,7 +110,7 @@ pub(crate) async fn boot_with_core( .lock() .unwrap_or_else(|p| p.into_inner()) .take(); - crate::bootstrap::cluster_ready::await_raft_ready(&state, ready) + crate::bootstrap::cluster_ready::await_raft_ready(&state, ready, TEST_RAFT_READY_TIMEOUT) .await .expect("the metadata group applies its first entry"); diff --git a/nodedb/src/control/security/auth_lease/status.rs b/nodedb/src/control/security/auth_lease/status.rs index 00727d625..620076982 100644 --- a/nodedb/src/control/security/auth_lease/status.rs +++ b/nodedb/src/control/security/auth_lease/status.rs @@ -82,6 +82,12 @@ pub(crate) fn lease_status_of( /// Wait until this node can plan permission-checked statements, or refuse /// once `timeout` passes. +/// +/// `timeout` is a hard deadline on the first lease grant, not the metadata +/// group's stall bound. The two waits end differently: the stall bound resets +/// every time the group applies an entry, and this one never does, so a caller +/// that has no reason of its own passes its own constant rather than the one +/// the group uses. pub async fn await_planning_admitted(state: &SharedState, timeout: Duration) -> crate::Result<()> { if planning_admitted_within(state, timeout).await { return Ok(()); diff --git a/nodedb/src/main.rs b/nodedb/src/main.rs index 5d90d43eb..4251004ad 100644 --- a/nodedb/src/main.rs +++ b/nodedb/src/main.rs @@ -291,6 +291,8 @@ async fn server_main() -> anyhow::Result<()> { health_loop_gate, gateway_enable_gate, }, + config.tuning.startup.raft_ready_timeout(), + config.tuning.startup.data_group_recovery_timeout(), ) .await?; From a8fdd9ead61f3bf1e39339f07c337ea4845b2399 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Fri, 9 Oct 2026 04:58:41 +0800 Subject: [PATCH 2/3] fix(wal): name the WAL format version when startup refuses a store A store written by another build failed every boot with a message calling its segments corrupt. The reader knew the format version and flattened it, along with every other header validation failure, into StopReason::Corruption, so the operator was sent looking for damage that was not there. The version error now travels out of the three readers that decide a record header. It is judged once per segment, at its first record: one writer never mixes formats inside one segment, so a mismatch there is another build's format, and a mismatch anywhere later is damage that torn_tail classifies. A zeroed version is neither, because no build writes format zero, so it stays a torn-write stop. torn_tail is deliberately not one of the three. A version mismatch past a corruption stop is damage there, never a resync point, and its own test pins that. The callers that know the segment path attach it, so the refusal names the file and both versions. Startup validation stays a single pass after the open: the open recovers the newest segment and fails there with the path attached, which is the same refusal without reading every segment twice. --- nodedb-wal/src/error.rs | 36 +++ nodedb-wal/src/lazy_reader.rs | 100 ++++++++- nodedb-wal/src/mmap_reader/reader.rs | 96 +++++++- nodedb-wal/src/mmap_reader/replay.rs | 10 +- nodedb-wal/src/reader.rs | 141 +++++++++++- nodedb-wal/src/recovery.rs | 60 ++++- nodedb-wal/src/segmented.rs | 40 +++- nodedb-wal/src/torn_tail.rs | 48 ++++ nodedb/src/bootstrap/wal_init.rs | 21 +- nodedb/src/ctl/restore/cut_rule.rs | 29 ++- nodedb/src/ctl/restore/segment.rs | 122 +++++++++- nodedb/src/wal/manager/replay.rs | 322 ++++++++++++++++++++++++++- 12 files changed, 994 insertions(+), 31 deletions(-) diff --git a/nodedb-wal/src/error.rs b/nodedb-wal/src/error.rs index 58c38dd1e..b092e9eae 100644 --- a/nodedb-wal/src/error.rs +++ b/nodedb-wal/src/error.rs @@ -31,6 +31,22 @@ pub enum WalError { #[error("unsupported WAL format version {version} (supported: {supported})")] UnsupportedVersion { version: u16, supported: u16 }, + /// A segment holds records in a WAL format version this build cannot read. + /// + /// The same finding as [`WalError::UnsupportedVersion`], reported by the + /// callers that know which file carries it. A reader sees one header at a + /// time and holds no path, and a version alone does not tell an operator + /// which segment to act on. + #[error( + "WAL segment '{path}' holds records in WAL format version {version}; this build reads \ + version {supported}" + )] + SegmentFormatVersion { + path: String, + version: u16, + supported: u16, + }, + /// Unknown required record type encountered during replay. /// Optional unknown record types are safely skipped. #[error("unknown required record type {record_type} at LSN {lsn}")] @@ -214,4 +230,24 @@ pub enum WalError { }, } +impl WalError { + /// Name the segment a version gap was found in. + /// + /// A reader reports [`WalError::UnsupportedVersion`] from header validation + /// alone and holds no path. The callers that opened the file do, and only + /// they can turn the finding into something an operator can act on. Every + /// other error passes through unchanged. + #[cfg(not(target_arch = "wasm32"))] + pub(crate) fn with_segment_path(self, path: &std::path::Path) -> Self { + match self { + Self::UnsupportedVersion { version, supported } => Self::SegmentFormatVersion { + path: path.display().to_string(), + version, + supported, + }, + other => other, + } + } +} + pub type Result = std::result::Result; diff --git a/nodedb-wal/src/lazy_reader.rs b/nodedb-wal/src/lazy_reader.rs index 3fb315548..071d0c60e 100644 --- a/nodedb-wal/src/lazy_reader.rs +++ b/nodedb-wal/src/lazy_reader.rs @@ -55,6 +55,10 @@ fn checked_offset_add(offset: u64, len: u64) -> Result { pub struct LazyWalReader { file: File, offset: u64, + /// Offset of the segment's first record, past any preamble. The format + /// version is judged here and nowhere else: one writer never mixes formats + /// inside one segment. + records_start: u64, /// Preamble read from offset 0 of this segment (present when encryption is /// active). The epoch is part of the AAD used to decrypt payloads. segment_preamble: Option, @@ -94,6 +98,7 @@ impl LazyWalReader { Ok(Self { file, offset: start_offset, + records_start: start_offset, segment_preamble, double_write, stop_reason: None, @@ -123,9 +128,10 @@ impl LazyWalReader { /// Read the next record header without reading the payload. /// - /// Returns `None` at EOF or first corruption. After this call, use - /// either `read_payload()` to get the payload or `skip_payload()` to - /// seek past it. + /// Returns `None` at EOF or first corruption. Returns `Err` for a format + /// version this build cannot read at the segment's first record, and the + /// caller must stop at that `Err`. After a header, use either + /// `read_payload()` to get the payload or `skip_payload()` to seek past it. /// /// Alignment padding records are consumed internally — the caller only /// ever sees real records. @@ -153,6 +159,15 @@ impl LazyWalReader { match header.validate(record_offset) { Ok(()) => {} Err(error @ WalError::PayloadTooLarge { .. }) => return Err(error), + // Format is decided at the segment's first record; a mismatch + // later in the segment is damage, and `torn_tail` classifies + // the stop that records it. A zeroed version is a torn write at + // either offset, not a format any build wrote. + Err(error @ WalError::UnsupportedVersion { version, .. }) + if version != 0 && record_offset == self.records_start => + { + return Err(error); + } Err(_) => { return self.stop(StopReason::Corruption { offset: record_offset, @@ -345,7 +360,10 @@ where { let mut reader = LazyWalReader::open(path, keys)?; let mut last_lsn = 0u64; - while let Some(header) = reader.next_header()? { + while let Some(header) = reader + .next_header() + .map_err(|error| error.with_segment_path(path))? + { last_lsn = header.lsn; handler(&mut reader, &header)?; } @@ -541,4 +559,78 @@ mod tests { let mut reader = LazyWalReader::open(&path, None).unwrap(); assert!(reader.next_header().unwrap().is_none()); } + + /// A version this build cannot read at the segment's first record is a + /// typed error that names the version. + #[test] + fn an_unknown_format_version_at_the_first_record_is_an_error() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("unknown.wal"); + { + let mut writer = WalWriter::open_without_direct_io(&path).unwrap(); + writer + .append(RecordType::Put as u32, 1, 0, 0, b"x") + .unwrap(); + writer.sync().unwrap(); + } + + // The version field is bytes 4..6 of the record header. + let mut bytes = std::fs::read(&path).unwrap(); + bytes[4..6].copy_from_slice(&1u16.to_le_bytes()); + std::fs::write(&path, &bytes).unwrap(); + + let mut reader = LazyWalReader::open(&path, None).unwrap(); + match reader.next_header() { + Err(WalError::UnsupportedVersion { version, supported }) => { + assert_eq!(version, 1); + assert_eq!(supported, crate::record::WAL_FORMAT_VERSION); + } + other => panic!("expected UnsupportedVersion, got {other:?}"), + } + } + + /// A version this build cannot read, found after an intact record, is + /// damage: the stream stops there as corruption, not as a format gap. + #[test] + fn a_version_mismatch_after_the_first_record_is_damage() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("two-records.wal"); + { + let mut writer = WalWriter::open_without_direct_io(&path).unwrap(); + writer + .append(RecordType::Put as u32, 1, 0, 0, b"first") + .unwrap(); + writer + .append(RecordType::Put as u32, 1, 0, 0, b"second") + .unwrap(); + writer.sync().unwrap(); + } + + let mut reader = LazyWalReader::open(&path, None).unwrap(); + let first = reader.next_header().unwrap().expect("the first header"); + reader.skip_payload(&first).unwrap(); + let second_at = reader.offset(); + + let mut bytes = std::fs::read(&path).unwrap(); + let version_at = usize::try_from(second_at).unwrap() + 4; + bytes[version_at..version_at + 2].copy_from_slice(&1u16.to_le_bytes()); + std::fs::write(&path, &bytes).unwrap(); + + let mut reader = LazyWalReader::open(&path, None).unwrap(); + let first = reader + .next_header() + .unwrap() + .expect("the first record is still readable"); + assert_eq!(first.lsn, 1); + reader.skip_payload(&first).unwrap(); + assert!( + reader.next_header().unwrap().is_none(), + "the mismatch stops the stream instead of failing the read" + ); + assert_eq!( + reader.stop_reason(), + Some(StopReason::Corruption { offset: second_at }), + "a mismatch after the first record is damage, not a format gap" + ); + } } diff --git a/nodedb-wal/src/mmap_reader/reader.rs b/nodedb-wal/src/mmap_reader/reader.rs index d60b43f57..fca6ee9c5 100644 --- a/nodedb-wal/src/mmap_reader/reader.rs +++ b/nodedb-wal/src/mmap_reader/reader.rs @@ -113,6 +113,10 @@ fn fadv_dontneed(fd: &std::fs::File, len: usize, path: &Path) { pub struct MmapWalReader { mmap: Mmap, offset: usize, + /// Offset of the segment's first record, past any preamble. The format + /// version is judged here and nowhere else: one writer never mixes formats + /// inside one segment. + records_start: usize, /// Preamble mapped at offset 0 of this segment (present when encryption is /// active). The epoch is part of the AAD used to decrypt payloads. segment_preamble: Option, @@ -194,6 +198,7 @@ impl MmapWalReader { Ok(Self { mmap, offset, + records_start: offset, segment_preamble, file, path: path.to_path_buf(), @@ -245,7 +250,9 @@ impl MmapWalReader { /// Read the next record from the mmap'd region. /// - /// Returns `None` at EOF or at the first corruption point. + /// Returns `None` at EOF or at the first corruption point. Returns `Err` + /// for a format version this build cannot read at the segment's first + /// record, and the caller must stop at that `Err`. /// Header parsing avoids extra copies, but the payload is copied out of /// the mmap'd region into an owned `Vec` on `WalRecord`: records must /// outlive the reader (they cross thread boundaries in parallel replay @@ -299,6 +306,15 @@ impl MmapWalReader { match header.validate(header_offset) { Ok(()) => {} Err(error @ WalError::PayloadTooLarge { .. }) => return Err(error), + // Format is decided at the segment's first record; a mismatch + // later in the segment is damage, and `torn_tail` classifies + // the stop that records it. A zeroed version is a torn write at + // either offset, not a format any build wrote. + Err(error @ WalError::UnsupportedVersion { version, .. }) + if version != 0 && self.offset == self.records_start => + { + return Err(error); + } Err(_) => { return self.stop(StopReason::Corruption { offset: header_offset, @@ -517,6 +533,84 @@ mod tests { )); } + /// A version this build cannot read at the segment's first record is a + /// typed error that names the version. + #[test] + fn an_unknown_format_version_at_the_first_record_is_an_error() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("unknown.wal"); + { + let mut writer = test_writer(&path); + writer + .append(RecordType::Put as u32, 1, 0, 0, b"x") + .unwrap(); + writer.sync().unwrap(); + } + + // The version field is bytes 4..6 of the record header. + let mut bytes = std::fs::read(&path).unwrap(); + bytes[4..6].copy_from_slice(&1u16.to_le_bytes()); + std::fs::write(&path, &bytes).unwrap(); + + let mut reader = MmapWalReader::open(&path, None).unwrap(); + match reader.next_record() { + Err(WalError::UnsupportedVersion { version, supported }) => { + assert_eq!(version, 1); + assert_eq!(supported, crate::record::WAL_FORMAT_VERSION); + } + other => panic!("expected UnsupportedVersion, got {other:?}"), + } + } + + /// A version this build cannot read, found after an intact record, is + /// damage: the stream stops there as corruption, not as a format gap. + #[test] + fn a_version_mismatch_after_the_first_record_is_damage() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("two-records.wal"); + { + let mut writer = test_writer(&path); + writer + .append(RecordType::Put as u32, 1, 0, 0, b"first") + .unwrap(); + writer + .append(RecordType::Put as u32, 1, 0, 0, b"second") + .unwrap(); + writer.sync().unwrap(); + } + + let second_at = { + let mut reader = MmapWalReader::open(&path, None).unwrap(); + reader + .next_record() + .unwrap() + .expect("the first record reads"); + reader.offset() + }; + + let mut bytes = std::fs::read(&path).unwrap(); + let version_at = second_at + 4; + bytes[version_at..version_at + 2].copy_from_slice(&1u16.to_le_bytes()); + std::fs::write(&path, &bytes).unwrap(); + + let mut reader = MmapWalReader::open(&path, None).unwrap(); + let first = reader + .next_record() + .unwrap() + .expect("the first record is still readable"); + assert_eq!(first.header.lsn, 1); + assert!( + reader.next_record().unwrap().is_none(), + "the mismatch stops the stream instead of failing the read" + ); + let second_at = u64::try_from(second_at).unwrap(); + assert_eq!( + reader.stop_reason(), + Some(StopReason::Corruption { offset: second_at }), + "a mismatch after the first record is damage, not a format gap" + ); + } + #[test] fn checked_range_end_rejects_overflow_and_short_ranges() { assert!(checked_range_end(usize::MAX, 1, usize::MAX).is_err()); diff --git a/nodedb-wal/src/mmap_reader/replay.rs b/nodedb-wal/src/mmap_reader/replay.rs index 7c4e172fb..4d82e16c2 100644 --- a/nodedb-wal/src/mmap_reader/replay.rs +++ b/nodedb-wal/src/mmap_reader/replay.rs @@ -103,7 +103,10 @@ fn scan_segment( let mut reader = MmapWalReader::open(&segment.path, keys)?; let mut records = Vec::new(); let mut last_lsn = 0u64; - while let Some(record) = reader.next_record()? { + while let Some(record) = reader + .next_record() + .map_err(|error| error.with_segment_path(&segment.path))? + { last_lsn = record.header.lsn; if record.header.lsn >= from_lsn { records.push(record); @@ -211,7 +214,10 @@ pub fn replay_segments_mmap_limit( continuity.check(seg)?; let mut reader = MmapWalReader::open(&seg.path, keys)?; let mut last_lsn = 0u64; - while let Some(record) = reader.next_record()? { + while let Some(record) = reader + .next_record() + .map_err(|error| error.with_segment_path(&seg.path))? + { last_lsn = record.header.lsn; if record.header.lsn >= from_lsn { records.push(record); diff --git a/nodedb-wal/src/reader.rs b/nodedb-wal/src/reader.rs index ce46141ab..676f1da80 100644 --- a/nodedb-wal/src/reader.rs +++ b/nodedb-wal/src/reader.rs @@ -58,8 +58,9 @@ pub enum StopReason { /// The segment ended cleanly on a record boundary. Eof, - /// A record at this offset failed structural validation (bad magic, - /// unsupported version, short read, or checksum mismatch). + /// A record at this offset failed structural validation (bad magic, a + /// zeroed format version, a short read, a checksum mismatch, or a format + /// version this build cannot read found after the segment's first record). Corruption { offset: u64 }, } @@ -67,6 +68,10 @@ pub enum StopReason { pub struct WalReader { file: File, offset: u64, + /// Offset of the segment's first record, past any preamble. The format + /// version is judged here and nowhere else: one writer never mixes formats + /// inside one segment. + records_start: u64, /// Preamble read from offset 0 of this segment (present when encryption /// is active). The epoch is used as part of the AAD for decryption. segment_preamble: Option, @@ -133,6 +138,7 @@ impl WalReader { Ok(Self { file, offset: start_offset, + records_start: start_offset, segment_preamble, double_write, stop_reason: None, @@ -185,7 +191,10 @@ impl WalReader { /// Read the next record from the WAL. /// /// Returns `None` at EOF (clean end) or at the first corruption point. - /// Returns `Err` only for I/O errors or unknown required record types. + /// Returns `Err` for I/O errors, unknown required record types, an + /// oversized payload declaration, and a format version this build cannot + /// read at the segment's first record. A reader that returns `Err` for a + /// header has consumed that header, so a caller must stop at the `Err`. pub fn next_record(&mut self) -> Result> { loop { // Read header. @@ -230,6 +239,22 @@ impl WalReader { match header.validate(header_offset) { Ok(()) => {} Err(error @ WalError::PayloadTooLarge { .. }) => return Err(error), + // The format is judged once per segment, at its first record. A + // mismatch there is another build's format, and the version has + // to travel out so the refusal can name it. A mismatch anywhere + // later cannot be a format gap, because a writer never mixes + // formats inside one segment. It is damage, and this reader + // stops at it for `torn_tail` to classify. + // + // A zeroed version at the first record is neither: no build + // writes format 0, so those are uninitialised bytes from a torn + // write, and the writer paths that share this reader would + // refuse a store whose tail is merely torn. It stays a stop. + Err(error @ WalError::UnsupportedVersion { version, .. }) + if version != 0 && header_offset == self.records_start => + { + return Err(error); + } Err(_) => { return self.stop(StopReason::Corruption { offset: header_offset, @@ -525,4 +550,114 @@ mod tests { assert_eq!(records.len(), 1); assert_eq!(records[0].payload, b"keep-me"); } + + /// An unknown format version is a typed error, not a silent stop, and a + /// zeroed version is a stop rather than an error. + /// + /// A store written by another build is intact, so its format reaches the + /// caller as a version and never as damage. Version zero is the exception: + /// no build writes it, so it stays a corruption stop, which `torn_tail` + /// already bounds. + #[test] + fn an_unknown_format_version_is_an_error_and_a_zeroed_one_is_a_stop() { + let dir = tempfile::tempdir().unwrap(); + let framed = dir.path().join("framed.wal"); + { + let mut writer = WalWriter::open_without_direct_io(&framed).unwrap(); + writer + .append(RecordType::Put as u32, 1, 0, 0, b"x") + .unwrap(); + writer.sync().unwrap(); + } + + // The version field is bytes 4..6 of the record header. + let mut bytes = std::fs::read(&framed).unwrap(); + bytes[4..6].copy_from_slice(&1u16.to_le_bytes()); + let unknown = dir.path().join("unknown.wal"); + std::fs::write(&unknown, &bytes).unwrap(); + + let reader = WalReader::open_raw(&unknown).unwrap(); + match reader.records().collect::>>() { + Err(WalError::UnsupportedVersion { version, supported }) => { + assert_eq!(version, 1); + assert_eq!(supported, crate::record::WAL_FORMAT_VERSION); + } + other => panic!("expected UnsupportedVersion, got {other:?}"), + } + + bytes[4..6].copy_from_slice(&0u16.to_le_bytes()); + let zeroed = dir.path().join("zeroed.wal"); + std::fs::write(&zeroed, &bytes).unwrap(); + + let mut reader = WalReader::open_raw(&zeroed).unwrap(); + let mut records = Vec::new(); + while let Some(record) = reader + .next_record() + .expect("a zeroed version stops the stream, it does not fail it") + { + records.push(record); + } + assert!( + records.is_empty(), + "nothing before the stop is a readable record" + ); + // An empty stream is also what a silent skip or a clean end would + // return, so the reason has to be asserted for this to prove a stop. + assert!( + matches!(reader.stop_reason(), Some(StopReason::Corruption { .. })), + "a zeroed version must be a corruption stop, not a quiet end: {:?}", + reader.stop_reason() + ); + } + + /// A version this build cannot read, found after the segment's first + /// record, is damage rather than a format gap. + /// + /// One writer never mixes formats inside one segment, so a mismatch there + /// can only be a broken header. The stream stops and `torn_tail` classifies + /// where, instead of the reader declaring a format gap the store does not + /// have. + #[test] + fn a_version_mismatch_after_the_first_record_is_damage() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("two-records.wal"); + { + let mut writer = WalWriter::open_without_direct_io(&path).unwrap(); + writer + .append(RecordType::Put as u32, 1, 0, 0, b"first") + .unwrap(); + writer + .append(RecordType::Put as u32, 1, 0, 0, b"second") + .unwrap(); + writer.sync().unwrap(); + } + + let mut reader = WalReader::open_raw(&path).unwrap(); + let first = reader + .next_record() + .unwrap() + .expect("the first record reads"); + assert_eq!(first.header.lsn, 1, "records before the damage still read"); + let second_at = reader.committed_end(); + + let mut bytes = std::fs::read(&path).unwrap(); + let version_at = usize::try_from(second_at).unwrap() + 4; + bytes[version_at..version_at + 2].copy_from_slice(&1u16.to_le_bytes()); + std::fs::write(&path, &bytes).unwrap(); + + let mut reader = WalReader::open_raw(&path).unwrap(); + assert!( + reader.next_record().unwrap().is_some(), + "the first record is still readable" + ); + assert!( + reader.next_record().unwrap().is_none(), + "the mismatch stops the stream instead of failing the read" + ); + assert!( + matches!(reader.stop_reason(), Some(StopReason::Corruption { .. })), + "a mismatch after the first record is damage, not a format gap: {:?}", + reader.stop_reason() + ); + } } diff --git a/nodedb-wal/src/recovery.rs b/nodedb-wal/src/recovery.rs index 59ffa96df..ec6d91c15 100644 --- a/nodedb-wal/src/recovery.rs +++ b/nodedb-wal/src/recovery.rs @@ -22,7 +22,7 @@ use std::path::Path; -use crate::error::{Result, WalError}; +use crate::error::Result; use crate::reader::WalReader; /// Result of scanning a WAL file for recovery. @@ -76,11 +76,11 @@ pub fn recover(path: &Path) -> Result { // End of committed prefix (EOF or corruption). break; } - Err(e @ WalError::UnknownRequiredRecordType { .. }) => { - // Cannot proceed past unknown required records. - return Err(e); - } - Err(e) => return Err(e), + // Nothing past an unknown required record can be read, and a version + // gap is the one error that needs the segment it was found in: the + // reader sees one header at a time and holds no path. Every other + // error passes through unchanged. + Err(error) => return Err(error.with_segment_path(path)), } } @@ -101,6 +101,7 @@ pub fn recover(path: &Path) -> Result { #[cfg(test)] mod tests { use super::*; + use crate::error::WalError; use crate::record::RecordType; use crate::writer::WalWriter; @@ -188,4 +189,51 @@ mod tests { assert_eq!(info.record_count, 2); assert_eq!(info.next_lsn(), 3); } + + /// The version field is bytes 4..6 of a record header. Where each of + /// `count` records starts, and the file's bytes. + fn segment_with_records(path: &Path, count: u8) -> (Vec, Vec) { + { + let mut writer = WalWriter::open_without_direct_io(path).unwrap(); + for i in 0..count { + writer + .append(RecordType::Put as u32, 1, 0, 0, &[i; 8]) + .unwrap(); + } + writer.sync().unwrap(); + } + let mut reader = WalReader::open_raw(path).unwrap(); + let mut starts = Vec::new(); + loop { + let start = reader.committed_end(); + if reader.next_record().unwrap().is_none() { + break; + } + starts.push(start); + } + (starts, std::fs::read(path).unwrap()) + } + + #[test] + fn a_lone_first_record_in_another_format_is_a_format_gap() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("lone.wal"); + let (starts, mut bytes) = segment_with_records(&path, 1); + let at = usize::try_from(starts[0]).unwrap() + 4; + bytes[at..at + 2].copy_from_slice(&1u16.to_le_bytes()); + std::fs::write(&path, &bytes).unwrap(); + + match recover(&path) { + Err(WalError::SegmentFormatVersion { + path: named, + version, + supported, + }) => { + assert_eq!(named, path.display().to_string()); + assert_eq!(version, 1); + assert_eq!(supported, crate::record::WAL_FORMAT_VERSION); + } + other => panic!("expected SegmentFormatVersion, got {other:?}"), + } + } } diff --git a/nodedb-wal/src/segmented.rs b/nodedb-wal/src/segmented.rs index f4754c63b..9211f70fe 100644 --- a/nodedb-wal/src/segmented.rs +++ b/nodedb-wal/src/segmented.rs @@ -463,7 +463,10 @@ pub fn replay_from_limit_dir( // the `from_lsn` filter — continuity is a property of the log, not of // the caller's window into it. let mut last_lsn = 0u64; - while let Some(record) = reader.next_record()? { + while let Some(record) = reader + .next_record() + .map_err(|error| error.with_segment_path(&seg.path))? + { last_lsn = record.header.lsn; if record.header.lsn >= from_lsn { records.push(record); @@ -1169,4 +1172,39 @@ mod tests { Ok(_) => panic!("opened a WAL whose last segment reissues LSNs"), } } + + /// A newest segment whose first record has the record magic and a zeroed + /// version is a torn write, not a format this build cannot read. Opening + /// over it must succeed, and the writer must keep working, because the + /// writer's open is the first thing a boot does. + #[test] + fn open_accepts_a_newest_segment_with_a_zeroed_version() { + let dir = tempfile::tempdir().unwrap(); + { + let mut wal = SegmentedWal::open(test_config(dir.path())).unwrap(); + wal.append(RecordType::Put as u32, 1, 0, 0, b"a").unwrap(); + wal.sync().unwrap(); + } + + // Written after the close, so the open has to read this state back. + let newest = discover_segments(dir.path()) + .unwrap() + .last() + .unwrap() + .first_lsn; + let mut bytes = vec![0u8; 64 * 1024]; + bytes[0..4].copy_from_slice(&crate::record::WAL_MAGIC.to_le_bytes()); + std::fs::write(segment_path(dir.path(), newest + 1_000), &bytes).unwrap(); + + let mut wal = SegmentedWal::open(test_config(dir.path())) + .expect("a torn newest segment is not a format this build cannot read"); + let lsn = wal.append(RecordType::Put as u32, 1, 0, 0, b"b").unwrap(); + // The torn segment holds no record, so its filename's first LSN is the + // floor the resumed sequence starts at. + assert_eq!( + lsn, + newest + 1_000, + "the writer has to keep working over the torn tail, not just open" + ); + } } diff --git a/nodedb-wal/src/torn_tail.rs b/nodedb-wal/src/torn_tail.rs index 39d5e9773..79b0affa9 100644 --- a/nodedb-wal/src/torn_tail.rs +++ b/nodedb-wal/src/torn_tail.rs @@ -139,6 +139,10 @@ pub fn classify(path: &Path, corruption_offset: u64, last_lsn: u64) -> Result Result> { let header_end = match offset.checked_add(HEADER_SIZE as u64) { Some(end) if end <= file_len => end, @@ -303,6 +307,50 @@ mod tests { ); } + /// A record whose format version this build cannot read is damage, not a + /// resync point: the scan passes over it to the intact record behind it. + #[test] + fn a_version_mismatch_past_the_stop_is_not_a_record() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("version_mismatch.wal"); + write_segment(&path, 4); + + // Where each record starts: the reader's committed end after the one + // before it. + let mut starts = Vec::new(); + let mut reader = crate::reader::WalReader::open_raw(&path).unwrap(); + loop { + let start = reader.committed_end(); + if reader.next_record().unwrap().is_none() { + break; + } + starts.push(start); + } + + // Damage the second record, and give the third a version this build + // cannot read. The checksum covers the version, so it is recomputed: + // the third record is intact in every way except its version. The + // fourth stays intact. + smash(&path, starts[1], HEADER_SIZE); + let mut bytes = std::fs::read(&path).unwrap(); + let third = usize::try_from(starts[2]).unwrap(); + let head: &[u8; HEADER_SIZE] = bytes[third..third + HEADER_SIZE].try_into().unwrap(); + let mut header = RecordHeader::from_bytes(head); + let payload_end = third + HEADER_SIZE + header.payload_len as usize; + header.format_version = 1; + header.crc32c = header.compute_checksum(&bytes[third + HEADER_SIZE..payload_end]); + bytes[third..third + HEADER_SIZE].copy_from_slice(&header.to_bytes()); + std::fs::write(&path, &bytes).unwrap(); + + assert_eq!( + classify(&path, starts[1], 1).unwrap(), + TailVerdict::MidFileCorruption { + resync_offset: starts[3], + resync_lsn: 4, + } + ); + } + #[test] fn clean_eof_needs_no_scan() { let dir = tempfile::tempdir().unwrap(); diff --git a/nodedb/src/bootstrap/wal_init.rs b/nodedb/src/bootstrap/wal_init.rs index 7c1d79dbf..0a17aa35b 100644 --- a/nodedb/src/bootstrap/wal_init.rs +++ b/nodedb/src/bootstrap/wal_init.rs @@ -22,6 +22,7 @@ pub fn init_wal( )> { let wal_segment_target = config.checkpoint.wal_segment_target_bytes(); let wal_dir = config.wal_dir(); + let wal = { let mut mgr = WalManager::open_with_tuning(&wal_dir, wal_segment_target, &config.tuning.wal) @@ -37,9 +38,12 @@ pub fn init_wal( info!(next_lsn = %wal.next_lsn(), "WAL ready"); if let Err(e) = wal.validate_for_startup() { + // Neutral on purpose: this wraps every validation error. The wrapped + // error names the cause, such as a WAL format version this build + // cannot read. tracing::error!( error = %e, - "StartupError: WAL validation failed — cannot start with corrupted WAL segments" + "StartupError: WAL validation failed: cannot start with these WAL segments" ); std::process::exit(1); } @@ -104,6 +108,11 @@ pub fn init_wal( /// `O_DIRECT` by design, the server will not downgrade itself to buffered /// writes to get past this, and the operator is the only one who can decide /// between the two ways out. +/// +/// A store whose records are in a format this build cannot read gets one too. +/// The open reaches it first, because it recovers the newest segment, and the +/// WAL layer reports the segment and both versions; what it cannot know is which +/// actions exist, and those belong to the server. fn wal_open_error(wal_dir: &std::path::Path, error: crate::Error) -> anyhow::Error { if let crate::Error::Wal(nodedb_wal::WalError::DirectIoUnsupported { .. }) = error { return anyhow::Error::new(nodedb_types::NodeDbError::wal_at( @@ -111,6 +120,16 @@ fn wal_open_error(wal_dir: &std::path::Path, error: crate::Error) -> anyhow::Err direct_io_unsupported_message(wal_dir), )); } + if let crate::Error::Wal(nodedb_wal::WalError::SegmentFormatVersion { + path, + version, + supported, + }) = &error + { + return anyhow::Error::new(crate::Error::VersionCompat { + detail: crate::wal::manager::replay::version_gap_detail(path, *version, *supported), + }); + } anyhow::Error::new(error) } diff --git a/nodedb/src/ctl/restore/cut_rule.rs b/nodedb/src/ctl/restore/cut_rule.rs index 93eb83683..678904750 100644 --- a/nodedb/src/ctl/restore/cut_rule.rs +++ b/nodedb/src/ctl/restore/cut_rule.rs @@ -83,10 +83,10 @@ pub fn scan_cut( list_through: u64, marks_after: u64, ) -> Result { - let (_, frames) = Frames::open(bytes)?; + let (_, mut frames) = Frames::open(bytes)?; let mut scan = CutScan::default(); let mut last = None; - for frame in frames { + for frame in frames.by_ref() { if frame.is_padding() { continue; } @@ -114,6 +114,7 @@ pub fn scan_cut( } } } + frames.finish(key)?; Ok(scan) } @@ -242,4 +243,28 @@ mod tests { assert_eq!(scan.max_kept, Some(1)); assert!(scan.dropped.is_empty()); } + + /// A segment whose first record is in another WAL format is refused. Read + /// as empty, it would report that no record of the restore point exists. + #[test] + fn a_segment_in_another_format_is_refused_not_read_as_empty() { + let mut bytes = segment(&[(RecordType::Put, 7, 50, b"kept".to_vec())]); + let head: &[u8; nodedb_wal::record::HEADER_SIZE] = + bytes[..nodedb_wal::record::HEADER_SIZE].try_into().unwrap(); + let mut header = nodedb_wal::record::RecordHeader::from_bytes(head); + header.format_version = 1; + bytes[..nodedb_wal::record::HEADER_SIZE].copy_from_slice(&header.to_bytes()); + + match scan_cut("seg", &bytes, &rule(), 9, 0) { + Err(RestoreError::Wal(nodedb_wal::WalError::SegmentFormatVersion { + path, + version, + .. + })) => { + assert_eq!(path, "seg"); + assert_eq!(version, 1); + } + other => panic!("expected a format gap, got {other:?}"), + } + } } diff --git a/nodedb/src/ctl/restore/segment.rs b/nodedb/src/ctl/restore/segment.rs index 7cb69e233..ffa24dfae 100644 --- a/nodedb/src/ctl/restore/segment.rs +++ b/nodedb/src/ctl/restore/segment.rs @@ -9,6 +9,7 @@ //! alignment padding and carry no LSN. use nodedb_types::temporal::LsnTimeAnchor; +use nodedb_wal::WalError; use nodedb_wal::crypto::KeyRing; use nodedb_wal::preamble::{ PREAMBLE_SIZE, SegmentPreamble, WAL_PREAMBLE_MAGIC, parse_leading_preamble, @@ -78,7 +79,13 @@ impl Frame<'_> { pub(super) struct Frames<'a> { bytes: &'a [u8], offset: usize, + /// Where the segment's first record starts. The WAL format version is + /// judged here and nowhere else, as the segment reader does. + records_start: usize, done: bool, + /// Version and supported version of an unreadable first record. A walk + /// that ends here holds no records, and [`Frames::finish`] reports why. + refused: Option<(u16, u16)>, } impl<'a> Frames<'a> { @@ -90,11 +97,29 @@ impl<'a> Frames<'a> { Self { bytes, offset, + records_start: offset, done: false, + refused: None, }, )) } + /// Report a first record in a WAL format version this build cannot read. + /// + /// Call it after the walk. An archive from another build holds intact + /// records, so reading it as an empty segment would drop all of them + /// without an error. `key` names the segment in the error. + pub(super) fn finish(&self, key: &str) -> Result<(), RestoreError> { + match self.refused { + Some((version, supported)) => Err(RestoreError::Wal(WalError::SegmentFormatVersion { + path: key.to_string(), + version, + supported, + })), + None => Ok(()), + } + } + /// Byte offset where records begin: past the preamble, if any. fn records_start(preamble: Option<&SegmentPreamble>) -> usize { if preamble.is_some() { PREAMBLE_SIZE } else { 0 } @@ -118,17 +143,28 @@ impl<'a> Iterator for Frames<'a> { } impl<'a> Frames<'a> { - fn read_at(&self, start: usize) -> Option> { - let head: &[u8; HEADER_SIZE] = self - .bytes + fn read_at(&mut self, start: usize) -> Option> { + let bytes: &'a [u8] = self.bytes; + let head: &[u8; HEADER_SIZE] = bytes .get(start..start.checked_add(HEADER_SIZE)?)? .try_into() .ok()?; let header = RecordHeader::from_bytes(head); - header.validate(start as u64).ok()?; + match header.validate(start as u64) { + Ok(()) => {} + Err(WalError::UnsupportedVersion { version, supported }) => { + // A zeroed version is a torn write, and a mismatch after the + // first record is damage. Both end the walk without an error. + if version != 0 && start == self.records_start { + self.refused = Some((version, supported)); + } + return None; + } + Err(_) => return None, + } let payload_start = start + HEADER_SIZE; let end = payload_start.checked_add(usize::try_from(header.payload_len).ok()?)?; - let payload = self.bytes.get(payload_start..end)?; + let payload = bytes.get(payload_start..end)?; (header.compute_checksum(payload) == header.crc32c).then_some(Frame { header, payload, @@ -144,10 +180,10 @@ pub fn scan_segment( bytes: &[u8], ring: Option<&KeyRing>, ) -> Result { - let (preamble, frames) = Frames::open(bytes)?; + let (preamble, mut frames) = Frames::open(bytes)?; let decryptor = SegmentDecryptor::new(preamble.as_ref(), ring); let mut scan = SegmentScan::default(); - for frame in frames { + for frame in frames.by_ref() { if frame.is_padding() { continue; } @@ -188,6 +224,7 @@ pub fn scan_segment( _ => {} } } + frames.finish(key)?; Ok(scan) } @@ -215,12 +252,12 @@ pub fn cut_segment( if !alignment.is_power_of_two() { return Err(RestoreError::BadAlignment { alignment }); } - let (preamble, frames) = Frames::open(bytes)?; + let (preamble, mut frames) = Frames::open(bytes)?; let mut keep_end = Frames::records_start(preamble.as_ref()); let mut last_kept: Option = None; let mut last_seen: Option = None; let mut dropped_records = 0u64; - for frame in frames { + for frame in frames.by_ref() { if frame.is_padding() { continue; } @@ -233,6 +270,7 @@ pub fn cut_segment( dropped_records += 1; } } + frames.finish(key)?; let mut out = bytes .get(..keep_end) @@ -486,4 +524,70 @@ mod tests { assert_eq!(scan.last_lsn, Some(5)); assert_eq!(scan_segment("seg", &[], None).unwrap().last_lsn, None); } + + /// Where each record of `bytes` starts. + fn record_starts(bytes: &[u8]) -> Vec { + let (preamble, frames) = Frames::open(bytes).unwrap(); + let mut starts = vec![Frames::records_start(preamble.as_ref())]; + starts.extend(frames.map(|frame| frame.end)); + starts.pop(); + starts + } + + /// Declare `version` in the header of the record that starts at `start`. + fn set_format_version(bytes: &mut [u8], start: usize, version: u16) { + let head: &[u8; HEADER_SIZE] = bytes[start..start + HEADER_SIZE].try_into().unwrap(); + let mut header = RecordHeader::from_bytes(head); + header.format_version = version; + bytes[start..start + HEADER_SIZE].copy_from_slice(&header.to_bytes()); + } + + fn assert_format_gap(result: Result, version: u16) { + match result { + Err(RestoreError::Wal(WalError::SegmentFormatVersion { + path, + version: found, + supported, + })) => { + assert_eq!(path, "seg"); + assert_eq!(found, version); + assert_eq!(supported, nodedb_wal::record::WAL_FORMAT_VERSION); + } + other => panic!("expected a format gap, got {other:?}"), + } + } + + /// An archive from another build holds intact records. Reading it as an + /// empty segment would drop all of them without an error. + #[test] + fn an_archive_in_another_format_is_refused_not_read_as_empty() { + let mut bytes = segment(3, None); + let starts = record_starts(&bytes); + for &start in &starts { + set_format_version(&mut bytes, start, 1); + } + + assert_format_gap(scan_segment("seg", &bytes, None), 1); + assert_format_gap(cut_segment("seg", &bytes, 3, &[], None, ALIGN), 1); + } + + #[test] + fn a_version_mismatch_after_the_first_record_stops_the_walk() { + let mut bytes = segment(3, None); + let starts = record_starts(&bytes); + set_format_version(&mut bytes, starts[1], 1); + + let scan = scan_segment("seg", &bytes, None).unwrap(); + assert_eq!(scan.last_lsn, Some(1)); + } + + #[test] + fn a_zeroed_first_version_stops_the_walk() { + let mut bytes = segment(3, None); + let starts = record_starts(&bytes); + set_format_version(&mut bytes, starts[0], 0); + + let scan = scan_segment("seg", &bytes, None).unwrap(); + assert_eq!(scan.last_lsn, None); + } } diff --git a/nodedb/src/wal/manager/replay.rs b/nodedb/src/wal/manager/replay.rs index dac5cc7c8..cf0ed19d1 100644 --- a/nodedb/src/wal/manager/replay.rs +++ b/nodedb/src/wal/manager/replay.rs @@ -17,11 +17,16 @@ impl WalManager { /// does not parse as WAL records is treated as fatal corruption, not as an /// empty WAL. The WAL replay path is lenient (stops at the first invalid /// record) — this method is the complementary hard check run at startup. + /// + /// A segment whose records are in a WAL format this build cannot read is + /// reported as a version gap rather than as corruption. Its answer is to + /// start the store with the build that wrote the segment, or to re-create + /// the data directory; corruption calls for repair. pub fn validate_for_startup(&self) -> crate::Result<()> { let segments = nodedb_wal::segment::discover_segments(&self.wal_dir).map_err(crate::Error::Wal)?; - for seg in &segments { + for (idx, seg) in segments.iter().enumerate() { let file_len = std::fs::metadata(&seg.path).map(|m| m.len()).unwrap_or(0); if file_len == 0 { @@ -47,7 +52,16 @@ impl WalManager { continue; } - let info = nodedb_wal::recovery::recover(&seg.path).map_err(crate::Error::Wal)?; + let info = nodedb_wal::recovery::recover(&seg.path).map_err(|error| match error { + nodedb_wal::error::WalError::SegmentFormatVersion { + path, + version, + supported, + } => crate::Error::VersionCompat { + detail: version_gap_detail(&path, version, supported), + }, + other => crate::Error::Wal(other), + })?; if info.end_offset == 0 { return Err(crate::Error::SegmentCorrupted { @@ -207,11 +221,28 @@ impl WalManager { } } +/// The operator-facing text for a store whose records are in a WAL format this +/// build cannot read. +/// +/// One text, because the same store is refused from two places: the open fails +/// on the newest segment, and startup validation finds the gap in any segment +/// behind it. It names the segment, both versions, and the actions that exist. +/// No WAL format migration tool exists, so naming one would send an operator +/// looking for something they cannot find. +pub(crate) fn version_gap_detail(path: &str, version: u16, supported: u16) -> String { + format!( + "WAL segment '{path}' holds records in WAL format version {version}; this build reads \ + version {supported}. Start the store with the build that wrote that segment, or re-create \ + the data directory and load the data again" + ) +} + #[cfg(test)] mod tests { use super::*; use crate::types::{DatabaseId, TenantId, VShardId}; use crate::wal::manager::NO_APPLY_KEY; + use nodedb_wal::record::{RecordType, WAL_FORMAT_VERSION, WAL_MAGIC, WalRecordArgs}; fn put(wal: &WalManager, payload: &[u8]) -> Lsn { wal.appender(NO_APPLY_KEY) @@ -261,4 +292,291 @@ mod tests { ); assert!(wal.time_anchors().lsn_at_or_before(first_ns - 1).is_err()); } + + /// One record that declares `version`, framed the way this build frames + /// one: header, then payload. + /// + /// Only the version field is under test. The framing is this build's, which + /// a version-1 writer never produced, because its record type field was 16 + /// bits wide and this one is 32, so the file is a store this build cannot + /// read rather than a faithful copy of one it once wrote. The reader decides + /// on the version field before it reaches the checksum, so no assertion here + /// rests on the rest of the header. + fn record_with_format_version(lsn: u64, version: u16) -> Vec { + let mut record = WalRecord::new(WalRecordArgs { + record_type: RecordType::Put as u32, + lsn, + tenant_id: 1, + vshard_id: 0, + database_id: 0, + payload: vec![7u8; 16], + encryption_key: None, + preamble_bytes: None, + }) + .expect("build a record"); + record.header.format_version = version; + let mut bytes = record.header.to_bytes().to_vec(); + bytes.extend_from_slice(&record.payload); + bytes + } + + /// A segment whose first record declares `version`. + fn write_segment_with_format_version( + path: &std::path::Path, + first_lsn: u64, + version: u16, + ) -> std::path::PathBuf { + let bytes = record_with_format_version(first_lsn, version); + let segment = nodedb_wal::segment::segment_path(path, first_lsn); + std::fs::write(&segment, &bytes).expect("write the segment"); + segment + } + + /// A valid store, then a newest segment holding no parseable record: the + /// zero-filled pages a crash can leave behind, which the next boot resumes. + fn wal_with_a_recordless_newest_segment(path: &std::path::Path) -> WalManager { + { + let wal = WalManager::open_for_testing(path).expect("open wal"); + put(&wal, b"a"); + wal.sync().expect("sync"); + } + + let wal = WalManager::open_for_testing(path).expect("reopen wal"); + let newest = nodedb_wal::segment::discover_segments(path) + .expect("discover segments") + .last() + .expect("the boot's own segment") + .first_lsn; + // Written after open on purpose: open resumes the newest segment and may + // rewrite its preamble, and the state under test is the one the gate + // then reads back. + std::fs::write( + nodedb_wal::segment::segment_path(path, newest + 1_000), + vec![0u8; 64 * 1024], + ) + .expect("write a zero-filled segment"); + wal + } + + #[test] + fn a_segment_with_no_records_behind_the_newest_still_blocks_startup() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("wal"); + let wal = wal_with_a_recordless_newest_segment(&path); + let segments = nodedb_wal::segment::discover_segments(&path).expect("discover segments"); + let recordless = segments.last().expect("the zero-filled segment"); + // A newer zero-filled segment becomes the newest, so the first one sits + // behind it: nothing resumes that file, and the gate refuses it. + std::fs::write( + nodedb_wal::segment::segment_path(&path, recordless.first_lsn + 2_000), + vec![0u8; 64 * 1024], + ) + .expect("write a newer zero-filled segment"); + + let error = wal + .validate_for_startup() + .expect_err("a record-less segment behind the newest must be refused"); + assert!( + matches!(error, crate::Error::SegmentCorrupted { .. }), + "a record-less segment behind the newest is corruption: {error}" + ); + let text = format!("{error}"); + assert!( + text.contains(&recordless.path.display().to_string()), + "the refusal must name the record-less segment: {text}" + ); + } + + /// An encrypted store keeps its first record behind the preamble, so the + /// version has to be read from behind it too. + #[test] + fn a_store_with_a_preamble_says_which_version() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("wal"); + { + let wal = WalManager::open_for_testing(&path).expect("open wal"); + put(&wal, b"a"); + wal.sync().expect("sync"); + } + let wal = WalManager::open_for_testing(&path).expect("reopen wal"); + let newest = nodedb_wal::segment::discover_segments(&path) + .expect("discover segments") + .last() + .expect("the boot's own segment") + .first_lsn; + + // A real preamble, not a zeroed one: the reader validates its own + // version field and would reject a malformed preamble for the wrong + // reason. + let mut bytes = nodedb_wal::preamble::SegmentPreamble::new_wal([9, 9, 9, 9]) + .to_bytes() + .to_vec(); + bytes.extend_from_slice(&record_with_format_version(newest + 1_000, 1)); + std::fs::write( + nodedb_wal::segment::segment_path(&path, newest + 1_000), + &bytes, + ) + .expect("write an encrypted-format segment"); + // A newer segment, so the one above is not the tail the gate tolerates. + std::fs::write( + nodedb_wal::segment::segment_path(&path, newest + 2_000), + vec![0u8; 64 * 1024], + ) + .expect("write a newer segment"); + + let error = wal + .validate_for_startup() + .expect_err("an older format behind a preamble must be refused"); + let text = format!("{error}"); + assert!( + text.contains("version 1"), + "the version sits behind the preamble and the message must still name it: {text}" + ); + assert!( + text.contains(&format!("reads version {WAL_FORMAT_VERSION}")), + "the message must name the version this build reads: {text}" + ); + } + + /// A valid store, then a newer segment written in an older format. + fn store_with_an_older_newest_segment( + path: &std::path::Path, + version: u16, + ) -> (WalManager, std::path::PathBuf) { + { + let wal = WalManager::open_for_testing(path).expect("open wal"); + put(&wal, b"a"); + wal.sync().expect("sync"); + } + let wal = WalManager::open_for_testing(path).expect("reopen wal"); + let newest = nodedb_wal::segment::discover_segments(path) + .expect("discover segments") + .last() + .expect("the boot's own segment") + .first_lsn; + let segment = write_segment_with_format_version(path, newest + 1_000, version); + (wal, segment) + } + + /// The newest segment is the one the open recovers, so a store in another + /// format fails there, before validation runs. The refusal names the + /// segment and both versions: that is what tells a store written by another + /// build apart from a damaged one. + #[test] + fn the_open_refuses_a_newest_segment_in_another_format() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("wal"); + let (_wal, segment) = store_with_an_older_newest_segment(&path, 1); + + let error = match WalManager::open_for_testing(&path) { + Ok(_) => panic!("a store written in another format must not open"), + Err(error) => error, + }; + let text = format!("{error}"); + + assert!( + text.contains(&segment.display().to_string()), + "the refusal must name the segment: {text}" + ); + assert!( + text.contains("version 1"), + "the refusal must name the version it found: {text}" + ); + assert!( + text.contains(&format!("reads version {WAL_FORMAT_VERSION}")), + "the refusal must name the version this build reads: {text}" + ); + assert!( + !text.contains("corrupted"), + "a format gap is not corruption: {text}" + ); + } + + /// The same gap behind the newest segment is out of the open's reach, so + /// startup validation is what finds it. That message also names the actions + /// that exist, which the WAL layer cannot know. + #[test] + fn validation_names_a_gap_behind_the_newest_segment() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("wal"); + let (wal, segment) = store_with_an_older_newest_segment(&path, 1); + let newest = nodedb_wal::segment::discover_segments(&path) + .expect("discover segments") + .last() + .expect("a segment") + .first_lsn; + std::fs::write( + nodedb_wal::segment::segment_path(&path, newest + 1_000), + vec![0u8; 64 * 1024], + ) + .expect("write a newer segment"); + + let error = wal + .validate_for_startup() + .expect_err("a gap behind the tail is still a gap"); + let text = format!("{error}"); + + assert!( + text.contains(&segment.display().to_string()), + "the message must name the segment: {text}" + ); + assert!( + text.contains("version 1"), + "the message must name the version it found, and it named none: {text}" + ); + assert!( + text.contains(&format!("reads version {WAL_FORMAT_VERSION}")), + "the message must name the version this build reads: {text}" + ); + assert!( + !text.contains("corrupted"), + "a format gap is reported as a version gap, never as corruption: {text}" + ); + assert!( + text.contains("re-create the data directory"), + "the message must name an action that exists: {text}" + ); + } + + /// The other side of the same branch: a segment that is not a WAL at all is + /// still corruption, and still says so. + #[test] + fn a_segment_that_is_not_a_wal_still_says_corrupted() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("wal"); + { + let wal = WalManager::open_for_testing(&path).expect("open wal"); + put(&wal, b"a"); + wal.sync().expect("sync"); + } + let wal = WalManager::open_for_testing(&path).expect("reopen wal"); + let newest = nodedb_wal::segment::discover_segments(&path) + .expect("discover segments") + .last() + .expect("a segment") + .first_lsn; + // No magic at all, which is what a truncated or overwritten file looks + // like, as opposed to a file another version wrote. A newer segment + // follows it so this one is not the tail, because the tail is the one + // case the gate tolerates. + std::fs::write( + nodedb_wal::segment::segment_path(&path, newest + 1_000), + vec![0u8; 64 * 1024], + ) + .expect("write a segment with no framing"); + std::fs::write( + nodedb_wal::segment::segment_path(&path, newest + 2_000), + vec![0u8; 64 * 1024], + ) + .expect("write a newer zero-filled segment"); + + let error = wal + .validate_for_startup() + .expect_err("a segment with no framing must be refused"); + let text = format!("{error}"); + assert!( + text.contains("corrupted"), + "a segment that is not a WAL is corruption: {text}" + ); + } } From b112e60c62f6675e1ffbfdc269b6d290ad173ee4 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Fri, 9 Oct 2026 04:58:41 +0800 Subject: [PATCH 3/3] fix(wal): start over a newest segment that holds no parseable record A crash can leave the newest segment non-empty and holding nothing that parses, such as zero-filled pages, and the next boot resumes that same file. Refusing it leaves the store permanently unbootable, because every later boot reads the same bytes and refuses them again. The newest segment is now the one exemption from the empty-segment check. A record-less segment anywhere else stays fatal, since nothing resumes it and replay would silently skip whatever it held. Only segments without a preamble reach that check at all: a segment that opens with one reports its end at the end of the preamble, never at zero. The version check runs first, so a segment written by another build is still refused by name rather than accepted as an empty tail. --- nodedb/src/wal/manager/replay.rs | 68 +++++++++++++++++++++++++++++++- 1 file changed, 67 insertions(+), 1 deletion(-) diff --git a/nodedb/src/wal/manager/replay.rs b/nodedb/src/wal/manager/replay.rs index cf0ed19d1..ccdaafa5e 100644 --- a/nodedb/src/wal/manager/replay.rs +++ b/nodedb/src/wal/manager/replay.rs @@ -11,7 +11,9 @@ impl WalManager { /// /// Returns `Err` if any non-empty segment contains no valid WAL records — /// a reliable signal that the segment was corrupted (wrong magic, truncated - /// header, etc.) rather than simply rolled over empty. + /// header, etc.) rather than simply rolled over empty. The newest segment is + /// the exception: a crash can leave it with no parseable record, and the + /// next boot resumes that same file, so it passes. /// /// This check is intentionally strict: a segment file with content that /// does not parse as WAL records is treated as fatal corruption, not as an @@ -64,6 +66,23 @@ impl WalManager { })?; if info.end_offset == 0 { + // This check covers only segments without a preamble. A + // segment that opens with one reports its end at the end of the + // preamble, never at 0, even when no record follows it. + // + // The newest segment is the one exception. A crash can leave + // zero-filled pages in it, and the next boot resumes that same + // file. Refusing it here makes the store permanently + // unbootable, because every later boot reads the same bytes and + // refuses them again. The version check runs before this one, + // so a segment written by another build is still refused above. + // + // Any other segment without a preamble and without a record is + // refused: nothing resumes it, and replay would silently skip + // whatever it held. + if idx + 1 == segments.len() { + continue; + } return Err(crate::Error::SegmentCorrupted { detail: format!( "WAL segment '{}' is non-empty ({file_len} bytes) but contains no valid \ @@ -358,6 +377,20 @@ mod tests { wal } + #[test] + fn a_recordless_newest_segment_does_not_block_startup() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("wal"); + let wal = wal_with_a_recordless_newest_segment(&path); + + assert!( + wal.validate_for_startup().is_ok(), + "the newest segment holds no committed record, which torn_tail already warns \ + about, and the next boot resumes it. Refusing it here makes the store \ + permanently unbootable" + ); + } + #[test] fn a_segment_with_no_records_behind_the_newest_still_blocks_startup() { let dir = tempfile::tempdir().expect("tempdir"); @@ -438,6 +471,39 @@ mod tests { ); } + /// A torn write that persists the magic and leaves the version zeroed is + /// not a format gap, and must not turn a bootable store into a refused one. + #[test] + fn a_zeroed_version_is_not_a_format_gap() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("wal"); + { + let wal = WalManager::open_for_testing(&path).expect("open wal"); + put(&wal, b"a"); + wal.sync().expect("sync"); + } + let wal = WalManager::open_for_testing(&path).expect("reopen wal"); + let newest = nodedb_wal::segment::discover_segments(&path) + .expect("discover segments") + .last() + .expect("the boot's own segment") + .first_lsn; + + // The magic of a real record, the version bytes never written. + let mut bytes = vec![0u8; 64 * 1024]; + bytes[0..4].copy_from_slice(&WAL_MAGIC.to_le_bytes()); + std::fs::write( + nodedb_wal::segment::segment_path(&path, newest + 1_000), + &bytes, + ) + .expect("write a torn newest segment"); + + assert!( + wal.validate_for_startup().is_ok(), + "a zeroed version is uninitialised bytes, not a format this build cannot read" + ); + } + /// A valid store, then a newer segment written in an older format. fn store_with_an_older_newest_segment( path: &std::path::Path,