use std::fmt;
use std::sync::Arc;
use hydracache::{ClusterMember, ClusterNodeId, ClusterRole, RaftMetadataSnapshot};
use serde::Serialize;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ClusterStatusRuntime {
pub ready: bool,
pub draining: bool,
}
impl ClusterStatusRuntime {
pub fn new(ready: bool, draining: bool) -> Self {
Self { ready, draining }
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum StatusSource {
Live,
Modeled,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum MemberRole {
Local,
Client,
Member,
}
impl From<ClusterRole> for MemberRole {
fn from(value: ClusterRole) -> Self {
match value {
ClusterRole::Local => Self::Local,
ClusterRole::Client => Self::Client,
ClusterRole::Member => Self::Member,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum Reachability {
Reachable,
Suspect,
Unreachable,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum ReshardPhase {
Idle,
Planning,
Moving,
Finalizing,
}
impl fmt::Display for ReshardPhase {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
let value = match self {
Self::Idle => "idle",
Self::Planning => "planning",
Self::Moving => "moving",
Self::Finalizing => "finalizing",
};
formatter.write_str(value)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct MemberStatus {
pub node_id: String,
pub role: MemberRole,
pub reachable: Reachability,
pub generation: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct ClusterStatus {
pub source: StatusSource,
pub leader: Option<String>,
pub term: u64,
pub epoch: u64,
pub quorum_ok: bool,
pub members: Vec<MemberStatus>,
pub reshard_phase: ReshardPhase,
pub draining: bool,
}
pub trait ClusterStatusProvider: fmt::Debug + Send + Sync {
fn cluster_status(&self, runtime: ClusterStatusRuntime) -> ClusterStatus;
}
#[derive(Debug, Clone, Default)]
pub struct ModeledClusterStatus;
impl ClusterStatusProvider for ModeledClusterStatus {
fn cluster_status(&self, runtime: ClusterStatusRuntime) -> ClusterStatus {
ClusterStatus {
source: StatusSource::Modeled,
leader: runtime.ready.then(|| "local".to_owned()),
term: u64::from(runtime.ready),
epoch: 0,
quorum_ok: runtime.ready && !runtime.draining,
members: Vec::new(),
reshard_phase: ReshardPhase::Idle,
draining: runtime.draining,
}
}
}
pub trait GridControlPlaneHandle: fmt::Debug + Send + Sync {
fn snapshot(&self) -> RaftMetadataSnapshot;
fn members(&self) -> Vec<ClusterMember>;
fn raft_leader_id(&self) -> Option<String>;
fn has_quorum(&self) -> bool;
fn reachability(&self, node: &ClusterNodeId) -> Reachability;
fn reshard_phase(&self) -> ReshardPhase;
fn is_draining(&self) -> bool;
}
#[derive(Debug, Clone)]
pub struct LiveClusterStatus {
grid: Arc<dyn GridControlPlaneHandle>,
}
impl LiveClusterStatus {
pub fn new(grid: Arc<dyn GridControlPlaneHandle>) -> Self {
Self { grid }
}
}
impl ClusterStatusProvider for LiveClusterStatus {
fn cluster_status(&self, runtime: ClusterStatusRuntime) -> ClusterStatus {
let snapshot = self.grid.snapshot();
let draining = runtime.draining || self.grid.is_draining();
let members = self
.grid
.members()
.into_iter()
.map(|member| MemberStatus {
node_id: member.node_id.to_string(),
role: MemberRole::from(member.role),
reachable: self.grid.reachability(&member.node_id),
generation: member.generation.value(),
})
.collect();
ClusterStatus {
source: StatusSource::Live,
leader: self.grid.raft_leader_id(),
term: snapshot.term,
epoch: snapshot.epoch.value(),
quorum_ok: runtime.ready && self.grid.has_quorum() && !draining,
members,
reshard_phase: self.grid.reshard_phase(),
draining,
}
}
}