diff --git a/nodedb/src/bootstrap/cluster_ready.rs b/nodedb/src/bootstrap/cluster_ready.rs index 731445079..6462821e6 100644 --- a/nodedb/src/bootstrap/cluster_ready.rs +++ b/nodedb/src/bootstrap/cluster_ready.rs @@ -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); diff --git a/nodedb/src/bootstrap/data_group_recovery.rs b/nodedb/src/bootstrap/data_group_recovery.rs index b050d30bd..d0da54919 100644 --- a/nodedb/src/bootstrap/data_group_recovery.rs +++ b/nodedb/src/bootstrap/data_group_recovery.rs @@ -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. diff --git a/nodedb/src/bridge/dispatch/dispatcher.rs b/nodedb/src/bridge/dispatch/dispatcher.rs index 7e1783f3e..6f531ce38 100644 --- a/nodedb/src/bridge/dispatch/dispatcher.rs +++ b/nodedb/src/bridge/dispatch/dispatcher.rs @@ -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, }); } } @@ -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 @@ -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 diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/busy.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/busy.rs new file mode 100644 index 000000000..b046bd7c9 --- /dev/null +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/busy.rs @@ -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) + } +} diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/commit_redo.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/commit_redo.rs index d9181f9bb..0c64aed57 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/commit_redo.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/commit_redo.rs @@ -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, diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/commit_resolution_dispatch.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/commit_resolution_dispatch.rs index 42dc11229..558d6e347 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/commit_resolution_dispatch.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/commit_resolution_dispatch.rs @@ -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, diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/dispatch/active_dispatch.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/dispatch/active_dispatch.rs index 206d24982..c84baf713 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/dispatch/active_dispatch.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/dispatch/active_dispatch.rs @@ -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; @@ -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, diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/dispatch/static_dispatch.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/dispatch/static_dispatch.rs index e8443cf30..d347f80f8 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/dispatch/static_dispatch.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/dispatch/static_dispatch.rs @@ -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; @@ -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, @@ -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(); diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/mod.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/mod.rs index e0f25ee2b..2f3cee115 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/mod.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/mod.rs @@ -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, @@ -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; diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs index 40aecbc67..1194f2c02 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs @@ -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; @@ -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, + /// 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 @@ -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, @@ -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 @@ -438,6 +461,7 @@ impl Scheduler { mod tests { use super::*; use std::collections::{BTreeSet, HashMap}; + use std::time::Duration; use nodedb_cluster::RoutingTable; @@ -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)); + } } diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/write_version_record.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/write_version_record.rs index 497cc8111..5b9b0f2ad 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/write_version_record.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/write_version_record.rs @@ -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, diff --git a/nodedb/src/control/cluster/calvin/scheduler/metrics.rs b/nodedb/src/control/cluster/calvin/scheduler/metrics.rs index d8ac0d837..7748d95ac 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/metrics.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/metrics.rs @@ -17,6 +17,8 @@ pub const EXECUTOR_TXN_DURATION_BUCKETS: &[u64] = &[1, 5, 10, 50, 100, 500, 1000 pub struct SchedulerMetrics { /// Total transactions dispatched to the Data Plane executor. pub dispatch_count: AtomicU64, + /// Total dispatches deferred because a per-tenant queue was full (issue 352). + pub dispatch_busy_count: AtomicU64, /// Total transactions that were blocked on lock acquisition. pub blocked_count: AtomicU64, /// Total lock-wait duration in milliseconds (sum across all txns). @@ -89,6 +91,11 @@ impl SchedulerMetrics { self.dispatch_count.fetch_add(1, Ordering::Relaxed); } + /// Record that a dispatch was deferred: the per-tenant queue was full. + pub fn record_dispatch_busy(&self) { + self.dispatch_busy_count.fetch_add(1, Ordering::Relaxed); + } + /// Record that a transaction was blocked on lock acquisition. pub fn record_blocked(&self) { self.blocked_count.fetch_add(1, Ordering::Relaxed); @@ -298,6 +305,7 @@ impl Default for SchedulerMetrics { fn default() -> Self { Self { dispatch_count: AtomicU64::new(0), + dispatch_busy_count: AtomicU64::new(0), blocked_count: AtomicU64::new(0), lock_wait_ms_total: AtomicU64::new(0), completed_count: AtomicU64::new(0), diff --git a/nodedb/src/control/server/resp/gateway_dispatch.rs b/nodedb/src/control/server/resp/gateway_dispatch.rs index cc33c82a5..c32b92340 100644 --- a/nodedb/src/control/server/resp/gateway_dispatch.rs +++ b/nodedb/src/control/server/resp/gateway_dispatch.rs @@ -373,7 +373,9 @@ fn gateway_payloads_to_response(payloads: Vec>) -> Response { /// which Redis clients handle with automatic retry (same as Redis Cluster BUSY). fn map_busy_error(e: crate::Error) -> crate::Error { match &e { - crate::Error::Bridge { .. } | crate::Error::Dispatch { .. } => crate::Error::Bridge { + crate::Error::Bridge { .. } + | crate::Error::Dispatch { .. } + | crate::Error::DispatchCapacityBusy { .. } => crate::Error::Bridge { detail: "BUSY NodeDB is processing requests, retry later".into(), }, _ => e, diff --git a/nodedb/src/error/types.rs b/nodedb/src/error/types.rs index 68c3c910b..371a006f9 100644 --- a/nodedb/src/error/types.rs +++ b/nodedb/src/error/types.rs @@ -378,6 +378,15 @@ pub enum Error { #[error("dispatch error: {detail}")] Dispatch { detail: String }, + /// Per-tenant dispatch capacity is saturated: the request was not enqueued, + /// so callers retry with backoff instead of reporting a terminal failure. + #[error("tenant {tenant_id}: dispatch capacity busy ({inflight}/{cap} in-flight); retry")] + DispatchCapacityBusy { + tenant_id: u64, + inflight: u32, + cap: u32, + }, + #[error("storage error ({engine}): {detail}")] Storage { engine: String, detail: String }, @@ -514,15 +523,12 @@ pub enum Error { )] SequencerUnavailable, - /// Active-session capacity reached. #[error("session cap ({cap}) exceeded — rejecting new login")] SessionCapExceeded { cap: usize }, - /// Session closed because the per-database idle timeout elapsed. #[error("session closed: idle timeout exceeded")] SessionIdleTimeout, - /// Session closed because the OIDC token expired. #[error("session closed: OIDC token expired")] SessionTokenExpired, @@ -530,29 +536,21 @@ pub enum Error { #[error("session terminated by administrator")] SessionKilledByAdmin, - /// Session closed because the associated user was dropped. #[error("session closed: user account was dropped")] SessionUserDropped, - /// OIDC bearer token rejected because the authenticated provider has no tenant binding. #[error("OIDC token rejected: authenticated provider has no tenant binding")] OidcProviderTenantUnbound, - /// OIDC provider tenant is absent or unreadable. #[error("OIDC token rejected: authenticated provider tenant is unavailable")] OidcProviderTenantUnavailable { tenant_id: u64 }, - /// OIDC bearer token rejected: claim mapping produced no default database. #[error("OIDC token rejected: claim mapping produced no default database for subject '{sub}'")] OidcNoDefaultDatabase { sub: String }, - /// Vector insert or index rejected: the vector dimension exceeds the - /// tenant's `max_vector_dim` quota. #[error("vector dimension {dim} exceeds tenant quota max_vector_dim={limit}")] TenantVectorDimExceeded { dim: u32, limit: u32 }, - /// Graph traversal rejected: the requested depth exceeds the tenant's - /// `max_graph_depth` quota. #[error("graph traversal depth {depth} exceeds tenant quota max_graph_depth={limit}")] TenantGraphDepthExceeded { depth: u32, limit: u32 }, @@ -576,7 +574,6 @@ pub enum Error { cause: super::ollp::OllpExhaustedCause, }, - /// Unpromoted mirrors are read-only. #[error("database '{database}' is a read-only mirror; promote it before writing")] MirrorReadOnly { database: String }, diff --git a/nodedb/src/error_classify.rs b/nodedb/src/error_classify.rs index afc808472..8f942f36b 100644 --- a/nodedb/src/error_classify.rs +++ b/nodedb/src/error_classify.rs @@ -198,6 +198,11 @@ pub(crate) fn classify(e: &Error) -> NodeDbError { Error::Wal(wal_err) => NodeDbError::wal(wal_err), Error::Dispatch { detail } => NodeDbError::dispatch(detail), + // Retryable capacity condition (issue 352): classified like a dispatch + // failure, and the pgwire gateway maps it to BUSY so clients retry. + Error::DispatchCapacityBusy { tenant_id, .. } => { + NodeDbError::dispatch(format!("tenant {tenant_id}: dispatch capacity busy; retry")) + } Error::Storage { detail, .. } => NodeDbError::storage(detail), Error::ColdStorage { detail } => NodeDbError::cold_storage(detail), Error::Serialization { format, detail } => {