use crate::sync::topology::SyncNodeId;
pub const CLUSTER_MEMBERS_KEY: &[u8] = b"\0cluster/members";
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub enum MemberStatus {
Joining,
Active,
Leaving,
Down,
}
impl MemberStatus {
#[must_use]
pub const fn counts_toward_denominator(self) -> bool {
matches!(
self,
Self::Joining | Self::Active | Self::Leaving | Self::Down
)
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct ClusterMember {
pub node_id: SyncNodeId,
pub status: MemberStatus,
}
impl ClusterMember {
#[must_use]
pub const fn active(node_id: SyncNodeId) -> Self {
Self {
node_id,
status: MemberStatus::Active,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ClusterMembersError {
Encode(String),
Decode(String),
EmptyMemberSet,
DuplicateMember(String),
EmptyClusterIdentity,
}
impl std::fmt::Display for ClusterMembersError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Encode(reason) => write!(formatter, "cluster/members encode failed: {reason}"),
Self::Decode(reason) => write!(formatter, "cluster/members decode failed: {reason}"),
Self::EmptyMemberSet => formatter.write_str("cluster/members has an empty member set"),
Self::DuplicateMember(node) => {
write!(formatter, "cluster/members has duplicate member: {node}")
}
Self::EmptyClusterIdentity => {
formatter.write_str("cluster/members has an empty cluster identity (cn)")
}
}
}
}
impl std::error::Error for ClusterMembersError {}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct ClusterMembers {
#[serde(rename = "cn")]
pub cluster_identity: String,
pub config_epoch: u64,
pub members: Vec<ClusterMember>,
}
impl ClusterMembers {
pub fn genesis(
cluster_identity: impl Into<String>,
node_id: SyncNodeId,
) -> Result<Self, ClusterMembersError> {
let record = Self {
cluster_identity: cluster_identity.into(),
config_epoch: 0,
members: vec![ClusterMember::active(node_id)],
};
record.validate()?;
Ok(record)
}
#[must_use]
pub fn denominator(&self) -> usize {
self.members
.iter()
.filter(|member| member.status.counts_toward_denominator())
.count()
}
pub fn validate(&self) -> Result<(), ClusterMembersError> {
if self.cluster_identity.is_empty() {
return Err(ClusterMembersError::EmptyClusterIdentity);
}
if self.members.is_empty() {
return Err(ClusterMembersError::EmptyMemberSet);
}
let mut seen = std::collections::BTreeSet::new();
for member in &self.members {
if !seen.insert(member.node_id.clone()) {
return Err(ClusterMembersError::DuplicateMember(
member.node_id.as_str().to_owned(),
));
}
}
Ok(())
}
pub fn encode(&self) -> Result<Vec<u8>, ClusterMembersError> {
serde_json::to_vec(self).map_err(|error| ClusterMembersError::Encode(error.to_string()))
}
pub fn decode(bytes: &[u8]) -> Result<Self, ClusterMembersError> {
let record: Self = serde_json::from_slice(bytes)
.map_err(|error| ClusterMembersError::Decode(error.to_string()))?;
record.validate()?;
Ok(record)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::error::Error;
fn node(name: &str) -> SyncNodeId {
SyncNodeId::from(name)
}
#[test]
fn genesis_is_single_active_member_at_epoch_zero() -> Result<(), Box<dyn Error>> {
let record = ClusterMembers::genesis("cluster-x", node("a"))?;
assert_eq!(record.config_epoch, 0);
assert_eq!(record.denominator(), 1, "lone node self-quorum denominator");
assert_eq!(record.members, vec![ClusterMember::active(node("a"))]);
assert_eq!(record.cluster_identity, "cluster-x");
Ok(())
}
#[test]
fn genesis_rejects_empty_cluster_identity() {
assert_eq!(
ClusterMembers::genesis("", node("a")),
Err(ClusterMembersError::EmptyClusterIdentity)
);
}
#[test]
fn encode_decode_round_trips_a_multi_member_record() -> Result<(), Box<dyn Error>> {
let record = ClusterMembers {
cluster_identity: "prod".to_owned(),
config_epoch: 7,
members: vec![
ClusterMember::active(node("a")),
ClusterMember {
node_id: node("b"),
status: MemberStatus::Joining,
},
ClusterMember {
node_id: node("c"),
status: MemberStatus::Down,
},
],
};
let bytes = record.encode()?;
let decoded = ClusterMembers::decode(&bytes)?;
assert_eq!(decoded, record, "byte round-trip is lossless");
assert_eq!(decoded.denominator(), 3, "all statuses count in CSOT-1");
Ok(())
}
#[test]
fn decode_rejects_garbage_bytes() {
let error = ClusterMembers::decode(b"not json");
assert!(matches!(error, Err(ClusterMembersError::Decode(_))));
}
#[test]
fn decode_rejects_empty_member_set() -> Result<(), Box<dyn Error>> {
let bytes = serde_json::to_vec(&serde_json::json!({
"cn": "prod",
"config_epoch": 0,
"members": [],
}))?;
assert_eq!(
ClusterMembers::decode(&bytes),
Err(ClusterMembersError::EmptyMemberSet)
);
Ok(())
}
#[test]
fn validate_rejects_duplicate_members() {
let record = ClusterMembers {
cluster_identity: "prod".to_owned(),
config_epoch: 0,
members: vec![
ClusterMember::active(node("dup")),
ClusterMember::active(node("dup")),
],
};
assert!(matches!(
record.validate(),
Err(ClusterMembersError::DuplicateMember(node)) if node == "dup"
));
}
#[test]
fn cluster_identity_serialises_as_cn() -> Result<(), Box<dyn Error>> {
let record = ClusterMembers::genesis("ident", node("a"))?;
let json: serde_json::Value = serde_json::from_slice(&record.encode()?)?;
assert_eq!(
json.get("cn").and_then(serde_json::Value::as_str),
Some("ident")
);
Ok(())
}
}