From 1efb77ec7796b9bcc0c248d11a891687fc2a72b6 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sat, 19 Sep 2026 14:02:17 +0800 Subject: [PATCH 1/4] fix(lease): bound the descriptor-lease expiry comparison by clock skew MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The drain wait counted a holder's lease as expired by comparing the holder's stamped expires_at against this node's wall clock. A peer stamps from its own clock and its HLC is never merged in, so a fast local clock could retire a live holder's lease early — the exact window the drain exists to close. lease_expired gives the stamped deadline the same allowance the HLC ingress enforces (MAX_CLOCK_SKEW_NS): a lease is retired only once even a clock that far ahead would agree it lapsed. A deadline whose addition saturates is never retired — the filter drops holds only where the lapse is certain. --- nodedb/src/control/lease/drain_propose.rs | 81 +++++++++++++++++++++-- 1 file changed, 76 insertions(+), 5 deletions(-) diff --git a/nodedb/src/control/lease/drain_propose.rs b/nodedb/src/control/lease/drain_propose.rs index bf87d865c..2a4f2277c 100644 --- a/nodedb/src/control/lease/drain_propose.rs +++ b/nodedb/src/control/lease/drain_propose.rs @@ -187,10 +187,13 @@ fn wait_for_lease_drain( /// Expiry compares against wall time, not [`HlcClock::peek`]: `peek` stays /// frozen on a quiet cluster, which would find every lease unexpired and /// reinstate the wedge — and an idle cluster is exactly when a crashed node's -/// leases are the only ones left. `expires_at.wall_ns` is stamped from -/// `HlcClock::now()`, local wall time held monotonic; a peer's HLC is never -/// merged in, so for a lease granted elsewhere the comparison carries that -/// node's clock offset, unbounded. +/// leases are the only ones left. `expires_at.wall_ns` is stamped from the +/// holder's own `HlcClock::now()`, so a raw comparison would carry that node's +/// clock offset unbounded. [`lease_expired`] gives the deadline the same skew +/// allowance the HLC ingress enforces ([`nodedb_types::MAX_CLOCK_SKEW_NS`]): +/// a lease is retired only once even a clock that far ahead of ours would +/// agree it lapsed, so a fast local clock cannot hand a still-live holder's +/// descriptor to the DDL early. /// /// `own_holds` excludes that many local refcount units — the requester's own — /// from both the refcount safety net and this node's replicated cache entry, @@ -220,7 +223,7 @@ fn count_matching_leases( .filter(|((lid, holder), l)| { lid == id && l.version <= up_to_version - && l.expires_at.wall_ns > now_wall_ns + && !lease_expired(l.expires_at.wall_ns, now_wall_ns) && lease_holder_is_member(shared, *holder) && !(self_only && *holder == shared.node_id) }) @@ -234,6 +237,28 @@ fn count_matching_leases( } } +/// Whether a lease stamped `expires_at_wall_ns` on its holder's clock is +/// certainly past expiry at this node's `now_wall_ns`. +/// +/// The holder stamped its deadline from its own wall clock and its HLC is +/// never merged into ours, so the raw comparison would carry that node's clock +/// offset. Retiring a lease early is the unsafe direction: the drain would +/// declare the descriptor free while the holder still believes its lease runs, +/// which is the window the drain exists to close. The comparison therefore +/// waits out [`nodedb_types::MAX_CLOCK_SKEW_NS`] past the stamped deadline — +/// the same bound the HLC ingress refuses observations beyond — so only a +/// clock ahead by more than the cluster's tolerance retires a lease early. +/// +/// A deadline so far future that adding the bound saturates is never retired: +/// the filter drops holds only where the lapse is certain, matching +/// [`lease_holder_is_member`]. +fn lease_expired(expires_at_wall_ns: u64, now_wall_ns: u64) -> bool { + match expires_at_wall_ns.checked_add(nodedb_types::MAX_CLOCK_SKEW_NS) { + Some(deadline) => deadline <= now_wall_ns, + None => false, + } +} + /// Whether `node_id` is a current cluster member. Missing topology treats /// every holder as a member. fn lease_holder_is_member(shared: &SharedState, node_id: u64) -> bool { @@ -406,6 +431,52 @@ mod tests { ) } + #[test] + fn lease_expired_waits_out_the_skew_bound() { + let deadline = 1_000_000_000_000u64; + assert!(!lease_expired(deadline, deadline)); + assert!(!lease_expired( + deadline, + deadline + nodedb_types::MAX_CLOCK_SKEW_NS - 1 + )); + assert!(lease_expired( + deadline, + deadline + nodedb_types::MAX_CLOCK_SKEW_NS + )); + assert!(lease_expired( + deadline, + deadline + nodedb_types::MAX_CLOCK_SKEW_NS + 1 + )); + } + + #[test] + fn lease_expired_saturates_at_the_clock_ceiling() { + assert!(!lease_expired(u64::MAX, u64::MAX)); + assert!(lease_expired(0, u64::MAX)); + } + + /// A deadline that just passed in our frame is still inside the skew + /// window: the drain keeps waiting, because the holder's clock may be + /// behind ours and its lease is still live there. + #[tokio::test] + async fn lease_inside_the_skew_window_still_blocks_the_drain() { + let directory = tempfile::tempdir().expect("create drain count test directory"); + let wal = Arc::new( + WalManager::open_for_testing(&directory.path().join("drain-count.wal")) + .expect("open drain count test WAL"), + ); + let (dispatcher, _data_sides) = Dispatcher::new(1, 64); + let mut state = SharedState::new(dispatcher, wal).expect("construct drain count state"); + Arc::get_mut(&mut state) + .expect("single owner in test") + .cluster_topology = Some(Arc::new(std::sync::RwLock::new(topo_with(&[1])))); + let descriptor = DescriptorId::new(0, 1, DescriptorKind::Collection, "orders".to_string()); + + let just_past = nodedb_types::Hlc::new(super::super::wall_now_ns().saturating_sub(1), 0); + insert_lease(&state, &descriptor, 1, 1, just_past); + assert_eq!(count_matching_leases(&state, &descriptor, 1, 0), 1); + } + #[tokio::test] async fn non_member_lease_does_not_block_drain_count() { let directory = tempfile::tempdir().expect("create drain count test directory"); From 1ce4d520a45ca225789be4fba9738059338f1a86 Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sat, 19 Sep 2026 14:16:49 +0800 Subject: [PATCH 2/4] feat(cluster): add lease-holder liveness fed by SWIM verdicts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The lease GC releases a holder's descriptor leases on topology absence, which lags a crash: the node stays a member until the removal path runs, and its leases block DDL drains for the whole lease duration. This adds the faster signal as a lease-specific input — routing, placement, and rebalancing never read it. Dead and Left mark a holder; Alive revives it, because a verdict is refutable by a higher incarnation and a refuted node is serving again. Suspect stays a no-op: it is transient and the holder may still hold its lease. The subscriber runs on the detector task and only touches a set. --- nodedb-cluster/src/lease_liveness.rs | 138 +++++++++++++++++++++++++++ nodedb-cluster/src/lib.rs | 2 + 2 files changed, 140 insertions(+) create mode 100644 nodedb-cluster/src/lease_liveness.rs diff --git a/nodedb-cluster/src/lease_liveness.rs b/nodedb-cluster/src/lease_liveness.rs new file mode 100644 index 000000000..94ace51e0 --- /dev/null +++ b/nodedb-cluster/src/lease_liveness.rs @@ -0,0 +1,138 @@ +// SPDX-License-Identifier: BUSL-1.1 + +//! Lease-holder liveness fed by SWIM verdicts. +//! +//! The lease GC releases a holder's descriptor leases once the holder is gone. +//! Topology absence is the stable signal, but it lags a crash: the node stays +//! in the topology until the removal path runs, and until then its leases block +//! DDL drains for the whole lease duration. This module carries the faster +//! signal — SWIM's verdict. +//! +//! The set is a lease-specific input only: routing, placement, and rebalancing +//! never read it. `Dead` and `Left` mark a holder; an `Alive` transition +//! revives it, because a verdict is refutable by a higher incarnation and a +//! refuted node is serving again. `Suspect` is deliberately a no-op — it is +//! transient and a suspect holder may still hold its lease. +//! +//! Subscribers run on the detector task and must not block; marking an id in a +//! set is exactly that cost. The proposal side stays in the raft loop, where +//! the GC already proposes releases. + +use std::collections::BTreeSet; +use std::sync::{Arc, RwLock}; + +use nodedb_types::NodeId; + +use crate::routing_liveness::NodeIdResolver; +use crate::swim::member::MemberState; +use crate::swim::subscriber::MembershipSubscriber; + +/// Node ids whose SWIM verdict is `Dead` or `Left`. +/// +/// `Left` is terminal and never revives; `Dead` is cleared by an `Alive` +/// transition for the same node (a refutation at a higher incarnation). +#[derive(Debug, Default)] +pub struct LeaseHolderLiveness { + gone: RwLock>, +} + +impl LeaseHolderLiveness { + pub fn new() -> Self { + Self::default() + } + + /// Mark `node_id` as gone: its leases may be released without waiting for + /// expiry or topology removal. + pub fn mark_gone(&self, node_id: u64) { + self.gone + .write() + .unwrap_or_else(|poison| poison.into_inner()) + .insert(node_id); + } + + /// Clear the mark for `node_id`: a refuted verdict means the holder is + /// serving again and its leased descriptors are live. + pub fn revive(&self, node_id: u64) { + self.gone + .write() + .unwrap_or_else(|poison| poison.into_inner()) + .remove(&node_id); + } + + /// Whether `node_id` is currently marked gone. + pub fn is_gone(&self, node_id: u64) -> bool { + self.gone + .read() + .unwrap_or_else(|poison| poison.into_inner()) + .contains(&node_id) + } +} + +/// [`MembershipSubscriber`] that feeds [`LeaseHolderLiveness`]. +pub struct LeaseHolderLivenessHook { + liveness: Arc, + resolver: NodeIdResolver, +} + +impl LeaseHolderLivenessHook { + pub fn new(liveness: Arc, resolver: NodeIdResolver) -> Self { + Self { liveness, resolver } + } +} + +impl MembershipSubscriber for LeaseHolderLivenessHook { + fn on_state_change(&self, node_id: &NodeId, _old: Option, new: MemberState) { + let Some(numeric_id) = (self.resolver)(node_id) else { + // SWIM knows a node the numeric registry does not — a seed + // placeholder or a learner mid-join. Nothing to mark. + return; + }; + match new { + MemberState::Dead | MemberState::Left => self.liveness.mark_gone(numeric_id), + MemberState::Alive => self.liveness.revive(numeric_id), + MemberState::Suspect => {} + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn nid(id: &str) -> NodeId { + NodeId::try_new(id.to_owned()).expect("valid node id fixture") + } + + #[test] + fn dead_and_left_mark_then_alive_revives() { + let liveness = Arc::new(LeaseHolderLiveness::new()); + let hook = LeaseHolderLivenessHook::new(liveness.clone(), Arc::new(|_| Some(7))); + + assert!(!liveness.is_gone(7)); + hook.on_state_change(&nid("n7"), Some(MemberState::Alive), MemberState::Dead); + assert!(liveness.is_gone(7)); + hook.on_state_change(&nid("n7"), Some(MemberState::Dead), MemberState::Alive); + assert!(!liveness.is_gone(7)); + hook.on_state_change(&nid("n7"), Some(MemberState::Alive), MemberState::Left); + assert!(liveness.is_gone(7)); + } + + #[test] + fn suspect_is_not_a_release_signal() { + let liveness = Arc::new(LeaseHolderLiveness::new()); + let hook = LeaseHolderLivenessHook::new(liveness.clone(), Arc::new(|_| Some(3))); + + hook.on_state_change(&nid("n3"), Some(MemberState::Alive), MemberState::Suspect); + assert!(!liveness.is_gone(3)); + } + + #[test] + fn an_unresolvable_node_is_ignored() { + let liveness = Arc::new(LeaseHolderLiveness::new()); + let hook = LeaseHolderLivenessHook::new(liveness.clone(), Arc::new(|_| None)); + + hook.on_state_change(&nid("seed"), None, MemberState::Dead); + assert!(!liveness.is_gone(0)); + assert!(!liveness.is_gone(u64::MAX)); + } +} diff --git a/nodedb-cluster/src/lib.rs b/nodedb-cluster/src/lib.rs index e17427a5a..9d16fdffa 100644 --- a/nodedb-cluster/src/lib.rs +++ b/nodedb-cluster/src/lib.rs @@ -44,6 +44,7 @@ pub mod ghost; pub mod ghost_sweeper; pub mod health; pub mod install_snapshot; +pub mod lease_liveness; pub mod lifecycle; pub mod lifecycle_state; pub mod loop_metrics; @@ -100,6 +101,7 @@ pub use error::{ pub use forward::{ChunkSink, NoopPlanExecutor, PlanExecutor}; pub use ghost::{GhostStub, GhostTable}; pub use health::{HealthConfig, HealthMonitor}; +pub use lease_liveness::{LeaseHolderLiveness, LeaseHolderLivenessHook}; pub use lifecycle_state::{ClusterLifecycleState, ClusterLifecycleTracker}; pub use loop_metrics::{LoopMetrics, LoopMetricsRegistry}; pub use migration::{MigrationPhase, MigrationState}; From e49b9e7b6ce2bdb029fa0fd1513e298fd85150fd Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sat, 19 Sep 2026 14:43:57 +0800 Subject: [PATCH 3/4] feat(cluster): release a dead holder's leases without waiting for expiry The periodic lease GC released only holders absent from ClusterTopology, so a crashed holder blocked every DDL drain on its descriptors until the lease expired. The GC now also releases a holder SWIM marked Dead/Left: the liveness set is created with the node's cluster handle, registered as a SWIM subscriber at subsystem start (with the same topology-backed resolver the routing hook uses), and handed to the RaftLoop through with_lease_liveness. Suspect is not a release signal, and a refuted verdict clears the mark again. --- nodedb-cluster/src/bootstrap/start.rs | 7 ++- nodedb-cluster/src/lease_liveness.rs | 16 +++++++ nodedb-cluster/src/raft_loop/builder.rs | 13 ++++++ nodedb-cluster/src/raft_loop/lease_gc.rs | 46 +++++++++++++++---- nodedb-cluster/src/raft_loop/loop_core.rs | 7 +++ nodedb-cluster/src/subsystem/context.rs | 7 +++ nodedb/src/control/cluster/handle.rs | 5 ++ nodedb/src/control/cluster/init.rs | 1 + .../control/cluster/start_raft/loop_build.rs | 2 + 9 files changed, 94 insertions(+), 10 deletions(-) diff --git a/nodedb-cluster/src/bootstrap/start.rs b/nodedb-cluster/src/bootstrap/start.rs index 00d26b1eb..17e1e2535 100644 --- a/nodedb-cluster/src/bootstrap/start.rs +++ b/nodedb-cluster/src/bootstrap/start.rs @@ -112,7 +112,10 @@ pub fn register_default_subsystems( swim_cfg, Arc::clone(&ctx.routing), Arc::clone(&ctx.topology), - vec![], + vec![Arc::new(crate::LeaseHolderLivenessHook::new( + Arc::clone(&ctx.lease_liveness), + crate::lease_liveness::topology_resolver(Arc::clone(&ctx.topology)), + ))], ))); registry.register(Arc::new(ReachabilitySubsystem::new( @@ -208,6 +211,7 @@ pub async fn start_cluster_subsystems( transport: Arc, raft_multi_raft: Arc>, catalog: &Arc, + lease_liveness: Arc, ) -> Result { let health = ClusterHealth::new(); let (decommission_signal, _) = tokio::sync::watch::channel(false); @@ -218,6 +222,7 @@ pub async fn start_cluster_subsystems( Arc::clone(&raft_multi_raft), health, decommission_signal, + lease_liveness, ); let executor = Arc::new(MigrationExecutor::new( diff --git a/nodedb-cluster/src/lease_liveness.rs b/nodedb-cluster/src/lease_liveness.rs index 94ace51e0..a70a910dc 100644 --- a/nodedb-cluster/src/lease_liveness.rs +++ b/nodedb-cluster/src/lease_liveness.rs @@ -80,6 +80,22 @@ impl LeaseHolderLivenessHook { } } +/// Build the same SWIM-id → numeric-id resolver the routing hook uses: +/// parse the decimal node id and require it to exist in the live topology. +/// Placeholder entries (`seed:…`) resolve to `None` and are ignored. +pub fn topology_resolver( + topology: Arc>, +) -> NodeIdResolver { + Arc::new(move |node_id| { + let topo = topology.read().unwrap_or_else(|poison| poison.into_inner()); + node_id + .as_str() + .parse::() + .ok() + .filter(|&id| topo.get_node(id).is_some()) + }) +} + impl MembershipSubscriber for LeaseHolderLivenessHook { fn on_state_change(&self, node_id: &NodeId, _old: Option, new: MemberState) { let Some(numeric_id) = (self.resolver)(node_id) else { diff --git a/nodedb-cluster/src/raft_loop/builder.rs b/nodedb-cluster/src/raft_loop/builder.rs index db681ea80..335d944bd 100644 --- a/nodedb-cluster/src/raft_loop/builder.rs +++ b/nodedb-cluster/src/raft_loop/builder.rs @@ -68,6 +68,7 @@ impl RaftLoop { // fresh `Notify` has no pending permit, so this loses nothing. reconcile_notify: tokio::sync::Notify::new(), metadata_cache: self.metadata_cache, + lease_liveness: self.lease_liveness, } } @@ -275,6 +276,18 @@ impl RaftLoop { self } + /// Wire the lease-holder liveness set fed by SWIM verdicts. The host + /// passes the same `Arc` the `LeaseHolderLivenessHook` subscriber writes + /// into, so the periodic lease-GC sweep releases a holder's leases as + /// soon as its verdict is Dead/Left. + pub fn with_lease_liveness( + mut self, + liveness: Arc, + ) -> Self { + self.lease_liveness = Some(liveness); + self + } + pub fn with_tick_interval(mut self, interval: Duration) -> Self { self.tick_interval = interval; self diff --git a/nodedb-cluster/src/raft_loop/lease_gc.rs b/nodedb-cluster/src/raft_loop/lease_gc.rs index 6a497e4ff..3a3f95cd4 100644 --- a/nodedb-cluster/src/raft_loop/lease_gc.rs +++ b/nodedb-cluster/src/raft_loop/lease_gc.rs @@ -17,17 +17,20 @@ use crate::topology::ClusterTopology; use super::loop_core::{CommitApplier, RaftLoop}; -/// Pure collection: `(node_id, descriptor_ids)` for every lease holder that -/// is not in `topology`. Sorted by `node_id` for deterministic proposal -/// order. Extracted so the sweep's decision logic is unit-testable without -/// a full `RaftLoop`. +/// Pure collection: `(node_id, descriptor_ids)` for every lease holder that is +/// gone — not in `topology`, or marked gone by the SWIM-fed +/// [`LeaseHolderLiveness`](crate::lease_liveness::LeaseHolderLiveness) set. +/// Sorted by `node_id` for deterministic proposal order. Extracted so the +/// sweep's decision logic is unit-testable without a full `RaftLoop`. pub(super) fn collect_non_member_lease_releases( topology: &ClusterTopology, cache: &MetadataCache, + liveness: Option<&crate::lease_liveness::LeaseHolderLiveness>, ) -> Vec<(u64, Vec)> { let mut by_holder: HashMap> = HashMap::new(); for (id, holder) in cache.leases.keys() { - if !topology.contains(*holder) { + let gone = !topology.contains(*holder) || liveness.is_some_and(|l| l.is_gone(*holder)); + if gone { by_holder.entry(*holder).or_default().push(id.clone()); } } @@ -50,7 +53,7 @@ impl RaftLoop { } let topo = self.topology.read().unwrap_or_else(|p| p.into_inner()); let cache = cache.read().unwrap_or_else(|p| p.into_inner()); - collect_non_member_lease_releases(&topo, &cache) + collect_non_member_lease_releases(&topo, &cache, self.lease_liveness.as_deref()) }; for (node_id, descriptor_ids) in to_release { @@ -132,7 +135,7 @@ mod tests { .leases .insert((metrics.clone(), 3), lease(&metrics, 3)); - let collected = collect_non_member_lease_releases(&topo, &cache); + let collected = collect_non_member_lease_releases(&topo, &cache, None); assert_eq!(collected.len(), 2); assert_eq!(collected[0].0, 2); assert_eq!(collected[0].1, vec![orders.clone()]); @@ -153,13 +156,38 @@ mod tests { .insert((orders.clone(), holder), lease(&orders, holder)); } - assert!(collect_non_member_lease_releases(&topo, &cache).is_empty()); + assert!(collect_non_member_lease_releases(&topo, &cache, None).is_empty()); } #[test] fn gc_collects_empty_cache() { let topo = topo_with(&[1]); let cache = MetadataCache::new(); - assert!(collect_non_member_lease_releases(&topo, &cache).is_empty()); + assert!(collect_non_member_lease_releases(&topo, &cache, None).is_empty()); + } + + /// A holder SWIM marked Dead is collectible while it is still a topology + /// member; a refuted (Alive) verdict clears the mark again. + #[test] + fn gc_collects_a_dead_member_from_liveness() { + use crate::lease_liveness::LeaseHolderLiveness; + + let topo = topo_with(&[1, 2]); + let mut cache = MetadataCache::new(); + let orders = DescriptorId::new(0, 1, DescriptorKind::Collection, "orders".to_string()); + cache.leases.insert((orders.clone(), 1), lease(&orders, 1)); + cache.leases.insert((orders.clone(), 2), lease(&orders, 2)); + + let liveness = LeaseHolderLiveness::new(); + assert!(collect_non_member_lease_releases(&topo, &cache, Some(&liveness)).is_empty()); + + liveness.mark_gone(2); + let collected = collect_non_member_lease_releases(&topo, &cache, Some(&liveness)); + assert_eq!(collected.len(), 1); + assert_eq!(collected[0].0, 2); + assert_eq!(collected[0].1, vec![orders]); + + liveness.revive(2); + assert!(collect_non_member_lease_releases(&topo, &cache, Some(&liveness)).is_empty()); } } diff --git a/nodedb-cluster/src/raft_loop/loop_core.rs b/nodedb-cluster/src/raft_loop/loop_core.rs index a6ff5b419..c3756c744 100644 --- a/nodedb-cluster/src/raft_loop/loop_core.rs +++ b/nodedb-cluster/src/raft_loop/loop_core.rs @@ -291,6 +291,12 @@ pub struct RaftLoop { /// periodic lease-GC sweep can read committed lease state directly. /// `None` in cluster-only tests that don't wire it. pub(super) metadata_cache: Option>>, + + /// Optional lease-holder liveness fed by SWIM verdicts. When set, the + /// periodic lease-GC sweep releases a holder's leases as soon as SWIM + /// marks it Dead/Left, instead of waiting for topology removal or lease + /// expiry. `None` in cluster-only tests that don't wire it. + pub(super) lease_liveness: Option>, } impl RaftLoop { @@ -344,6 +350,7 @@ impl RaftLoop { tick_count: std::sync::atomic::AtomicU64::new(0), reconcile_notify: tokio::sync::Notify::new(), metadata_cache: None, + lease_liveness: None, } } } diff --git a/nodedb-cluster/src/subsystem/context.rs b/nodedb-cluster/src/subsystem/context.rs index a364a1167..fde9a84bc 100644 --- a/nodedb-cluster/src/subsystem/context.rs +++ b/nodedb-cluster/src/subsystem/context.rs @@ -55,6 +55,11 @@ pub struct BootstrapCtx { /// the receiver exposed on `RunningCluster::decommission_signal` /// from this sender via `subscribe()`. pub decommission_signal: watch::Sender, + + /// Lease-holder liveness fed by SWIM verdicts. The host creates it, + /// registers [`crate::LeaseHolderLivenessHook`] as a SWIM subscriber, + /// and hands the same `Arc` to the raft loop's lease GC. + pub lease_liveness: Arc, } impl BootstrapCtx { @@ -66,6 +71,7 @@ impl BootstrapCtx { multi_raft: Arc>, health: ClusterHealth, decommission_signal: watch::Sender, + lease_liveness: Arc, ) -> Self { Self { topology, @@ -74,6 +80,7 @@ impl BootstrapCtx { multi_raft, health, decommission_signal, + lease_liveness, } } } diff --git a/nodedb/src/control/cluster/handle.rs b/nodedb/src/control/cluster/handle.rs index 677531d49..cb1b50b86 100644 --- a/nodedb/src/control/cluster/handle.rs +++ b/nodedb/src/control/cluster/handle.rs @@ -66,4 +66,9 @@ pub struct ClusterHandle { /// calls [`nodedb_cluster::start_cluster_subsystems`] with the /// loop's shared `multi_raft` handle. pub pending_subsystems: Mutex>, + /// Lease-holder liveness fed by SWIM verdicts. Registered as a SWIM + /// subscriber at subsystem start and handed to the `RaftLoop`'s lease GC, + /// which releases a Dead/Left holder's descriptor leases without waiting + /// for topology removal or expiry. + pub lease_liveness: Arc, } diff --git a/nodedb/src/control/cluster/init.rs b/nodedb/src/control/cluster/init.rs index 059eb608f..e51535be5 100644 --- a/nodedb/src/control/cluster/init.rs +++ b/nodedb/src/control/cluster/init.rs @@ -164,6 +164,7 @@ pub async fn init_cluster_with_transport( pending_subsystems: Mutex::new(Some(crate::control::cluster::handle::PendingSubsystems { config: cluster_config, })), + lease_liveness: Arc::new(nodedb_cluster::LeaseHolderLiveness::new()), }) } diff --git a/nodedb/src/control/cluster/start_raft/loop_build.rs b/nodedb/src/control/cluster/start_raft/loop_build.rs index a1b591409..97d873d07 100644 --- a/nodedb/src/control/cluster/start_raft/loop_build.rs +++ b/nodedb/src/control/cluster/start_raft/loop_build.rs @@ -94,6 +94,7 @@ pub(super) fn build_raft_loop( .with_plan_executor(plan_executor) .with_metadata_applier(metadata_applier) .with_metadata_cache(shared.metadata_cache.clone()) + .with_lease_liveness(Arc::clone(&handle.lease_liveness)) .with_vshard_handler(vshard_handler) .with_tick_interval(tick_interval) .with_group_watchers(handle.group_watchers.clone()) @@ -173,6 +174,7 @@ pub(super) fn build_raft_loop( Arc::clone(&handle.transport), raft_loop_handle, &handle.catalog, + Arc::clone(&handle.lease_liveness), )) }) .map_err(|e| crate::Error::Config { From 2a85df48ca0f6069a5d730aa8250a289facc694a Mon Sep 17 00:00:00 2001 From: EnRaiha <15997552+EnRaiha@users.noreply.github.com> Date: Sat, 19 Sep 2026 15:08:20 +0800 Subject: [PATCH 4/4] feat(cluster): fence descriptor-lease releases by the holder's incarnation MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A SWIM Dead verdict releases a holder's leases, but a node that restarted (re-acquired at a higher incarnation) could still be fenced by a stale verdict. The fenced grant variant carries the holder's incarnation: DescriptorLeaseGrantFenced is appended last — zerompk numbers variants by position — and is proposed only once the cluster reports LEASE_FENCING_VERSION, so mixed-version clusters keep the unfenced grant and run without the fence. The cache records the stamp, the liveness hook stores the incarnation the verdict landed at, and the GC releases a lease only when its stamp is <= that incarnation. An unstamped lease (pre-fencing or mixed-version) releases as before. --- nodedb-cluster/src/lease_liveness.rs | 75 +++++++++++++++-- nodedb-cluster/src/metadata_group/cache.rs | 81 +++++++++++++++++++ nodedb-cluster/src/metadata_group/entry.rs | 17 ++++ nodedb-cluster/src/raft_loop/lease_gc.rs | 57 ++++++++++++- nodedb-cluster/src/swim/detector/runner.rs | 6 ++ nodedb-cluster/src/swim/subscriber.rs | 12 +++ nodedb/src/bootstrap/state_wiring.rs | 1 + nodedb/src/control/cluster/handle.rs | 3 + nodedb/src/control/cluster/init.rs | 11 +++ nodedb/src/control/lease/propose.rs | 12 ++- nodedb/src/control/rolling_upgrade/mod.rs | 2 +- .../src/control/rolling_upgrade/versions.rs | 6 ++ nodedb/src/control/state/fields.rs | 12 +-- nodedb/src/control/state/init.rs | 18 ++--- nodedb/src/control/state/init_prod/open.rs | 1 + 15 files changed, 290 insertions(+), 24 deletions(-) diff --git a/nodedb-cluster/src/lease_liveness.rs b/nodedb-cluster/src/lease_liveness.rs index a70a910dc..2aac237ff 100644 --- a/nodedb-cluster/src/lease_liveness.rs +++ b/nodedb-cluster/src/lease_liveness.rs @@ -18,7 +18,7 @@ //! set is exactly that cost. The proposal side stays in the raft loop, where //! the GC already proposes releases. -use std::collections::BTreeSet; +use std::collections::BTreeMap; use std::sync::{Arc, RwLock}; use nodedb_types::NodeId; @@ -29,11 +29,17 @@ use crate::swim::subscriber::MembershipSubscriber; /// Node ids whose SWIM verdict is `Dead` or `Left`. /// +/// The value is the incarnation the verdict landed at, when the detector +/// supplied one: a lease stamped at incarnation `<=` that value was granted +/// before the verdict and may be released; a lease re-acquired later (a +/// restart bumped the incarnation) is not fenced by it. `None` means gone +/// without a fencing value — release unconditionally, the phase-1 behaviour. +/// /// `Left` is terminal and never revives; `Dead` is cleared by an `Alive` /// transition for the same node (a refutation at a higher incarnation). #[derive(Debug, Default)] pub struct LeaseHolderLiveness { - gone: RwLock>, + gone: RwLock>>, } impl LeaseHolderLiveness { @@ -41,13 +47,22 @@ impl LeaseHolderLiveness { Self::default() } - /// Mark `node_id` as gone: its leases may be released without waiting for - /// expiry or topology removal. + /// Mark `node_id` as gone without a fencing value: its leases may be + /// released without waiting for expiry or topology removal. pub fn mark_gone(&self, node_id: u64) { self.gone .write() .unwrap_or_else(|poison| poison.into_inner()) - .insert(node_id); + .insert(node_id, None); + } + + /// Mark `node_id` as gone at `incarnation`: leases stamped at an + /// incarnation `<=` this value are fenced; later re-acquisitions are not. + pub fn mark_gone_at(&self, node_id: u64, incarnation: u64) { + self.gone + .write() + .unwrap_or_else(|poison| poison.into_inner()) + .insert(node_id, Some(incarnation)); } /// Clear the mark for `node_id`: a refuted verdict means the holder is @@ -64,7 +79,20 @@ impl LeaseHolderLiveness { self.gone .read() .unwrap_or_else(|poison| poison.into_inner()) - .contains(&node_id) + .contains_key(&node_id) + } + + /// The incarnation the gone verdict landed at, when one was supplied. + /// + /// Callers check [`is_gone`](Self::is_gone) first; `None` then means the + /// mark carries no fence and the release is unconditional. + pub fn fence_at(&self, node_id: u64) -> Option { + self.gone + .read() + .unwrap_or_else(|poison| poison.into_inner()) + .get(&node_id) + .copied() + .flatten() } } @@ -109,6 +137,24 @@ impl MembershipSubscriber for LeaseHolderLivenessHook { MemberState::Suspect => {} } } + + fn on_state_change_with_incarnation( + &self, + node_id: &NodeId, + new: MemberState, + incarnation: u64, + ) { + let Some(numeric_id) = (self.resolver)(node_id) else { + return; + }; + match new { + MemberState::Dead | MemberState::Left => { + self.liveness.mark_gone_at(numeric_id, incarnation) + } + MemberState::Alive => self.liveness.revive(numeric_id), + MemberState::Suspect => {} + } + } } #[cfg(test)] @@ -133,6 +179,23 @@ mod tests { assert!(liveness.is_gone(7)); } + #[test] + fn a_dead_verdict_carries_the_fencing_incarnation() { + let liveness = Arc::new(LeaseHolderLiveness::new()); + let hook = LeaseHolderLivenessHook::new(liveness.clone(), Arc::new(|_| Some(4))); + + // The detector calls both hooks per transition; the fencing value + // from the second must win over the fenceless mark from the first. + hook.on_state_change(&nid("n4"), Some(MemberState::Alive), MemberState::Dead); + hook.on_state_change_with_incarnation(&nid("n4"), MemberState::Dead, 12); + assert_eq!(liveness.fence_at(4), Some(12)); + + hook.on_state_change(&nid("n4"), Some(MemberState::Dead), MemberState::Alive); + hook.on_state_change_with_incarnation(&nid("n4"), MemberState::Alive, 13); + assert!(!liveness.is_gone(4)); + assert_eq!(liveness.fence_at(4), None); + } + #[test] fn suspect_is_not_a_release_signal() { let liveness = Arc::new(LeaseHolderLiveness::new()); diff --git a/nodedb-cluster/src/metadata_group/cache.rs b/nodedb-cluster/src/metadata_group/cache.rs index e9c9bde53..180554a63 100644 --- a/nodedb-cluster/src/metadata_group/cache.rs +++ b/nodedb-cluster/src/metadata_group/cache.rs @@ -35,6 +35,11 @@ pub struct MetadataCache { /// `(descriptor_id, node_id) -> lease`. pub leases: HashMap<(DescriptorId, u64), DescriptorLease>, + /// Fencing incarnation per `(descriptor_id, node_id)`, recorded by the + /// fenced grant variant. Absent for a grant that carried none — a + /// pre-fencing lease or a mixed-version cluster. + pub lease_incarnations: HashMap<(DescriptorId, u64), u64>, + /// Topology mutations applied so far. pub topology_log: Vec, pub routing_log: Vec, @@ -69,6 +74,13 @@ impl MetadataCache { } } + /// The fencing incarnation recorded for `(descriptor_id, holder)`, when + /// the grant carried one. `None` means the lease is unfenced: release + /// decisions fall back to the topology/verdict alone. + pub fn lease_incarnation(&self, id: &DescriptorId, holder: u64) -> Option { + self.lease_incarnations.get(&(id.clone(), holder)).copied() + } + /// Apply a committed entry. Idempotent by `applied_index`: /// entries at or below the current watermark are ignored. pub fn apply(&mut self, index: u64, entry: &MetadataEntry) { @@ -114,12 +126,27 @@ impl MetadataCache { self.leases .insert((lease.descriptor_id.clone(), lease.node_id), lease.clone()); } + MetadataEntry::DescriptorLeaseGrantFenced { + lease, + holder_incarnation, + } => { + if lease.expires_at > self.last_applied_hlc { + self.last_applied_hlc = lease.expires_at; + } + self.leases + .insert((lease.descriptor_id.clone(), lease.node_id), lease.clone()); + self.lease_incarnations.insert( + (lease.descriptor_id.clone(), lease.node_id), + *holder_incarnation, + ); + } MetadataEntry::DescriptorLeaseRelease { node_id, descriptor_ids, } => { for id in descriptor_ids { self.leases.remove(&(id.clone(), *node_id)); + self.lease_incarnations.remove(&(id.clone(), *node_id)); } } // Drain state is host-side (lives in @@ -312,3 +339,57 @@ fn apply_compensation( ); Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + use crate::metadata_group::descriptors::DescriptorKind; + + fn lease(id: &DescriptorId, holder: u64) -> DescriptorLease { + DescriptorLease { + descriptor_id: id.clone(), + version: 1, + node_id: holder, + expires_at: Hlc::new(1_000_000, 0), + } + } + + /// A fenced grant records its incarnation; the release clears both the + /// lease and the fence. + #[test] + fn fenced_grant_records_and_release_clears_the_fence() { + let mut cache = MetadataCache::new(); + let orders = DescriptorId::new(0, 1, DescriptorKind::Collection, "orders".to_string()); + + cache.apply( + 1, + &MetadataEntry::DescriptorLeaseGrantFenced { + lease: lease(&orders, 2), + holder_incarnation: 9, + }, + ); + assert!(cache.leases.contains_key(&(orders.clone(), 2))); + assert_eq!(cache.lease_incarnation(&orders, 2), Some(9)); + + cache.apply( + 2, + &MetadataEntry::DescriptorLeaseRelease { + node_id: 2, + descriptor_ids: vec![orders.clone()], + }, + ); + assert!(!cache.leases.contains_key(&(orders.clone(), 2))); + assert_eq!(cache.lease_incarnation(&orders, 2), None); + } + + /// An unfenced grant records no incarnation — release decisions fall + /// back to the topology/verdict alone. + #[test] + fn unfenced_grant_records_no_incarnation() { + let mut cache = MetadataCache::new(); + let orders = DescriptorId::new(0, 1, DescriptorKind::Collection, "orders".to_string()); + + cache.apply(1, &MetadataEntry::DescriptorLeaseGrant(lease(&orders, 2))); + assert_eq!(cache.lease_incarnation(&orders, 2), None); + } +} diff --git a/nodedb-cluster/src/metadata_group/entry.rs b/nodedb-cluster/src/metadata_group/entry.rs index bd6552176..79c95f060 100644 --- a/nodedb-cluster/src/metadata_group/entry.rs +++ b/nodedb-cluster/src/metadata_group/entry.rs @@ -350,6 +350,23 @@ pub enum MetadataEntry { spki: [u8; 32], expires_at_ms: u64, }, + + /// A descriptor lease grant stamped with the holder's incarnation. + /// + /// Appended last on purpose: zerompk derives the variant index from + /// position, so inserting anywhere earlier would renumber every following + /// variant. Proposed only once the cluster reports the fencing version; + /// mixed-version clusters keep the plain + /// [`DescriptorLeaseGrant`](Self::DescriptorLeaseGrant) and run without + /// the fence. + /// + /// The value fences the release path: a holder marked gone at incarnation + /// `N` fences only leases stamped at `<= N`, so a node that restarted (its + /// incarnation bumped) and re-acquired is not affected by a stale verdict. + DescriptorLeaseGrantFenced { + lease: DescriptorLease, + holder_incarnation: u64, + }, } /// The direction of a join-token lifecycle transition. diff --git a/nodedb-cluster/src/raft_loop/lease_gc.rs b/nodedb-cluster/src/raft_loop/lease_gc.rs index 3a3f95cd4..cfdb4c395 100644 --- a/nodedb-cluster/src/raft_loop/lease_gc.rs +++ b/nodedb-cluster/src/raft_loop/lease_gc.rs @@ -29,7 +29,9 @@ pub(super) fn collect_non_member_lease_releases( ) -> Vec<(u64, Vec)> { let mut by_holder: HashMap> = HashMap::new(); for (id, holder) in cache.leases.keys() { - let gone = !topology.contains(*holder) || liveness.is_some_and(|l| l.is_gone(*holder)); + let gone = !topology.contains(*holder) + || liveness + .is_some_and(|l| l.is_gone(*holder) && fence_allows_release(cache, id, *holder, l)); if gone { by_holder.entry(*holder).or_default().push(id.clone()); } @@ -39,6 +41,26 @@ pub(super) fn collect_non_member_lease_releases( out } +/// Whether a liveness verdict releases this lease, per the fencing rule. +/// +/// A verdict at incarnation `N` releases leases stamped at `<= N`; a lease +/// with no stamped incarnation — a pre-fencing grant, or a grant from a +/// mixed-version cluster — is released too, matching the topology-only +/// behaviour. +fn fence_allows_release( + cache: &MetadataCache, + id: &DescriptorId, + holder: u64, + liveness: &crate::lease_liveness::LeaseHolderLiveness, +) -> bool { + match liveness.fence_at(holder) { + None => true, + Some(gone_at) => cache + .lease_incarnation(id, holder) + .map_or(true, |stamped| stamped <= gone_at), + } +} + impl RaftLoop { /// On the metadata-group leader, propose `DescriptorLeaseRelease` for /// every lease whose holder is no longer in `ClusterTopology`. @@ -190,4 +212,37 @@ mod tests { liveness.revive(2); assert!(collect_non_member_lease_releases(&topo, &cache, Some(&liveness)).is_empty()); } + + /// A fencing verdict spares a lease re-acquired at a newer incarnation, + /// releases one stamped at or below it, and releases unstamped leases + /// unconditionally. + #[test] + fn a_fenced_verdict_spares_a_newer_lease() { + use crate::lease_liveness::LeaseHolderLiveness; + + let topo = topo_with(&[2]); + let liveness = LeaseHolderLiveness::new(); + liveness.mark_gone_at(2, 4); + let orders = DescriptorId::new(0, 1, DescriptorKind::Collection, "orders".to_string()); + + let mut cache = MetadataCache::new(); + cache.leases.insert((orders.clone(), 2), lease(&orders, 2)); + + // Stamped above the verdict's incarnation: a restart re-acquired. + cache.lease_incarnations.insert((orders.clone(), 2), 5); + assert!(collect_non_member_lease_releases(&topo, &cache, Some(&liveness)).is_empty()); + + // Stamped at the verdict's incarnation: released. + cache.lease_incarnations.insert((orders.clone(), 2), 4); + let collected = collect_non_member_lease_releases(&topo, &cache, Some(&liveness)); + assert_eq!(collected.len(), 1); + assert_eq!(collected[0].0, 2); + + // No stamp (pre-fencing grant): unconditional release. + cache.lease_incarnations.clear(); + assert_eq!( + collect_non_member_lease_releases(&topo, &cache, Some(&liveness)).len(), + 1 + ); + } } diff --git a/nodedb-cluster/src/swim/detector/runner.rs b/nodedb-cluster/src/swim/detector/runner.rs index 3e12cb41b..7e447801b 100644 --- a/nodedb-cluster/src/swim/detector/runner.rs +++ b/nodedb-cluster/src/swim/detector/runner.rs @@ -122,8 +122,14 @@ impl FailureDetector { None => return outcome, }; if old_state != Some(new_state) { + let incarnation = self + .membership + .get(&update.node_id) + .map(|member| member.incarnation.get()) + .unwrap_or(0); for sub in &self.subscribers { sub.on_state_change(&update.node_id, old_state, new_state); + sub.on_state_change_with_incarnation(&update.node_id, new_state, incarnation); } } outcome diff --git a/nodedb-cluster/src/swim/subscriber.rs b/nodedb-cluster/src/swim/subscriber.rs index 20ccccca4..57fb101cc 100644 --- a/nodedb-cluster/src/swim/subscriber.rs +++ b/nodedb-cluster/src/swim/subscriber.rs @@ -29,4 +29,16 @@ pub trait MembershipSubscriber: Send + Sync { /// Called after the membership list has accepted a state change /// for `node_id`. `old` is `None` on first-time insert. fn on_state_change(&self, node_id: &NodeId, old: Option, new: MemberState); + + /// Called with the incarnation the transition landed at, immediately + /// after [`on_state_change`](Self::on_state_change). Hooks that need a + /// fencing value (lease-holder liveness) implement this; the default + /// keeps every existing hook unchanged. + fn on_state_change_with_incarnation( + &self, + _node_id: &NodeId, + _new: MemberState, + _incarnation: u64, + ) { + } } diff --git a/nodedb/src/bootstrap/state_wiring.rs b/nodedb/src/bootstrap/state_wiring.rs index 0b5034c87..f570a6957 100644 --- a/nodedb/src/bootstrap/state_wiring.rs +++ b/nodedb/src/bootstrap/state_wiring.rs @@ -62,6 +62,7 @@ pub async fn wire_state( && let Some(state) = Arc::get_mut(shared) { state.node_id = handle.node_id; + state.node_incarnation = handle.node_incarnation; state.cluster_topology = Some(Arc::clone(&handle.topology)); state.cluster_routing = Some(Arc::clone(&handle.routing)); state.cluster_transport = Some(Arc::clone(&handle.transport)); diff --git a/nodedb/src/control/cluster/handle.rs b/nodedb/src/control/cluster/handle.rs index cb1b50b86..5f7d7836e 100644 --- a/nodedb/src/control/cluster/handle.rs +++ b/nodedb/src/control/cluster/handle.rs @@ -71,4 +71,7 @@ pub struct ClusterHandle { /// which releases a Dead/Left holder's descriptor leases without waiting /// for topology removal or expiry. pub lease_liveness: Arc, + /// This node's SWIM incarnation, resolved at init (the persisted value + /// bumped, or zero on a fresh node). Stamped on fenced lease grants. + pub node_incarnation: u64, } diff --git a/nodedb/src/control/cluster/init.rs b/nodedb/src/control/cluster/init.rs index e51535be5..7173836a2 100644 --- a/nodedb/src/control/cluster/init.rs +++ b/nodedb/src/control/cluster/init.rs @@ -150,6 +150,16 @@ pub async fn init_cluster_with_transport( .into_inner() .unwrap_or_else(|p| p.into_inner()); + // Mirrors the incarnation `start_cluster_subsystems` resolves for SWIM: + // the persisted value bumped, or zero on a fresh node. Resolved before + // the handle literal moves `catalog` into it. + let node_incarnation = catalog + .load_swim_incarnation() + .ok() + .flatten() + .map(|value| nodedb_cluster::Incarnation::new(value).bump().get()) + .unwrap_or(0); + Ok(ClusterHandle { transport, topology, @@ -165,6 +175,7 @@ pub async fn init_cluster_with_transport( config: cluster_config, })), lease_liveness: Arc::new(nodedb_cluster::LeaseHolderLiveness::new()), + node_incarnation, }) } diff --git a/nodedb/src/control/lease/propose.rs b/nodedb/src/control/lease/propose.rs index b37f8a419..8bc542446 100644 --- a/nodedb/src/control/lease/propose.rs +++ b/nodedb/src/control/lease/propose.rs @@ -209,7 +209,17 @@ fn refresh_lease_after_admission( // Cluster path: encode + propose + block on apply via the // shared `propose_and_wait` helper. - let entry = MetadataEntry::DescriptorLeaseGrant(lease.clone()); + let entry = if shared + .cluster_version_view() + .can_activate_feature(crate::control::rolling_upgrade::LEASE_FENCING_VERSION) + { + MetadataEntry::DescriptorLeaseGrantFenced { + lease: lease.clone(), + holder_incarnation: shared.node_incarnation, + } + } else { + MetadataEntry::DescriptorLeaseGrant(lease.clone()) + }; super::propose_and_wait(shared, &entry, "grant")?; // Re-read the cache. Under normal conditions the apply path diff --git a/nodedb/src/control/rolling_upgrade/mod.rs b/nodedb/src/control/rolling_upgrade/mod.rs index df3eb1fdc..d89ae809a 100644 --- a/nodedb/src/control/rolling_upgrade/mod.rs +++ b/nodedb/src/control/rolling_upgrade/mod.rs @@ -23,6 +23,6 @@ pub mod view; pub use versions::{ DESCRIPTOR_DRAIN_VERSION, DESCRIPTOR_VERSIONING_VERSION, DISTRIBUTED_CATALOG_VERSION, - accept_message, should_compat_mode, + LEASE_FENCING_VERSION, accept_message, should_compat_mode, }; pub use view::{ClusterVersionView, compute_from_topology}; diff --git a/nodedb/src/control/rolling_upgrade/versions.rs b/nodedb/src/control/rolling_upgrade/versions.rs index 19127d92c..d297c506a 100644 --- a/nodedb/src/control/rolling_upgrade/versions.rs +++ b/nodedb/src/control/rolling_upgrade/versions.rs @@ -62,6 +62,12 @@ pub const DESCRIPTOR_VERSIONING_VERSION: u16 = 1; /// compat-mode fallback in `drain_for_ddl`. pub const DESCRIPTOR_DRAIN_VERSION: u16 = 1; +/// Wire version that introduced the fenced descriptor-lease grant +/// (`DescriptorLeaseGrantFenced`, stamped with the holder's incarnation). +/// Mixed-version clusters below this version keep proposing the unfenced +/// grant and run without lease fencing. +pub const LEASE_FENCING_VERSION: u16 = 1; + /// Check if a message from a remote node should be accepted. /// /// Accepts only messages with the exact current wire format version. diff --git a/nodedb/src/control/state/fields.rs b/nodedb/src/control/state/fields.rs index 4d3236139..aa4bcc388 100644 --- a/nodedb/src/control/state/fields.rs +++ b/nodedb/src/control/state/fields.rs @@ -64,10 +64,8 @@ pub struct SharedState { pub usage_store: Arc, /// Quota manager (enforcement against scope quotas). pub quota_manager: crate::control::security::metering::quota::QuotaManager, - /// Usage metering configuration (enabled flag, per-operation costs). - /// No `RwLock`/atomics — there is no live-mutation DDL for it, only the - /// bounds baked into `usage_store` / `quota_manager` at construction and - /// the `.enabled` / `.operation_costs` reads at dispatch time. + /// Usage metering configuration (enabled flag, per-operation costs); no live + /// mutation, so the bounds are baked into `usage_store`/`quota_manager`. pub metering_config: crate::control::security::metering::config::MeteringConfig, /// Auth-scoped API keys (nda_ format, bound to auth_users). pub auth_api_keys: crate::control::security::auth_apikey::AuthApiKeyStore, @@ -107,6 +105,9 @@ pub struct SharedState { pub cluster_transport: Option>, /// This node's ID (0 in single-node mode). pub node_id: u64, + + /// SWIM incarnation resolved at cluster init; stamped on fenced lease grants (0 standalone). + pub node_incarnation: u64, /// Live view of the replicated metadata catalog. Falls through to legacy redb in single-node mode. pub metadata_cache: Arc>, /// Broadcasts one event per committed metadata entry to subscribers (pgwire cache, CDC, etc.). @@ -115,8 +116,7 @@ pub struct SharedState { >, /// Per-Raft-group apply watermark registry for commit-wait and drain paths. pub group_watchers: Arc, - /// Serializes this node's attempts to acquire the replicated descriptor - /// preparation lease. + /// Serializes this node's descriptor preparation lease (replicated). pub metadata_ddl_lock: Mutex<()>, /// Replicated preparation owner plus local monotonic apply time. pub metadata_ddl_owner: Mutex>, diff --git a/nodedb/src/control/state/init.rs b/nodedb/src/control/state/init.rs index 0184c9413..282a8b06a 100644 --- a/nodedb/src/control/state/init.rs +++ b/nodedb/src/control/state/init.rs @@ -35,12 +35,11 @@ impl SharedState { /// Create shared state with a pre-built credential store (for tests that need catalog). /// - /// `is_cluster` is the static, deployment-time surrogate-registry mode - /// choice — same predicate as `SharedState::open`'s `is_cluster` - /// (whether this node's caller is about to wire it into a real Raft - /// cluster), not a property of the credential store. Almost every - /// caller is a single-process fixture and passes `false`; the cluster - /// test harness passes `true`. + /// `is_cluster` is the static, deployment-time surrogate-registry mode choice — + /// same predicate as `SharedState::open`'s `is_cluster` (whether this node's + /// caller is about to wire it into a real Raft cluster), not a property of the + /// credential store. Almost every caller is a single-process fixture and passes + /// `false`; the cluster test harness passes `true`. pub fn new_with_credentials( dispatcher: Dispatcher, wal: Arc, @@ -122,9 +121,9 @@ impl SharedState { Self::new_inner(dispatcher, wal) } - /// Create shared state whose risk scorer is built from `risk_config` - /// instead of the disabled default (for tests that exercise the risk - /// gate). Production wires the same configuration from `[auth.risk]`. + /// Create shared state whose risk scorer is built from `risk_config` instead + /// of the disabled default (tests exercising the risk gate; production wires + /// the same configuration from `[auth.risk]`). pub fn new_with_risk_config( dispatcher: Dispatcher, wal: Arc, @@ -233,6 +232,7 @@ impl SharedState { cluster_routing: None, cluster_transport: None, node_id: 0, + node_incarnation: 0, metadata_cache: Arc::new(std::sync::RwLock::new(nodedb_cluster::MetadataCache::new())), catalog_change_tx: tokio::sync::broadcast::channel( crate::control::cluster::metadata_applier::CATALOG_CHANNEL_CAPACITY, diff --git a/nodedb/src/control/state/init_prod/open.rs b/nodedb/src/control/state/init_prod/open.rs index 938a43cb0..e04e44324 100644 --- a/nodedb/src/control/state/init_prod/open.rs +++ b/nodedb/src/control/state/init_prod/open.rs @@ -171,6 +171,7 @@ impl SharedState { dispatcher: Mutex::new(dispatcher), tracker: RequestTracker::new(), wal, + node_incarnation: 0, quiesce, http_client, credentials: Arc::clone(&credentials),