use alloc::boxed::Box;
use alloc::string::String;
use alloc::sync::Arc;
use alloc::vec::Vec;
use core::fmt;
use core::future::Future;
use core::pin::Pin;
use core::time::Duration;
use super::effect::TabulaDirigendi;
use super::node::{InformationesNodi, InscriptioNodi, NodusIdentitas, StatusNodi};
pub const DEFAULT_CLUSTER_NAME: &str = "ordofp-cluster";
#[derive(Debug, Clone)]
pub struct ConfiguratioGregis {
pub nomen: String,
pub inventio: MethodusInventionis,
pub intervallum_pulsationis: Duration,
pub mora_nodi: Duration,
pub nodi_maximi: Option<usize>,
pub factor_replicationis: u8,
pub protocollum_consensus: ProtocollumConsensus,
pub tls: Option<ConfiguratioTls>,
}
impl Default for ConfiguratioGregis {
fn default() -> Self {
ConfiguratioGregis {
nomen: String::from(DEFAULT_CLUSTER_NAME),
inventio: MethodusInventionis::Staticus(Vec::new()),
intervallum_pulsationis: Duration::from_secs(5),
mora_nodi: Duration::from_secs(30),
nodi_maximi: None,
factor_replicationis: 1,
protocollum_consensus: ProtocollumConsensus::Raft,
tls: None,
}
}
}
impl ConfiguratioGregis {
pub fn new(nomen: impl Into<String>) -> Self {
ConfiguratioGregis {
nomen: nomen.into(),
..Default::default()
}
}
pub fn with_discovery(mut self, inventio: MethodusInventionis) -> Self {
self.inventio = inventio;
self
}
pub fn with_heartbeat(mut self, interval: Duration) -> Self {
self.intervallum_pulsationis = interval;
self
}
pub fn with_timeout(mut self, timeout: Duration) -> Self {
self.mora_nodi = timeout;
self
}
pub fn with_replication(mut self, factor: u8) -> Self {
self.factor_replicationis = factor;
self
}
pub fn with_consensus(mut self, protocol: ProtocollumConsensus) -> Self {
self.protocollum_consensus = protocol;
self
}
}
#[derive(Debug, Clone)]
pub enum MethodusInventionis {
Staticus(Vec<InscriptioNodi>),
Dns {
nomen: String,
portus: u16,
intervallum: Duration,
},
Kubernetes {
spatium: String,
servitium: String,
selector: Option<String>,
},
Consul {
inscriptio: InscriptioNodi,
servitium: String,
centrum: Option<String>,
},
Etcd {
endpoints: Vec<InscriptioNodi>,
praefixum: String,
},
Proprius {
nomen: String,
configuratio: String,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ProtocollumConsensus {
#[default]
Raft,
Paxos,
Zab,
Pbft,
Nullus,
}
#[derive(Debug, Clone)]
pub struct ConfiguratioTls {
pub certificatum: String,
pub clavis: String,
pub ca: Option<String>,
pub verifica_clientes: bool,
}
#[derive(Debug, Clone)]
pub struct StatusGregis {
pub configuratio: ConfiguratioGregis,
pub nodi: Vec<InformationesNodi>,
pub dux: Option<NodusIdentitas>,
pub salus: SalusGregis,
pub tabula_dirigendi: TabulaDirigendi,
pub generatio: u64,
}
impl StatusGregis {
#[inline]
pub fn new(configuratio: ConfiguratioGregis) -> Self {
StatusGregis {
configuratio,
nodi: Vec::with_capacity(16),
dux: None,
salus: SalusGregis::Unknown,
tabula_dirigendi: TabulaDirigendi::new(),
generatio: 0,
}
}
#[inline]
pub fn healthy_nodes(&self) -> impl Iterator<Item = &InformationesNodi> {
self.nodi.iter().filter(|n| n.is_healthy())
}
#[inline]
pub fn executor_nodes(&self) -> impl Iterator<Item = &InformationesNodi> {
self.nodi.iter().filter(|n| n.can_execute())
}
pub fn get_node(&self, id: NodusIdentitas) -> Option<&InformationesNodi> {
self.nodi.iter().find(|n| n.identitas == id)
}
pub fn get_node_mut(&mut self, id: NodusIdentitas) -> Option<&mut InformationesNodi> {
self.nodi.iter_mut().find(|n| n.identitas == id)
}
pub fn add_node(&mut self, node: InformationesNodi) {
self.tabula_dirigendi.remove_node(node.identitas);
for effect_id in &node.facultates.effectus_tractati {
self.tabula_dirigendi.add_route(*effect_id, node.identitas);
}
if let Some(existing) = self.nodi.iter_mut().find(|n| n.identitas == node.identitas) {
*existing = node;
} else {
self.nodi.push(node);
}
self.update_health();
}
pub fn remove_node(&mut self, id: NodusIdentitas) {
self.tabula_dirigendi.remove_node(id);
self.nodi.retain(|n| n.identitas != id);
self.update_health();
}
pub fn update_node_status(&mut self, id: NodusIdentitas, status: StatusNodi) {
if let Some(node) = self.get_node_mut(id) {
node.status = status;
self.update_health();
}
}
fn update_health(&mut self) {
let total = self.nodi.len();
let healthy = self.nodi.iter().filter(|n| n.is_healthy()).count();
self.salus = if total == 0 {
SalusGregis::Unknown
} else if healthy == total {
SalusGregis::Sanus
} else if healthy > total / 2 {
SalusGregis::Degradatus
} else if healthy > 0 {
SalusGregis::Criticus
} else {
SalusGregis::Mortuus
};
}
pub fn has_quorum(&self) -> bool {
let total = self.nodi.len();
let healthy = self.nodi.iter().filter(|n| n.is_healthy()).count();
healthy > total / 2
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum SalusGregis {
#[default]
Unknown,
Sanus,
Degradatus,
Criticus,
Mortuus,
}
impl fmt::Display for SalusGregis {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
SalusGregis::Unknown => write!(f, "unknown"),
SalusGregis::Sanus => write!(f, "healthy"),
SalusGregis::Degradatus => write!(f, "degraded"),
SalusGregis::Criticus => write!(f, "critical"),
SalusGregis::Mortuus => write!(f, "dead"),
}
}
}
pub struct AdministratorGregis {
pub nodus_localis: NodusIdentitas,
pub status: StatusGregis,
pub inventor: Option<Arc<dyn InventorNodorum>>,
}
impl AdministratorGregis {
pub fn new(nodus_localis: NodusIdentitas, configuratio: ConfiguratioGregis) -> Self {
AdministratorGregis {
nodus_localis,
status: StatusGregis::new(configuratio),
inventor: None,
}
}
pub fn with_discovery(mut self, inventor: Arc<dyn InventorNodorum>) -> Self {
self.inventor = Some(inventor);
self
}
pub fn is_leader(&self) -> bool {
self.status.dux == Some(self.nodus_localis)
}
pub fn leader(&self) -> Option<&InformationesNodi> {
self.status.dux.and_then(|id| self.status.get_node(id))
}
pub fn join(
&mut self,
info: InformationesNodi,
) -> Pin<Box<dyn Future<Output = Result<(), ErrorGregis>> + Send + '_>> {
Box::pin(async move {
self.status.add_node(info);
self.status.generatio += 1;
Ok(())
})
}
pub fn leave(&mut self) -> Pin<Box<dyn Future<Output = Result<(), ErrorGregis>> + Send + '_>> {
Box::pin(async move {
self.status.remove_node(self.nodus_localis);
self.status.generatio += 1;
Ok(())
})
}
pub fn refresh(
&mut self,
) -> Pin<Box<dyn Future<Output = Result<(), ErrorGregis>> + Send + '_>> {
Box::pin(async move {
if let Some(inventor) = &self.inventor {
let nodes = inventor.discover().await?;
for node in nodes {
self.status.add_node(node);
}
}
Ok(())
})
}
pub fn select_nodes(&self, required_effects: &[u64], count: usize) -> Vec<&InformationesNodi> {
let total = self.status.nodi.len();
let mut candidates: Vec<_> = Vec::with_capacity(total);
candidates.extend(self.status.executor_nodes().filter(|n| {
required_effects
.iter()
.all(|e| n.facultates.can_handle_effect(*e))
}));
candidates.sort_by(|a, b| {
a.facultates
.load_factor()
.partial_cmp(&b.facultates.load_factor())
.unwrap_or(core::cmp::Ordering::Equal)
});
candidates.truncate(count);
candidates
}
}
impl fmt::Debug for AdministratorGregis {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("AdministratorGregis")
.field("nodus_localis", &self.nodus_localis)
.field("status", &self.status)
.finish_non_exhaustive()
}
}
pub trait InventorNodorum: Send + Sync {
fn discover(
&self,
) -> Pin<Box<dyn Future<Output = Result<Vec<InformationesNodi>, ErrorGregis>> + Send + '_>>;
fn register(
&self,
info: &InformationesNodi,
) -> Pin<Box<dyn Future<Output = Result<(), ErrorGregis>> + Send + '_>>;
fn deregister(
&self,
id: NodusIdentitas,
) -> Pin<Box<dyn Future<Output = Result<(), ErrorGregis>> + Send + '_>>;
fn watch(
&self,
) -> Pin<Box<dyn Future<Output = Result<EventusGregis, ErrorGregis>> + Send + '_>>;
}
#[derive(Debug, Clone)]
pub enum EventusGregis {
NodusAddidit(Box<InformationesNodi>),
NodusAbscessit(NodusIdentitas),
StatusMutatus {
nodus: NodusIdentitas,
status: StatusNodi,
},
DuxMutatus(Option<NodusIdentitas>),
}
#[derive(Debug, Clone)]
pub enum ErrorGregis {
Inventio(String),
Consensus(String),
Rete(String),
Configuratio(String),
NodusNonInventus(NodusIdentitas),
SineQuorum,
IamConiunctus,
NonConiunctus,
}
impl fmt::Display for ErrorGregis {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
ErrorGregis::Inventio(msg) => write!(f, "Discovery error: {msg}"),
ErrorGregis::Consensus(msg) => write!(f, "Consensus error: {msg}"),
ErrorGregis::Rete(msg) => write!(f, "Network error: {msg}"),
ErrorGregis::Configuratio(msg) => write!(f, "Configuration error: {msg}"),
ErrorGregis::NodusNonInventus(id) => write!(f, "Node not found: {id:?}"),
ErrorGregis::SineQuorum => write!(f, "No quorum"),
ErrorGregis::IamConiunctus => write!(f, "Already joined cluster"),
ErrorGregis::NonConiunctus => write!(f, "Not joined to cluster"),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::distributed::{FacultatesNodi, MunusNodi};
fn test_node(id: u64) -> InformationesNodi {
InformationesNodi {
identitas: NodusIdentitas::new(0, id),
inscriptio: InscriptioNodi::new("localhost", 8080 + id as u16),
munus: MunusNodi::Executor,
status: StatusNodi::Sanus,
facultates: FacultatesNodi::default(),
tituli: Vec::new(),
ultima_pulsatio: None,
}
}
#[test]
fn test_configuratio_default() {
let config = ConfiguratioGregis::default();
assert_eq!(config.nomen, "ordofp-cluster");
assert_eq!(config.factor_replicationis, 1);
}
#[test]
fn add_node_updates_known_nodes() {
let mut status = StatusGregis::new(ConfiguratioGregis::default());
let mut node = test_node(1);
status.add_node(node.clone());
assert_eq!(
status.get_node(node.identitas).unwrap().status,
StatusNodi::Sanus
);
node.status = StatusNodi::Aegrotus;
status.add_node(node.clone());
assert_eq!(
status.get_node(node.identitas).unwrap().status,
StatusNodi::Aegrotus
);
assert_eq!(
status.nodi.len(),
1,
"re-discovery must not duplicate the node"
);
}
#[test]
fn test_status_gregis_add_node() {
let mut status = StatusGregis::new(ConfiguratioGregis::default());
status.add_node(test_node(1));
status.add_node(test_node(2));
assert_eq!(status.nodi.len(), 2);
assert_eq!(status.salus, SalusGregis::Sanus);
}
#[test]
fn test_status_gregis_remove_node() {
let mut status = StatusGregis::new(ConfiguratioGregis::default());
status.add_node(test_node(1));
status.add_node(test_node(2));
status.remove_node(NodusIdentitas::new(0, 1));
assert_eq!(status.nodi.len(), 1);
}
#[test]
fn test_status_gregis_health() {
let mut status = StatusGregis::new(ConfiguratioGregis::default());
for i in 1..=3 {
status.add_node(test_node(i));
}
assert_eq!(status.salus, SalusGregis::Sanus);
status.update_node_status(NodusIdentitas::new(0, 1), StatusNodi::Aegrotus);
assert_eq!(status.salus, SalusGregis::Degradatus);
assert!(status.has_quorum());
}
#[test]
fn test_administrator_select_nodes() {
let mut admin =
AdministratorGregis::new(NodusIdentitas::new(0, 0), ConfiguratioGregis::default());
let mut node1 = test_node(1);
node1.facultates.effectus_tractati = alloc::vec![1, 2];
node1.facultates.munera_maxima = 10;
node1.facultates.munera_currentia = 2;
let mut node2 = test_node(2);
node2.facultates.effectus_tractati = alloc::vec![1, 3];
node2.facultates.munera_maxima = 10;
node2.facultates.munera_currentia = 5;
admin.status.add_node(node1);
admin.status.add_node(node2);
let selected = admin.select_nodes(&[1], 2);
assert_eq!(selected.len(), 2);
let selected = admin.select_nodes(&[2], 2);
assert_eq!(selected.len(), 1);
}
}