use std::collections::BTreeSet;
use crate::db::DistributedDatabaseConfig;
use crate::sync::cluster_members::ClusterMembers;
use crate::sync::topology::SyncNodeId;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct WriteMembership {
pub total_nodes: usize,
pub send_targets: Vec<SyncNodeId>,
}
#[must_use]
pub fn resolve_membership(
config: &DistributedDatabaseConfig,
reachable: &BTreeSet<SyncNodeId>,
) -> WriteMembership {
resolve_membership_with_record(config, None, reachable)
}
#[must_use]
pub fn resolve_membership_with_record(
config: &DistributedDatabaseConfig,
record: Option<&ClusterMembers>,
reachable: &BTreeSet<SyncNodeId>,
) -> WriteMembership {
let total_nodes = record.map_or_else(|| config.nodes.len(), ClusterMembers::denominator);
let mut emitted = BTreeSet::new();
let mut send_targets = Vec::new();
for node in &config.nodes {
if node == &config.local_node {
continue;
}
if reachable.contains(node) && emitted.insert(node.clone()) {
send_targets.push(node.clone());
}
}
WriteMembership {
total_nodes,
send_targets,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::sync::consistency::{ConsistencyError, StrongConsistency, wait_for_quorum};
use std::time::Duration;
fn config(local: &str, nodes: &[&str]) -> DistributedDatabaseConfig {
DistributedDatabaseConfig {
local_node: SyncNodeId::from(local),
nodes: nodes.iter().map(|name| SyncNodeId::from(*name)).collect(),
topology: None,
sync_interval: 1,
}
}
fn reachable(names: &[&str]) -> BTreeSet<SyncNodeId> {
names.iter().map(|name| SyncNodeId::from(*name)).collect()
}
#[test]
fn total_nodes_is_full_membership_not_reachable_subset() {
let config = config("a", &["a", "b", "c"]);
let membership = resolve_membership(&config, &reachable(&["a"]));
assert_eq!(membership.total_nodes, 3, "denominator is full membership");
assert!(
membership.send_targets.is_empty(),
"no reachable peers to propose to"
);
}
#[test]
fn send_targets_are_reachable_peers_excluding_local() {
let config = config("a", &["a", "b", "c"]);
let membership = resolve_membership(&config, &reachable(&["a", "b", "c"]));
assert_eq!(membership.total_nodes, 3);
assert_eq!(
membership.send_targets,
vec![SyncNodeId::from("b"), SyncNodeId::from("c")],
"local node is never a send target; peers in config order"
);
}
#[test]
fn unknown_reachable_names_cannot_inflate_send_targets() {
let config = config("a", &["a", "b"]);
let membership = resolve_membership(&config, &reachable(&["b", "z"]));
assert_eq!(membership.total_nodes, 2);
assert_eq!(membership.send_targets, vec![SyncNodeId::from("b")]);
}
#[test]
fn minority_denominator_fences_against_self_quorum_q3() {
let config = config("c", &["a", "b", "c"]);
let membership = resolve_membership(&config, &reachable(&["c"]));
assert_eq!(membership.total_nodes, 3);
let strong = StrongConsistency::new(membership.total_nodes, Duration::from_millis(5));
let outcome = wait_for_quorum::<SyncNodeId, _>(strong, std::iter::empty());
assert!(
matches!(
outcome,
Err(ConsistencyError::QuorumTimeout { .. }
| ConsistencyError::QuorumUnavailable { .. })
),
"minority must be fenced, got {outcome:?}"
);
}
use crate::sync::cluster_members::{ClusterMember, ClusterMembers, MemberStatus};
fn record(cn: &str, epoch: u64, members: &[&str]) -> ClusterMembers {
ClusterMembers {
cluster_identity: cn.to_owned(),
config_epoch: epoch,
members: members
.iter()
.map(|name| ClusterMember::active(SyncNodeId::from(*name)))
.collect(),
}
}
#[test]
fn no_record_denominator_is_byte_identical_to_static_config() {
let config = config("a", &["a", "b", "c"]);
let reach = reachable(&["a", "b", "c"]);
let static_only = resolve_membership(&config, &reach);
let explicit_none = resolve_membership_with_record(&config, None, &reach);
assert_eq!(
static_only, explicit_none,
"None path == static config path"
);
assert_eq!(
static_only.total_nodes, 3,
"denominator = config.nodes.len()"
);
}
#[test]
fn present_record_denominator_wins_over_config_nodes_len() {
let config = config("a", &["a", "b", "c"]);
let rec = record("prod", 4, &["a", "b", "c", "d", "e"]);
let membership =
resolve_membership_with_record(&config, Some(&rec), &reachable(&["a", "b", "c"]));
assert_eq!(
membership.total_nodes, 5,
"record denominator (5) wins over config.nodes.len() (3)"
);
assert_ne!(
membership.total_nodes,
config.nodes.len(),
"must NOT fall back to static config when a record exists"
);
}
#[test]
fn single_node_genesis_record_yields_denominator_one()
-> Result<(), crate::sync::cluster_members::ClusterMembersError> {
let config = config("solo", &["solo", "ghost-b", "ghost-c"]);
let genesis = ClusterMembers::genesis("cluster-solo", SyncNodeId::from("solo"))?;
let membership =
resolve_membership_with_record(&config, Some(&genesis), &reachable(&["solo"]));
assert_eq!(membership.total_nodes, 1, "lone genesis denominator = 1");
assert_eq!(
quorum_size(membership.total_nodes),
Ok(1),
"self-quorum: one node satisfies its own quorum"
);
Ok(())
}
#[test]
fn record_does_not_change_send_targets_in_csot1() {
let config = config("a", &["a", "b", "c"]);
let rec = record("prod", 1, &["a", "b", "c", "d"]);
let with_record =
resolve_membership_with_record(&config, Some(&rec), &reachable(&["a", "b", "c"]));
assert_eq!(
with_record.send_targets,
vec![SyncNodeId::from("b"), SyncNodeId::from("c")],
"send_targets unchanged by the record"
);
}
#[test]
fn down_and_joining_members_still_count_in_csot1_denominator() {
let config = config("a", &["a", "b", "c"]);
let rec = ClusterMembers {
cluster_identity: "prod".to_owned(),
config_epoch: 2,
members: vec![
ClusterMember::active(SyncNodeId::from("a")),
ClusterMember {
node_id: SyncNodeId::from("b"),
status: MemberStatus::Joining,
},
ClusterMember {
node_id: SyncNodeId::from("c"),
status: MemberStatus::Down,
},
],
};
let membership = resolve_membership_with_record(&config, Some(&rec), &reachable(&["a"]));
assert_eq!(membership.total_nodes, 3, "all statuses count in CSOT-1");
}
use crate::sync::consistency::quorum_size;
}