From ff5868cc0c7ce943d3b406ca3b94939f32f9a70e Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 20 Sep 2026 13:52:12 +0800 Subject: [PATCH 1/8] fix(cluster): treat a full dispatch queue as retryable capacity, not a failure A full per-tenant bridge queue returned Error::Dispatch, which the Calvin scheduler logged at ERROR and completed as a failed transaction, releasing its locks. During the startup rebuild the epoch is re-driven, so a full queue spun: one ERROR and one failed transaction per attempt, the queue never drained, CPU pinned, and pgwire starved. - add Error::DispatchCapacityBusy { tenant_id, inflight, cap }: the request was not enqueued and nothing was applied, so it is retryable by type - the bridge dispatcher returns it for the per-tenant cap - the scheduler demotes the log to debug, counts it (dispatch_busy_count), and paces the catch-up drain with a bounded jittered backoff (busy_backoff, 5 ms doubling to a 250 ms cap) instead of spinning - classify the variant like a dispatch failure and map it to BUSY for clients Tests: full_queue_is_a_retryable_capacity_condition (new; failed before the fix), full_queue_returns_error (unchanged), busy_backoff_is_bounded_and_jittered. Issue 352. --- nodedb/src/bridge/dispatch/dispatcher.rs | 34 ++++++++-- .../driver/core/dispatch/static_dispatch.rs | 29 ++++++++- .../calvin/scheduler/driver/core/scheduler.rs | 64 +++++++++++++++++-- .../cluster/calvin/scheduler/metrics.rs | 8 +++ .../control/server/resp/gateway_dispatch.rs | 4 +- nodedb/src/error/types.rs | 12 ++++ nodedb/src/error_classify.rs | 5 ++ 7 files changed, 143 insertions(+), 13 deletions(-) diff --git a/nodedb/src/bridge/dispatch/dispatcher.rs b/nodedb/src/bridge/dispatch/dispatcher.rs index 7e1783f3e..2e7501e87 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 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..f06c40170 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,31 @@ impl Scheduler { }; if let Err(e) = dispatch_result { + if matches!(e, crate::Error::DispatchCapacityBusy { .. }) { + // Retryable capacity condition (issue 352): nothing was + // enqueued. Demote the log, count it, 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.metrics.record_dispatch_busy(); + let attempts = self.dispatch_busy_attempts.saturating_add(1); + self.dispatch_busy_attempts = attempts; + self.dispatch_busy_until = Some( + Instant::now() + + super::scheduler::Scheduler::busy_backoff( + attempts, + self.vshard_id, + epoch, + ), + ); + self.on_txn_complete(txn_id); + return; + } error!( vshard_id = self.vshard_id, epoch, @@ -283,6 +308,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/scheduler.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs index 40aecbc67..16d72ef4f 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::{Duration, 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, @@ -259,6 +268,22 @@ impl Scheduler { fully_applied >= self.rebuild_target_epoch } + /// Bounded, deterministic backoff for a capacity-busy dispatch (issue 352). + /// + /// 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) + } + /// Publish an advanced fully-applied watermark to the metrics gauge and the /// shared cross-shard snapshot anchor. /// @@ -371,11 +396,24 @@ 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. + 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 @@ -565,4 +603,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/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..fe66cf8c9 100644 --- a/nodedb/src/error/types.rs +++ b/nodedb/src/error/types.rs @@ -378,6 +378,18 @@ pub enum Error { #[error("dispatch error: {detail}")] Dispatch { detail: String }, + /// Per-tenant dispatch capacity is saturated: the request was NOT enqueued + /// and nothing was applied. The identical request is expected to succeed + /// once in-flight work drains, so callers must retry with backoff — never + /// treat this as a terminal dispatch failure (issue 352: a terminal read + /// turns the startup rebuild into an error storm). + #[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 }, 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 } => { From 15825ae5b39285386fef5c22e9398907baabdb21 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 20 Sep 2026 14:26:40 +0800 Subject: [PATCH 2/8] fix(cluster): handle a capacity-busy dispatch at every dispatch site The first change classified a full per-tenant queue as retryable and paced the epoch dispatch, but four other sites still logged ERROR and completed the work as failed: the CalvinResolve dispatch, the commit-resolution dispatch, the write-version record dispatch, and the second (active) epoch dispatch. On a large WAL replay those sites kept the catch-up loop busy and the readiness gate never opened. - note_dispatch_busy(&mut self, ..): counts the event and arms the bounded drain backoff; used by the epoch, CalvinResolve and commit-resolution sites - note_dispatch_busy_shared(&self, ..): counts only (no arming) for the one-way write-version record dispatch, which holds only a shared borrow - each site demotes the log to debug on the retryable class Issue 352. --- .../scheduler/driver/core/commit_redo.rs | 12 +++++++ .../driver/core/commit_resolution_dispatch.rs | 12 +++++++ .../calvin/scheduler/driver/core/scheduler.rs | 32 +++++++++++++++++++ .../driver/core/write_version_record.rs | 11 +++++++ 4 files changed, 67 insertions(+) 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/scheduler.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs index 16d72ef4f..f772eb80b 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs @@ -268,6 +268,38 @@ impl Scheduler { fully_applied >= self.rebuild_target_epoch } + /// 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 (issue 352). + 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 = + 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). Issue 352. + 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 (issue 352). /// /// 5 ms doubling on consecutive busy outcomes, capped at 250 ms, plus a 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, From 82e9d9f50f151e23adc05301d527280ed2c6e48a Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 20 Sep 2026 15:53:33 +0800 Subject: [PATCH 3/8] fix(bootstrap): allow a larger metadata apply backlog before failing startup The readiness gate failed startup when the metadata group applied no entry for 30 s. With a large apply backlog the group needs minutes, so a slow boot became a restart loop (2026-09-20 incident; issue 352). Raise the bound to 300 s and keep failing a group that never applies. Follow-up: expose the bound as configuration. --- nodedb/src/bootstrap/cluster_ready.rs | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) 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); From a71ec8eb884e581371a24ca193b9bc6261eaad9b Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 20 Sep 2026 16:40:58 +0800 Subject: [PATCH 4/8] fix(bootstrap): restore the 600 s data-group recovery bound The configurable [server] data_group_recovery_timeout_ms (600 s in production) was replaced by a hard-coded 60 s. On a backlogged recovery the data groups need longer: group 2 reached 261 of 262 committed entries inside 60 s and startup aborted into a restart loop (2026-09-20 deploy attempts of the HEAD+fixes builds). Restore the previous operator value as the constant; a follow-up should expose it as configuration again. --- nodedb/src/bootstrap/data_group_recovery.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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. From 699debbee75d30969fc594fd56a46adb6180878f Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 20 Sep 2026 17:22:10 +0800 Subject: [PATCH 5/8] chore(bridge): drop an issue reference from a touched comment The house preflight requires the issue number to live in the PR, not in the code, and this branch touches the file. --- nodedb/src/bridge/dispatch/dispatcher.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/nodedb/src/bridge/dispatch/dispatcher.rs b/nodedb/src/bridge/dispatch/dispatcher.rs index 2e7501e87..6f531ce38 100644 --- a/nodedb/src/bridge/dispatch/dispatcher.rs +++ b/nodedb/src/bridge/dispatch/dispatcher.rs @@ -661,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 From 0477305e8564c92dea6871b4416bcf2f1d05b3b4 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 20 Sep 2026 17:26:23 +0800 Subject: [PATCH 6/8] chore(cluster): mark the backoff clock reads for the determinism gate MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The Calvin determinism gate forbids unmarked Instant::now() in the write path. The backoff deadline and the drain gate read the wall clock for pacing only — scheduler observability, never Calvin WAL data. Comment-only change. --- .../calvin/scheduler/driver/core/dispatch/static_dispatch.rs | 1 + .../control/cluster/calvin/scheduler/driver/core/scheduler.rs | 2 ++ 2 files changed, 3 insertions(+) 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 f06c40170..d014db1d0 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 @@ -288,6 +288,7 @@ impl Scheduler { let attempts = self.dispatch_busy_attempts.saturating_add(1); self.dispatch_busy_attempts = attempts; self.dispatch_busy_until = Some( + // no-determinism: the backoff deadline is scheduler observability, not Calvin WAL data Instant::now() + super::scheduler::Scheduler::busy_backoff( attempts, 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 f772eb80b..4ad14b3aa 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs @@ -283,6 +283,7 @@ impl Scheduler { 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 } @@ -431,6 +432,7 @@ impl Scheduler { // 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!( From 2bbc44b128bb80073ea1c98c67a6c182ec57d022 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 20 Sep 2026 18:29:48 +0800 Subject: [PATCH 7/8] refactor(cluster): split the calvin busy helpers and active dispatch into modules The scheduler and dispatch files grew past the 500-line house limit. Move the capacity-busy accounting and backoff into core/busy.rs and the active dependent-read dispatch into core/active_dispatch.rs, and route the static dispatch busy path through the shared helper instead of duplicating it. --- .../calvin/scheduler/driver/core/busy.rs | 58 +++++++++++++++++++ .../driver/core/dispatch/active_dispatch.rs | 13 ++++- .../driver/core/dispatch/static_dispatch.rs | 20 ++----- .../calvin/scheduler/driver/core/mod.rs | 2 + .../calvin/scheduler/driver/core/scheduler.rs | 52 +---------------- 5 files changed, 78 insertions(+), 67 deletions(-) create mode 100644 nodedb/src/control/cluster/calvin/scheduler/driver/core/busy.rs 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/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 d014db1d0..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 @@ -273,10 +273,10 @@ impl Scheduler { }; if let Err(e) = dispatch_result { - if matches!(e, crate::Error::DispatchCapacityBusy { .. }) { - // Retryable capacity condition (issue 352): nothing was - // enqueued. Demote the log, count it, and pace the re-drive via - // the scheduler's drain gate instead of spinning on ERROR. + 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, @@ -284,18 +284,6 @@ impl Scheduler { error = %e, "calvin scheduler: dispatch deferred (capacity busy)" ); - self.metrics.record_dispatch_busy(); - let attempts = self.dispatch_busy_attempts.saturating_add(1); - self.dispatch_busy_attempts = attempts; - self.dispatch_busy_until = Some( - // no-determinism: the backoff deadline is scheduler observability, not Calvin WAL data - Instant::now() - + super::scheduler::Scheduler::busy_backoff( - attempts, - self.vshard_id, - epoch, - ), - ); self.on_txn_complete(txn_id); return; } 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 4ad14b3aa..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,7 +4,7 @@ use std::collections::BTreeMap; use std::sync::{Arc, Mutex}; -use std::time::{Duration, Instant}; +use std::time::Instant; use tokio::sync::mpsc; use tracing::{debug, info}; @@ -268,55 +268,6 @@ impl Scheduler { fully_applied >= self.rebuild_target_epoch } - /// 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 (issue 352). - 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). Issue 352. - 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 (issue 352). - /// - /// 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) - } - /// Publish an advanced fully-applied watermark to the metrics gauge and the /// shared cross-shard snapshot anchor. /// @@ -510,6 +461,7 @@ impl Scheduler { mod tests { use super::*; use std::collections::{BTreeSet, HashMap}; + use std::time::Duration; use nodedb_cluster::RoutingTable; From cbae0482ab149738c52afeace1fafd2188f55b30 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sun, 20 Sep 2026 18:29:50 +0800 Subject: [PATCH 8/8] refactor(error): drop variant docs that restate the display message Each removed doc line repeated the variant's #[error] text, so the message is now the single source for what the variant says. The enum had also grown past its base line count; this returns it below the unchanged-file bound. --- nodedb/src/error/types.rs | 19 ++----------------- 1 file changed, 2 insertions(+), 17 deletions(-) diff --git a/nodedb/src/error/types.rs b/nodedb/src/error/types.rs index fe66cf8c9..371a006f9 100644 --- a/nodedb/src/error/types.rs +++ b/nodedb/src/error/types.rs @@ -378,11 +378,8 @@ pub enum Error { #[error("dispatch error: {detail}")] Dispatch { detail: String }, - /// Per-tenant dispatch capacity is saturated: the request was NOT enqueued - /// and nothing was applied. The identical request is expected to succeed - /// once in-flight work drains, so callers must retry with backoff — never - /// treat this as a terminal dispatch failure (issue 352: a terminal read - /// turns the startup rebuild into an error storm). + /// 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, @@ -526,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, @@ -542,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 }, @@ -588,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 },