#![cfg_attr(not(test), deny(clippy::disallowed_methods))]
#![cfg_attr(
not(test),
deny(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
clippy::unreachable,
clippy::todo,
clippy::unimplemented,
clippy::indexing_slicing,
clippy::string_slice,
clippy::arithmetic_side_effects,
)
)]
use std::collections::BTreeSet;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use tokio_util::sync::CancellationToken;
use super::counter::CELL_SEPARATOR;
use super::membership::{ClusterState, LivenessOverlay, MemberRecord};
use super::transport::{IncomingFrames, PeerTransport};
use super::wire::{ClusterMessage, Envelope, FrameVerifier};
use super::{ClusterHandle, ClusterInner, ClusterMetrics, Incarnation, jittered, wire};
use crate::config::MIN_CLUSTER_SECRET_LEN;
use crate::entropy::Entropy;
use crate::time::{ClockSource, MonotonicInstant, clock_unix_duration};
use crate::{AutumnError, AutumnResult, cluster::NodeId};
pub const LEAVE_BUDGET: Duration = Duration::from_millis(250);
const LEAVE_FLUSH_POLL: Duration = Duration::from_millis(5);
const PUSH_FLOOR_DIVISOR: u32 = 4;
const PUSH_FLOOR_MIN: Duration = Duration::from_millis(50);
const UNSENDABLE_WARN_INTERVAL: Duration = Duration::from_secs(60);
#[derive(Clone)]
pub struct ClusterRuntimeConfig {
pub cluster_name: String,
pub secret: Vec<u8>,
pub node_id: Option<String>,
pub advertise_addr: Option<String>,
pub seed_peers: Vec<String>,
pub push_interval: Duration,
pub suspicion_timeout: Duration,
}
impl std::fmt::Debug for ClusterRuntimeConfig {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ClusterRuntimeConfig")
.field("cluster_name", &self.cluster_name)
.field("node_id", &self.node_id)
.field("advertise_addr", &self.advertise_addr)
.field("seed_peers", &self.seed_peers)
.field("push_interval", &self.push_interval)
.field("suspicion_timeout", &self.suspicion_timeout)
.finish_non_exhaustive()
}
}
pub struct ClusterNode;
impl ClusterNode {
pub fn start(
config: ClusterRuntimeConfig,
entropy: Arc<dyn Entropy>,
clock: Arc<dyn ClockSource>,
shutdown: CancellationToken,
transport: Arc<dyn PeerTransport>,
) -> AutumnResult<ClusterHandle> {
let runtime = tokio::runtime::Handle::try_current().map_err(|err| {
boot_error(format!(
"ClusterNode::start must be called from within a Tokio runtime: {err}"
))
})?;
let node_id = resolve_node_id(config.node_id.as_deref(), entropy.as_ref());
validate(&config, &node_id)?;
let incoming = transport.take_incoming().ok_or_else(|| {
boot_error("the transport's inbound frame stream has already been taken")
})?;
let local_addr = transport.local_addr();
let advertise_addr = config
.advertise_addr
.as_deref()
.map(str::trim)
.filter(|s| !s.is_empty())
.map_or_else(|| local_addr.to_string(), ToOwned::to_owned);
let incarnation = seed_incarnation(clock.as_ref());
let mut state = ClusterState::default();
state.members.insert(
node_id.clone(),
MemberRecord::alive(advertise_addr.clone(), incarnation),
);
let inner = Arc::new(ClusterInner {
node_id,
cluster_name: config.cluster_name,
local_addr,
advertise_addr,
secret: config.secret,
seed_peers: config.seed_peers,
push_interval: config.push_interval,
incarnation: AtomicU64::new(incarnation),
state: std::sync::Mutex::new(state),
overlay: std::sync::Mutex::new(LivenessOverlay::new(
config.push_interval,
config.suspicion_timeout,
)),
pruned_senders: std::sync::Mutex::new(BTreeSet::new()),
clock,
entropy,
transport,
shutdown,
notify: tokio::sync::Notify::new(),
metrics: ClusterMetrics::default(),
});
inner.transport.start(&inner.shutdown, &inner.entropy);
runtime.spawn(push_loop(Arc::clone(&inner), inner.shutdown.child_token()));
runtime.spawn(receive_loop(
Arc::clone(&inner),
incoming,
inner.shutdown.child_token(),
));
Ok(ClusterHandle::from_inner(inner))
}
}
fn boot_error(message: impl std::fmt::Display) -> AutumnError {
AutumnError::internal_server_error_msg(format!("cluster: {message}"))
}
fn validate(config: &ClusterRuntimeConfig, node_id: &str) -> AutumnResult<()> {
let secret_len = config.secret.len();
if secret_len < MIN_CLUSTER_SECRET_LEN {
return Err(boot_error(format!(
"the shared secret must be at least {MIN_CLUSTER_SECRET_LEN} bytes, got \
{secret_len}: there is no unauthenticated mode, and an absent secret must \
never fall back to an empty key"
)));
}
validate_ident("cluster_name", &config.cluster_name)?;
validate_ident("node_id", node_id)?;
if config.push_interval.is_zero() {
return Err(boot_error(
"push_interval must be greater than zero: a zero interval spins the push loop",
));
}
Ok(())
}
fn validate_ident(field: &str, value: &str) -> AutumnResult<()> {
if value.trim().is_empty() {
return Err(boot_error(format!("{field} must not be empty")));
}
if value.contains(CELL_SEPARATOR) {
return Err(boot_error(format!(
"{field} must not contain {CELL_SEPARATOR:?} ({value:?}): it separates the node id \
from the incarnation in every counter cell key"
)));
}
Ok(())
}
fn seed_incarnation(clock: &dyn ClockSource) -> Incarnation {
u64::try_from(clock_unix_duration(clock).as_millis()).unwrap_or(u64::MAX)
}
async fn push_loop(inner: Arc<ClusterInner>, shutdown: CancellationToken) {
let mut seq: u64 = 0;
let mut published = inner.incarnation.load(Ordering::Relaxed);
let floor = push_floor(inner.push_interval);
let mut last_push: Option<MonotonicInstant> = None;
loop {
let interval = jittered(inner.push_interval, inner.entropy.as_ref());
tokio::select! {
() = shutdown.cancelled() => {
depart(&inner, seq).await;
return;
}
() = inner.notify.notified() => {
let remaining = last_push
.and_then(|last| {
let since = inner.clock.monotonic().saturating_duration_since(last);
floor.checked_sub(since)
})
.filter(|remaining| !remaining.is_zero());
if let Some(remaining) = remaining {
tokio::select! {
() = tokio::time::sleep(remaining) => {}
() = shutdown.cancelled() => {
depart(&inner, seq).await;
return;
}
}
}
}
() = tokio::time::sleep(interval) => {}
}
push_round(&inner, &mut seq, &mut published);
last_push = Some(inner.clock.monotonic());
}
}
fn push_floor(push_interval: Duration) -> Duration {
push_interval
.checked_div(PUSH_FLOOR_DIVISOR)
.unwrap_or(push_interval)
.max(PUSH_FLOOR_MIN)
.min(push_interval)
}
fn push_round(inner: &Arc<ClusterInner>, seq: &mut u64, published: &mut Incarnation) {
let incarnation = inner.incarnation.load(Ordering::Relaxed);
if incarnation > *published {
*seq = 0;
*published = incarnation;
}
let now = inner.clock.monotonic();
let (document, targets) = {
let mut state = inner.lock_state();
let mut overlay = inner.lock_overlay();
let converted = state.convert_down_members(&inner.node_id, &mut overlay, now);
state.observe_tombstones(&mut overlay, now);
let pruned = state.prune_tombstones(&mut overlay, now);
drop(overlay);
if !converted.is_empty() {
tracing::debug!(
converted = converted.len(),
"cluster: members Down for a whole tombstone window recorded as Left"
);
}
if !pruned.is_empty() {
tracing::debug!(
pruned = pruned.len(),
"cluster: pruned expired Left tombstones"
);
inner.note_pruned_senders(pruned);
}
publish_self(&mut state, inner, incarnation);
let targets = push_targets(&state, inner);
let document = state.clone();
drop(state);
(document, targets)
};
inner.transport.retain_peers(&targets);
let message = ClusterMessage::StatePush { state: document };
let mut sent: u64 = 0;
for target in &targets {
if send_signed(inner, target, incarnation, *seq, &message) {
*seq = seq.saturating_add(1);
sent = sent.saturating_add(1);
}
}
inner.metrics.pushes_sent.fetch_add(sent, Ordering::Relaxed);
inner
.metrics
.frames_dropped
.store(inner.transport.dropped_frames(), Ordering::Relaxed);
inner
.metrics
.framing_rejected
.store(inner.transport.framing_rejections(), Ordering::Relaxed);
}
fn publish_self(state: &mut ClusterState, inner: &ClusterInner, incarnation: Incarnation) {
let own = MemberRecord::alive(inner.advertise_addr.clone(), incarnation);
match state.members.get_mut(&inner.node_id) {
Some(existing) if existing.incarnation <= incarnation => *existing = own,
Some(_) => {}
None => {
state.members.insert(inner.node_id.clone(), own);
}
}
}
fn push_targets(state: &ClusterState, inner: &ClusterInner) -> BTreeSet<String> {
target_addresses(
state,
&inner.node_id,
&inner.seed_peers,
&[inner.advertise_addr.as_str(), &inner.local_addr.to_string()],
)
}
fn target_addresses(
state: &ClusterState,
node_id: &str,
seeds: &[String],
own_addrs: &[&str],
) -> BTreeSet<String> {
let mut targets: BTreeSet<String> = seeds
.iter()
.map(|seed| seed.trim().to_owned())
.filter(|seed| !seed.is_empty())
.collect();
for (id, record) in &state.members {
if id == node_id || record.addr.is_empty() {
continue;
}
targets.insert(record.addr.clone());
}
for own in own_addrs {
targets.remove(*own);
}
targets
}
fn signed_frame(
inner: &ClusterInner,
incarnation: Incarnation,
seq: u64,
message: &ClusterMessage,
) -> Option<Vec<u8>> {
let Some(envelope) = wire::sign_envelope(
&inner.secret,
&inner.cluster_name,
&inner.node_id,
incarnation,
seq,
message,
) else {
note_unsendable(inner, None);
return None;
};
let Some(frame) = wire::encode_frame(&envelope) else {
note_unsendable(inner, wire::encoded_body_len(&envelope));
return None;
};
Some(frame)
}
fn send_signed(
inner: &ClusterInner,
target: &str,
incarnation: Incarnation,
seq: u64,
message: &ClusterMessage,
) -> bool {
let Some(frame) = signed_frame(inner, incarnation, seq, message) else {
return false;
};
inner.transport.send(target, frame);
true
}
fn send_farewell_signed(
inner: &ClusterInner,
target: &str,
incarnation: Incarnation,
seq: u64,
message: &ClusterMessage,
) -> bool {
let Some(frame) = signed_frame(inner, incarnation, seq, message) else {
return false;
};
inner.transport.send_farewell(target, frame)
}
fn note_unsendable(inner: &ClusterInner, serialized_bytes: Option<usize>) {
inner
.metrics
.pushes_unsendable
.fetch_add(1, Ordering::Relaxed);
let now_ms = u64::try_from(inner.clock.monotonic().since_origin().as_millis())
.unwrap_or(u64::MAX)
.saturating_add(1);
let due = inner
.metrics
.unsendable_warned_at_ms
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |last| {
let elapsed_ms =
u64::try_from(UNSENDABLE_WARN_INTERVAL.as_millis()).unwrap_or(u64::MAX);
(last == 0 || now_ms.saturating_sub(last) >= elapsed_ms).then_some(now_ms)
})
.is_ok();
if !due {
return;
}
let Some(bytes) = serialized_bytes else {
tracing::warn!(
node_id = %inner.node_id,
unsendable_total = inner.metrics.pushes_unsendable.load(Ordering::Relaxed),
"cluster: an outbound message could not be serialized at all and was \
dropped before the transport"
);
return;
};
tracing::warn!(
node_id = %inner.node_id,
serialized_bytes = bytes,
frame_cap_bytes = wire::MAX_FRAME_BYTES,
unsendable_total = inner.metrics.pushes_unsendable.load(Ordering::Relaxed),
"cluster: this node's replicated document no longer fits one frame, so \
every state push (which is also the heartbeat) is discarded before the \
transport; peers will evict this node at the suspicion timeout. Counter \
cells are never pruned — set a stable node_id and keep counter names a \
small fixed set (see docs/guide/clustering.md, State growth)"
);
}
async fn depart(inner: &Arc<ClusterInner>, mut seq: u64) {
let incarnation = inner.incarnation.load(Ordering::Relaxed);
let (document, targets) = {
let mut state = inner.lock_state();
state.members.insert(
inner.node_id.clone(),
MemberRecord::left(inner.advertise_addr.clone(), incarnation),
);
let targets = push_targets(&state, inner);
let document = state.clone();
drop(state);
(document, targets)
};
let farewell = ClusterMessage::StatePush { state: document };
let mut lost: usize = 0;
for target in &targets {
if send_farewell_signed(inner, target, incarnation, seq, &farewell) {
seq = seq.saturating_add(1);
} else {
lost = lost.saturating_add(1);
}
if send_farewell_signed(inner, target, incarnation, seq, &ClusterMessage::Leave) {
seq = seq.saturating_add(1);
}
}
inner
.metrics
.frames_dropped
.store(inner.transport.dropped_frames(), Ordering::Relaxed);
if lost > 0 {
tracing::warn!(
node_id = %inner.node_id,
peers = lost,
frames_dropped = inner.metrics.frames_dropped.load(Ordering::Relaxed),
"cluster: this node's final document could not be handed to the \
transport for one or more peers, so its departure is lost: those \
peers keep this node in view until the suspicion timeout, and any \
increment accepted since the last push round dies with this process"
);
}
let flushed = tokio::time::timeout(LEAVE_BUDGET, async {
while inner.transport.pending_frames() > 0 {
tokio::time::sleep(LEAVE_FLUSH_POLL).await;
}
})
.await;
if flushed.is_err() {
tracing::debug!(
budget_ms = LEAVE_BUDGET.as_millis(),
"cluster: departure notice not fully flushed inside its budget; \
peers converge on the suspicion timeout instead"
);
}
}
async fn receive_loop(
inner: Arc<ClusterInner>,
mut incoming: IncomingFrames,
shutdown: CancellationToken,
) {
let mut verifier = FrameVerifier::new(
inner.cluster_name.clone(),
inner.node_id.clone(),
inner.secret.clone(),
);
loop {
let received = tokio::select! {
item = incoming.recv() => item,
() = shutdown.cancelled() => None,
};
let Some((from, frame)) = received else {
return;
};
for pruned in inner.take_pruned_senders() {
verifier.forget(&pruned);
}
match verifier.accept(&frame) {
Ok((envelope, message)) => apply(&inner, &envelope, &from, &message),
Err(reason) => {
inner.metrics.record_rejection(reason);
if !reason.authenticated() {
inner.transport.note_unauthenticated_frame(&from);
}
tracing::debug!(
peer = %from,
reason = reason.label(),
closes_connection = reason.closes_connection(),
rejected_total = verifier.rejected_total(),
"cluster: inbound frame rejected"
);
}
}
}
}
fn apply(
inner: &Arc<ClusterInner>,
envelope: &Envelope,
peer_addr: &str,
message: &ClusterMessage,
) {
let now = inner.clock.monotonic();
let refuted = {
let mut state = inner.lock_state();
let outcome = match message {
ClusterMessage::StatePush { state: theirs } => {
let before = state.clone();
{
let overlay = inner.lock_overlay();
state.merge(theirs, &overlay, now);
}
if *state != before {
inner.metrics.merges_applied.fetch_add(1, Ordering::Relaxed);
}
inner
.metrics
.pushes_received
.fetch_add(1, Ordering::Relaxed);
let own = MemberRecord::alive(
inner.advertise_addr.clone(),
inner.incarnation.load(Ordering::Relaxed),
);
let bumped = state.refute(&inner.node_id, &own);
if let Some(bumped) = bumped {
inner.incarnation.store(bumped, Ordering::Relaxed);
}
bumped
}
ClusterMessage::Leave => {
apply_leave(&mut state, envelope, peer_addr);
None
}
};
drop(state);
inner.lock_overlay().record_receipt(&envelope.sender, now);
outcome
};
if let Some(bumped) = refuted {
tracing::info!(
node_id = %inner.node_id,
incarnation = bumped,
"cluster: refuting a stale record about this node at a higher incarnation"
);
inner.notify.notify_one();
}
}
fn apply_leave(state: &mut ClusterState, envelope: &Envelope, peer_addr: &str) {
let addr = state
.members
.get(&envelope.sender)
.map_or_else(|| peer_addr.to_owned(), |record| record.addr.clone());
let tombstone = MemberRecord::left(addr, envelope.incarnation);
match state.members.get_mut(&envelope.sender) {
Some(existing) => existing.merge(&tombstone),
None => {
state.members.insert(envelope.sender.clone(), tombstone);
}
}
}
pub fn resolve_node_id(configured: Option<&str>, entropy: &dyn Entropy) -> NodeId {
configured
.map(str::trim)
.filter(|s| !s.is_empty())
.map_or_else(
|| format!("node-{}", entropy.uuid_v4().simple()),
ToOwned::to_owned,
)
}
#[cfg(test)]
mod tests {
use super::{PUSH_FLOOR_MIN, push_floor, seed_incarnation, target_addresses};
use crate::cluster::membership::{ClusterState, MemberRecord};
use crate::time::FixedClock;
use std::time::Duration;
fn document(records: &[(&str, MemberRecord)]) -> ClusterState {
let mut state = ClusterState::default();
for (id, record) in records {
state.members.insert((*id).to_owned(), record.clone());
}
state
}
#[test]
fn tombstoned_members_stay_push_targets() {
let state = document(&[
("node-a", MemberRecord::alive("127.0.0.1:7001", 4)),
("node-b", MemberRecord::left("127.0.0.1:7002", 9)),
]);
let targets: Vec<String> =
target_addresses(&state, "node-a", &[], &["127.0.0.1:7001", "127.0.0.1:7001"])
.into_iter()
.collect();
assert_eq!(
targets,
vec!["127.0.0.1:7002".to_owned()],
"a tombstoned member must stay a push target until it is pruned, \
and this node must never gossip at itself; observed {targets:?}"
);
}
#[test]
fn own_addresses_and_empty_addresses_are_never_targets() {
let state = document(&[
("node-a", MemberRecord::alive("10.0.0.1:7946", 1)),
("node-c", MemberRecord::alive("", 1)),
]);
let seeds = vec![
" 127.0.0.1:7100 ".to_owned(),
String::new(),
"10.0.0.1:7946".to_owned(),
];
let targets: Vec<String> =
target_addresses(&state, "node-a", &seeds, &["10.0.0.1:7946", "127.0.0.1:9"])
.into_iter()
.collect();
assert_eq!(
targets,
vec!["127.0.0.1:7100".to_owned()],
"seeds must be trimmed, blanks and address-less members skipped, and \
every address of this node removed; observed {targets:?}"
);
}
#[test]
fn runtime_config_debug_never_prints_the_secret() {
let config = super::ClusterRuntimeConfig {
cluster_name: "orchard".to_owned(),
secret: b"a-shared-cluster-secret-value-32".to_vec(),
node_id: Some("node-a".to_owned()),
advertise_addr: Some("10.0.1.7:7946".to_owned()),
seed_peers: vec!["10.0.1.8:7946".to_owned()],
push_interval: Duration::from_millis(500),
suspicion_timeout: Duration::from_millis(2_500),
};
let rendered = format!("{config:?}");
assert!(
!rendered.contains("secret"),
"the secret field must not appear in Debug output at all; got {rendered}"
);
assert!(
!rendered.contains("97, 45, 115"),
"…and certainly not as the byte array a derived Debug would print \
(`a-s` = 97, 45, 115); got {rendered}"
);
assert!(
rendered.contains("orchard") && rendered.contains("10.0.1.7:7946"),
"the diagnosable fields must still be there, or redaction has cost \
the struct its usefulness; got {rendered}"
);
}
#[test]
fn incarnation_is_seeded_from_unix_milliseconds() {
let clock = FixedClock::at(
chrono::DateTime::<chrono::Utc>::from_timestamp(1_765_430_000, 250_000_000)
.unwrap_or_default(),
);
assert_eq!(
seed_incarnation(&clock),
1_765_430_000_250,
"the incarnation must be the clock's Unix MILLISECONDS: seconds \
granularity (or truncation of the sub-second part) would let two \
boots inside one second mint byte-identical self-records"
);
let one_ms_later = FixedClock::at(
chrono::DateTime::<chrono::Utc>::from_timestamp(1_765_430_000, 251_000_000)
.unwrap_or_default(),
);
assert!(
seed_incarnation(&one_ms_later) > seed_incarnation(&clock),
"a boot one millisecond later must come back strictly higher"
);
let pre_epoch = FixedClock::at(
chrono::DateTime::<chrono::Utc>::from_timestamp(-10, 0).unwrap_or_default(),
);
assert_eq!(
seed_incarnation(&pre_epoch),
0,
"a pre-epoch clock must clamp to 0, never wrap to a huge incarnation \
no later boot could beat"
);
}
#[test]
fn push_floor_is_bounded_at_both_ends() {
assert_eq!(
push_floor(Duration::from_millis(500)),
Duration::from_millis(125),
"the floor is a quarter of a comfortable push interval"
);
assert_eq!(
push_floor(Duration::from_millis(80)),
PUSH_FLOOR_MIN,
"a quarter of a short interval is floored, so a write-heavy handler \
cannot gossip at request rate"
);
assert_eq!(
push_floor(Duration::from_millis(10)),
Duration::from_millis(10),
"…but never above the interval itself: a prompt push must not be \
rarer than the periodic one"
);
}
}