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-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/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/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/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/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/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?; diff --git a/nodedb/src/wal/manager/replay.rs b/nodedb/src/wal/manager/replay.rs index dac5cc7c8..ccdaafa5e 100644 --- a/nodedb/src/wal/manager/replay.rs +++ b/nodedb/src/wal/manager/replay.rs @@ -11,17 +11,24 @@ 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 /// 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,9 +54,35 @@ 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 { + // 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 \ @@ -207,11 +240,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 +311,338 @@ 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_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"); + 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 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, + 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}" + ); + } }