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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion nodedb-cluster/src/bootstrap/start.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -208,6 +211,7 @@ pub async fn start_cluster_subsystems(
transport: Arc<NexarTransport>,
raft_multi_raft: Arc<std::sync::Mutex<crate::multi_raft::MultiRaft>>,
catalog: &Arc<ClusterCatalog>,
lease_liveness: Arc<crate::lease_liveness::LeaseHolderLiveness>,
) -> Result<RunningCluster> {
let health = ClusterHealth::new();
let (decommission_signal, _) = tokio::sync::watch::channel(false);
Expand All @@ -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(
Expand Down
217 changes: 217 additions & 0 deletions nodedb-cluster/src/lease_liveness.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,217 @@
// 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::BTreeMap;
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`.
///
/// 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<BTreeMap<u64, Option<u64>>>,
}

impl LeaseHolderLiveness {
pub fn new() -> Self {
Self::default()
}

/// 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, 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
/// 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_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<u64> {
self.gone
.read()
.unwrap_or_else(|poison| poison.into_inner())
.get(&node_id)
.copied()
.flatten()
}
}

/// [`MembershipSubscriber`] that feeds [`LeaseHolderLiveness`].
pub struct LeaseHolderLivenessHook {
liveness: Arc<LeaseHolderLiveness>,
resolver: NodeIdResolver,
}

impl LeaseHolderLivenessHook {
pub fn new(liveness: Arc<LeaseHolderLiveness>, resolver: NodeIdResolver) -> Self {
Self { liveness, resolver }
}
}

/// 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<RwLock<crate::topology::ClusterTopology>>,
) -> NodeIdResolver {
Arc::new(move |node_id| {
let topo = topology.read().unwrap_or_else(|poison| poison.into_inner());
node_id
.as_str()
.parse::<u64>()
.ok()
.filter(|&id| topo.get_node(id).is_some())
})
}

impl MembershipSubscriber for LeaseHolderLivenessHook {
fn on_state_change(&self, node_id: &NodeId, _old: Option<MemberState>, 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 => {}
}
}

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)]
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 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());
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));
}
}
2 changes: 2 additions & 0 deletions nodedb-cluster/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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};
Expand Down
81 changes: 81 additions & 0 deletions nodedb-cluster/src/metadata_group/cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<TopologyChange>,
pub routing_log: Vec<RoutingChange>,
Expand Down Expand Up @@ -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<u64> {
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) {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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);
}
}
17 changes: 17 additions & 0 deletions nodedb-cluster/src/metadata_group/entry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading
Loading