Skip to content
Closed
10 changes: 9 additions & 1 deletion nodedb/src/bootstrap/cluster_ready.rs
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,15 @@ pub async fn await_cluster_ready(
/// 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 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 Down
2 changes: 1 addition & 1 deletion nodedb/src/bootstrap/data_group_recovery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ 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);
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
36 changes: 30 additions & 6 deletions nodedb/src/bridge/dispatch/dispatcher.rs
Original file line number Diff line number Diff line change
Expand Up @@ -171,11 +171,13 @@ impl Dispatcher {
if self.max_per_tenant_inflight > 0 {
let inflight = self.tenant_inflight.get(&tenant_id).copied().unwrap_or(0);
if inflight >= self.max_per_tenant_inflight {
return Err(crate::Error::Dispatch {
detail: format!(
"tenant {tenant_id}: queue full ({inflight}/{} in-flight)",
self.max_per_tenant_inflight
),
// Retryable by type (issue 352): the request was not enqueued
// and nothing was applied. The Calvin scheduler backs off and
// re-drives instead of logging an error per attempt.
return Err(crate::Error::DispatchCapacityBusy {
tenant_id,
inflight,
cap: self.max_per_tenant_inflight,
});
}
}
Expand Down Expand Up @@ -559,6 +561,28 @@ mod tests {
assert_eq!(&*responses[0].payload, b"result");
}

/// A full per-tenant queue must surface as a *retryable capacity*
/// condition, never as a generic `Error::Dispatch` that callers treat as a
/// terminal failure. Issue 352: the Calvin scheduler completes the
/// transaction as failed and logs ERROR on `Error::Dispatch`, which turns
/// the startup rebuild into an error storm.
#[test]
fn full_queue_is_a_retryable_capacity_condition() {
let (mut dispatcher, _data_sides) = Dispatcher::new(1, 4);
for i in 0..4u64 {
dispatcher
.dispatch(make_request_for_db(0, i + 1, i + 1))
.unwrap();
}
let err = dispatcher
.dispatch(make_request_for_db(0, 99, 99))
.unwrap_err();
assert!(
matches!(err, crate::Error::DispatchCapacityBusy { .. }),
"a full queue must be a retryable capacity signal (issue 352); got {err:?}"
);
}

#[test]
fn full_queue_returns_error() {
// With WFQ capacity == ring capacity, filling WFQ should eventually
Expand Down Expand Up @@ -637,7 +661,7 @@ mod tests {
let _ = dispatcher.db_pressure_on_core(0, 2);
}

// --- Dead-core request loss (GitHub #265) ---
// --- Dead-core request loss ---
//
// When a Data Plane core's consumer/producer is dropped (the core thread
// died), `Dispatcher` must synthesize an error `Response` for every
Expand Down
58 changes: 58 additions & 0 deletions nodedb/src/control/cluster/calvin/scheduler/driver/core/busy.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
// SPDX-License-Identifier: BUSL-1.1

//! Retryable capacity-busy accounting and the bounded re-drive backoff.

use std::time::{Duration, Instant};

use super::scheduler::Scheduler;

impl Scheduler {
/// Note a retryable capacity-busy dispatch outcome: count it, arm the
/// drain backoff, and return true. Shared by every dispatch site so none of
/// them logs ERROR or drives the epoch unpaced.
pub(in crate::control::cluster::calvin::scheduler::driver::core) fn note_dispatch_busy(
&mut self,
e: &crate::Error,
epoch: u64,
) -> bool {
if !matches!(e, crate::Error::DispatchCapacityBusy { .. }) {
return false;
}
self.metrics.record_dispatch_busy();
let attempts = self.dispatch_busy_attempts.saturating_add(1);
self.dispatch_busy_attempts = attempts;
self.dispatch_busy_until =
// no-determinism: the backoff deadline is scheduler observability, not Calvin WAL data
Some(Instant::now() + Self::busy_backoff(attempts, self.vshard_id, epoch));
true
}

/// Count a capacity-busy dispatch from a `&self` context (no backoff
/// arming: the drain gates on the `&mut self` sites).
pub(in crate::control::cluster::calvin::scheduler::driver::core) fn note_dispatch_busy_shared(
&self,
e: &crate::Error,
) -> bool {
if !matches!(e, crate::Error::DispatchCapacityBusy { .. }) {
return false;
}
self.metrics.record_dispatch_busy();
true
}

/// Bounded, deterministic backoff for a capacity-busy dispatch.
///
/// 5 ms doubling on consecutive busy outcomes, capped at 250 ms, plus a
/// small deterministic jitter derived from the vShard and epoch so peers do
/// not re-drive in lockstep. No RNG: the value is reproducible in tests.
pub(in crate::control::cluster::calvin::scheduler::driver::core) fn busy_backoff(
attempts: u32,
vshard_id: u32,
epoch: u64,
) -> Duration {
let base = 5u64.saturating_mul(1u64 << attempts.min(6));
let capped = base.min(250);
let jitter = (u64::from(vshard_id).wrapping_add(epoch)) % 7;
Duration::from_millis(capped + jitter)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,18 @@ impl Scheduler {
};
if let Err(e) = dispatch_result {
self.shared.tracker.cancel(&request_id);
if self.note_dispatch_busy(&e, epoch) {
// Retryable capacity condition (issue 352): the catch-up drain
// re-drives after the armed backoff; no error log per attempt.
tracing::debug!(
vshard_id = self.vshard_id,
epoch,
position,
error = %e,
"calvin: CalvinResolve dispatch deferred (capacity busy)"
);
return false;
}
tracing::error!(
vshard_id = self.vshard_id,
epoch,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,18 @@ impl Scheduler {
};
if let Err(error) = dispatch_result {
self.shared.tracker.cancel(&request_id);
if self.note_dispatch_busy(&error, epoch) {
// Retryable capacity condition (issue 352): paced by the drain.
tracing::debug!(
vshard_id = self.vshard_id,
epoch,
position,
committed,
%error,
"calvin: commit resolution dispatch deferred (capacity busy)"
);
return false;
}
tracing::error!(
vshard_id = self.vshard_id,
epoch,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@

use std::time::Instant;

use tracing::error;
use tracing::{debug, error};

use nodedb_cluster::calvin::types::SequencedTxn;
use nodedb_physical::physical_plan::PhysicalPlan;
Expand Down Expand Up @@ -121,6 +121,17 @@ impl Scheduler {
};

if let Err(e) = dispatch_result {
if self.note_dispatch_busy(&e, epoch) {
debug!(
vshard_id = self.vshard_id,
epoch,
position,
error = %e,
"calvin scheduler: active dispatch deferred (capacity busy)"
);
self.on_txn_complete(txn_id);
return;
}
error!(
vshard_id = self.vshard_id,
epoch,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@

use std::time::Instant;

use tracing::error;
use tracing::{debug, error};

use nodedb_cluster::calvin::types::SequencedTxn;
use nodedb_physical::physical_plan::PhysicalPlan;
Expand Down Expand Up @@ -273,6 +273,20 @@ impl Scheduler {
};

if let Err(e) = dispatch_result {
if self.note_dispatch_busy(&e, epoch) {
// Retryable capacity condition: nothing was enqueued. Demote
// the log and pace the re-drive via the scheduler's drain gate
// instead of spinning on ERROR.
debug!(
vshard_id = self.vshard_id,
epoch,
position,
error = %e,
"calvin scheduler: dispatch deferred (capacity busy)"
);
self.on_txn_complete(txn_id);
return;
}
error!(
vshard_id = self.vshard_id,
epoch,
Expand All @@ -283,6 +297,8 @@ impl Scheduler {
self.on_txn_complete(txn_id);
return;
}
self.dispatch_busy_attempts = 0;
self.dispatch_busy_until = None;

self.metrics.record_dispatch();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
//! Sub-modules (one concern per file):
//!
//! - [`scheduler`] — `Scheduler` struct, ctor, run loop.
//! - [`busy`] — capacity-busy accounting and the bounded re-drive backoff.
//! - [`completion_route`] — routes each executor response (disconnect, OLLP
//! mismatch, staged commit-resolution state, or direct apply) to its handler.
//! - [`process`] — new-txn processing, dependent-read barrier setup,
Expand Down Expand Up @@ -45,6 +46,7 @@
//!
//! Never used for WAL-influencing values.

pub mod busy;
pub mod catch_up;
pub mod commit_redo;
pub mod commit_resolution_dispatch;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,10 @@

use std::collections::BTreeMap;
use std::sync::{Arc, Mutex};
use std::time::Instant;

use tokio::sync::mpsc;
use tracing::info;
use tracing::{debug, info};

use nodedb_cluster::MultiRaft;
use nodedb_cluster::calvin::types::SchedulerInput;
Expand Down Expand Up @@ -95,6 +96,12 @@ pub struct Scheduler {
/// Rebuild target epoch (highest applied epoch from the initial recovery
/// scan).
pub(in crate::control::cluster::calvin::scheduler::driver::core) rebuild_target_epoch: u64,
/// Pacing gate for the catch-up drain while dispatch capacity is saturated
/// (issue 352): the stall tick defers its re-drive until this instant.
pub(in crate::control::cluster::calvin::scheduler::driver::core) dispatch_busy_until:
Option<Instant>,
/// Consecutive capacity-busy dispatch outcomes; drives the bounded backoff.
pub(in crate::control::cluster::calvin::scheduler::driver::core) dispatch_busy_attempts: u32,
/// Highest replicated epoch observed across all scheduler inputs so far.
/// Advances monotonically as `process_scheduler_input` sees new inputs; the
/// lease-based reservation reap uses it (minus `LEASE_EPOCHS`) as the
Expand Down Expand Up @@ -218,6 +225,8 @@ impl Scheduler {
read_result_rx,
applied: AppliedGate::new(fully_applied_epoch, applied_tail),
rebuild_target_epoch,
dispatch_busy_until: None,
dispatch_busy_attempts: 0,
max_input_epoch: 0,
config,
metrics,
Expand Down Expand Up @@ -371,11 +380,25 @@ impl Scheduler {
}

_ = stall_tick.tick() => {
// Replay any sequencer-fan-out inputs dropped on this replica
// (channel Full/Closed) so a missed `SchedulerInput` never
// permanently diverges this vShard's lock table from its peers.
// O(1) common case (no pending catch-up). See `drain_catch_up`.
self.drain_catch_up();
// Issue 352: while dispatch capacity is saturated, defer the
// catch-up drain's re-drive instead of spinning on a full
// per-tenant queue. The next tick after the backoff retries.
// no-determinism: the drain gate reads the wall clock; not Calvin WAL data
let now = Instant::now();
if self.dispatch_busy_until.map(|t| now < t).unwrap_or(false) {
debug!(
vshard_id = self.vshard_id,
"calvin scheduler: drain deferred (capacity busy)"
);
} else {
self.dispatch_busy_until = None;
// Replay any sequencer-fan-out inputs dropped on this
// replica (channel Full/Closed) so a missed
// `SchedulerInput` never permanently diverges this
// vShard's lock table from its peers. O(1) common case
// (no pending catch-up). See `drain_catch_up`.
self.drain_catch_up();
}
// The top-of-loop check_awaiting_verdict_stalls /
// check_dependent_barrier_timeouts do the stall work on every
// wake; this arm guarantees the loop wakes to run them (and the
Expand Down Expand Up @@ -438,6 +461,7 @@ impl Scheduler {
mod tests {
use super::*;
use std::collections::{BTreeSet, HashMap};
use std::time::Duration;

use nodedb_cluster::RoutingTable;

Expand Down Expand Up @@ -565,4 +589,18 @@ mod tests {
"no rebuild target (greenfield node) must report caught-up"
);
}

#[test]
fn busy_backoff_is_bounded_and_jittered() {
assert_eq!(Scheduler::busy_backoff(0, 0, 0), Duration::from_millis(5));
assert_eq!(Scheduler::busy_backoff(3, 0, 0), Duration::from_millis(40));
assert!(Scheduler::busy_backoff(30, 0, 0) <= Duration::from_millis(257));
// jitter varies with vshard + epoch and stays inside the 0..=6 ms band
assert_ne!(
Scheduler::busy_backoff(2, 1, 0),
Scheduler::busy_backoff(2, 1, 1)
);
assert!(Scheduler::busy_backoff(2, 0, 0) >= Duration::from_millis(20));
assert!(Scheduler::busy_backoff(2, 0, 6) <= Duration::from_millis(26));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,17 @@ impl Scheduler {
};
if let Err(e) = dispatch_result {
self.shared.tracker.cancel(&request_id);
if self.note_dispatch_busy_shared(&e) {
// Retryable capacity condition (issue 352): counted + paced.
tracing::debug!(
vshard_id = self.vshard_id,
epoch,
position,
error = %e,
"calvin: write-version record dispatch deferred (capacity busy)"
);
return;
}
tracing::warn!(
vshard_id = self.vshard_id,
epoch,
Expand Down
Loading
Loading