Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions docs/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -239,6 +239,13 @@ pgwire = true
http = true
resp = true
ilp = false # Example: disable TLS for ILP ingest

[tuning.startup]
# Boot bounds for the readiness gates (defaults shown). Raise them when a
# backlogged restart needs minutes of replay; a bound set under [server] is
# rejected at load with the [tuning.startup] path named.
raft_ready_timeout_ms = 300000 # metadata group may stall this long (5 min)
data_group_recovery_timeout_ms = 600000 # data groups must replay within (10 min)
```

Every listener binds during startup, before the server accepts any
Expand Down
5 changes: 5 additions & 0 deletions nodedb-types/src/config/tuning/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
///
Expand Down Expand Up @@ -48,6 +49,8 @@ pub struct TuningConfig {
#[serde(default)]
pub shutdown: ShutdownTuning,
#[serde(default)]
pub startup: StartupTuning,
#[serde(default)]
pub bitemporal: BitemporalTuning,
#[serde(default)]
pub maintenance: MaintenanceTuning,
Expand All @@ -73,6 +76,8 @@ mod tests {
let parsed: TuningConfig = toml::from_str(&toml_str).expect("deserialize");

assert_eq!(parsed.data_plane.idle_poll_timeout_ms, 100);
assert_eq!(parsed.startup.raft_ready_timeout_ms, 300_000);
assert_eq!(parsed.startup.data_group_recovery_timeout_ms, 600_000);
assert_eq!(parsed.query.sort_run_size, 100_000);
assert_eq!(parsed.vector.flat_index_threshold, 10_000);
assert_eq!(parsed.sparse.bm25_k1, 1.2);
Expand Down
2 changes: 2 additions & 0 deletions nodedb-types/src/config/tuning/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ mod memory;
mod network;
mod scheduler;
mod shutdown;
mod startup;

pub use bitemporal::BitemporalTuning;
pub use config::TuningConfig;
Expand All @@ -23,3 +24,4 @@ pub use memory::MemoryTuning;
pub use network::{BridgeTuning, ClusterTransportTuning, NetworkTuning, WalTuning};
pub use scheduler::SchedulerTuning;
pub use shutdown::ShutdownTuning;
pub use startup::StartupTuning;
82 changes: 82 additions & 0 deletions nodedb-types/src/config/tuning/startup.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
// SPDX-License-Identifier: Apache-2.0

//! Startup tuning — boot-time bounds applied by the readiness gates.

use std::time::Duration;

use serde::{Deserialize, Serialize};

fn default_raft_ready_timeout_ms() -> u64 {
// 5 minutes: 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 turned a slow boot into a restart loop.
300_000
}

fn default_data_group_recovery_timeout_ms() -> u64 {
// 10 minutes: the value production carried before the bound became a
// hard-coded constant.
600_000
}

/// Boot-time bounds for the readiness gates in `bootstrap::cluster_ready`.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StartupTuning {
/// How long the metadata raft group may 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 as a `Duration`.
pub fn raft_ready_timeout(&self) -> Duration {
Duration::from_millis(self.raft_ready_timeout_ms)
}

/// Data-group recovery bound as a `Duration`.
pub fn data_group_recovery_timeout(&self) -> Duration {
Duration::from_millis(self.data_group_recovery_timeout_ms)
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn defaults_are_five_and_ten_minutes() {
let t = StartupTuning::default();
assert_eq!(t.raft_ready_timeout(), Duration::from_secs(300));
assert_eq!(t.data_group_recovery_timeout(), Duration::from_secs(600));
}

#[test]
fn both_bounds_take_overrides() {
let parsed: StartupTuning =
toml::from_str("raft_ready_timeout_ms = 1000\ndata_group_recovery_timeout_ms = 2000\n")
.unwrap();
assert_eq!(parsed.raft_ready_timeout(), Duration::from_millis(1000));
assert_eq!(
parsed.data_group_recovery_timeout(),
Duration::from_millis(2000)
);
}
}
30 changes: 14 additions & 16 deletions nodedb/src/bootstrap/cluster_ready.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,11 +28,16 @@ pub struct ClusterReadyGates {
///
/// In single-node mode `raft_ready_rx` is `None` and the raft-ready wait is
/// skipped. Gate fires are always performed regardless of cluster mode.
///
/// `raft_ready_timeout` and `data_group_recovery_timeout` are the boot bounds
/// from `[tuning.startup]`.
pub async fn await_cluster_ready(
shared: &Arc<SharedState>,
raft_ready_rx: Option<tokio::sync::watch::Receiver<bool>>,
data_plane_replay_done: Vec<tokio::sync::oneshot::Receiver<()>>,
gates: ClusterReadyGates,
raft_ready_timeout: Duration,
data_group_recovery_timeout: Duration,
) -> anyhow::Result<()> {
let ClusterReadyGates {
raft_gate,
Expand All @@ -55,7 +60,7 @@ pub async fn await_cluster_ready(
shared,
&mut ready_rx,
&raft_gate,
RAFT_READY_STALL_TIMEOUT,
raft_ready_timeout,
RAFT_READY_POLL_INTERVAL,
)
.await?;
Expand Down Expand Up @@ -143,7 +148,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);
}
Expand Down Expand Up @@ -197,19 +207,6 @@ 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.
/// How long the metadata group may go without *any* applied entry before the
/// readiness gate fails startup.
///
/// Raised from 30 s after the 2026-09-20 incident: with a large apply backlog
/// (tens of thousands of entries from a burst of cross-shard writes) the group
/// needs minutes, and aborting the start turned a slow boot into a restart
/// loop. The gate still fails a group that never applies anything.
/// Follow-up: make this configurable.
const RAFT_READY_STALL_TIMEOUT: Duration = Duration::from_secs(300);

/// How often the stall check samples the applied index while waiting.
const RAFT_READY_POLL_INTERVAL: Duration = Duration::from_secs(1);

Expand All @@ -218,7 +215,8 @@ const RAFT_READY_POLL_INTERVAL: Duration = Duration::from_secs(1);
/// Bounds the wait on lack of PROGRESS, not on total elapsed time. A node
/// replaying a large log keeps advancing `applied_index` and must be allowed
/// to finish; a group that is genuinely stuck advances nothing and fails
/// after [`RAFT_READY_STALL_TIMEOUT`].
/// after `stall_timeout` (`[tuning.startup] raft_ready_timeout_ms`, five
/// minutes by default).
async fn wait_for_raft_ready(
shared: &Arc<SharedState>,
ready_rx: &mut tokio::sync::watch::Receiver<bool>,
Expand Down
19 changes: 11 additions & 8 deletions nodedb/src/bootstrap/data_group_recovery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,11 +31,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(600);

/// 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.
///
Expand Down Expand Up @@ -170,14 +165,22 @@ fn pending_groups(statuses: Vec<nodedb_cluster::GroupStatus>) -> Vec<PendingGrou
/// Hold startup until every locally hosted data group has replayed its retained
/// Raft log.
///
/// `timeout` (from `[tuning.startup] data_group_recovery_timeout_ms`) bounds
/// 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.
///
/// A node with no Raft status source (a deployment with no cluster handle
/// installed) hosts no data groups and returns immediately.
pub async fn await_data_group_recovery(shared: &Arc<SharedState>) -> anyhow::Result<()> {
pub async fn await_data_group_recovery(
shared: &Arc<SharedState>,
timeout: Duration,
) -> anyhow::Result<()> {
let Some(status_fn) = shared.raft_status_fn.get() else {
return Ok(());
};
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());
Expand All @@ -193,7 +196,7 @@ pub async fn await_data_group_recovery(shared: &Arc<SharedState>) -> anyhow::Res
.collect::<Vec<_>>()
.join("; ");
return Err(anyhow::anyhow!(
"data raft group recovery timeout after {DATA_GROUP_RECOVERY_TIMEOUT:?}: {detail}"
"data raft group recovery timeout after {timeout:?}: {detail}"
));
}

Expand Down
68 changes: 68 additions & 0 deletions nodedb/src/config/server/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,30 @@ pub struct ServerConfig {
pub scheduler: SchedulerConfig,
}

/// Rejects startup bounds written where they used to live, naming the moved-to
/// path. `deny_unknown_fields` catches the key anyway, but its message lists
/// every valid field instead of the replacement.
fn reject_moved_startup_bounds(content: &str) -> crate::Result<()> {
let Ok(doc) = toml::from_str::<toml::Table>(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}` moved to `[tuning.startup] {key}`; set the bound there"
),
});
}
}
Ok(())
}

impl ServerConfig {
/// Load configuration from a TOML file, falling back to defaults.
pub fn from_file(path: &std::path::Path) -> crate::Result<Self> {
Expand All @@ -131,6 +155,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_moved_startup_bounds(&content)?;
let parsed: Self = toml::from_str(&content).map_err(|e| crate::Error::Config {
detail: format!("invalid TOML config: {e}"),
})?;
Expand Down Expand Up @@ -387,6 +412,49 @@ mod tests {
assert!(msg.contains("positive integer"), "{msg}");
}

/// The pre-tuning location of a startup bound is rejected with the new
/// path named, not serde's full valid-field list.
#[test]
fn from_file_rejects_a_startup_bound_at_its_old_server_path() {
let path = write_temp_config(
"nodedb-moved-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 moved key is rejected under either legacy name.
#[test]
fn from_file_rejects_the_legacy_data_group_recovery_bound() {
let path = write_temp_config(
"nodedb-moved-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);
}

#[test]
fn from_file_rejects_scope_expiry_below_the_floor() {
let path = write_temp_config(
Expand Down
3 changes: 3 additions & 0 deletions nodedb/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -260,6 +260,7 @@ async fn server_main() -> anyhow::Result<()> {
};

// Wait for raft readiness, run catalog sanity check, warm peer cache, fire gates.
// The two boot bounds come from [tuning.startup].
nodedb::bootstrap::cluster_ready::await_cluster_ready(
&shared,
raft_ready_rx,
Expand All @@ -274,6 +275,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?;

Expand Down
Loading