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 super::node::{AffinitasNodi, InformationesNodi, NodusIdentitas};
use super::serializable::NodusSerializabilis;
pub trait EffectusDistributus: Send + Sync {
const EFFECTUS_ID: u64;
fn nomen() -> &'static str;
fn affinitas(&self) -> AffinitasNodi {
AffinitasNodi::Quodlibet
}
fn serialize(&self) -> Result<EffectusSerializatus, ErrorSerializationis>;
fn requires_ordering(&self) -> bool {
true
}
fn batchable(&self) -> bool {
false
}
}
#[derive(Debug, Clone)]
pub struct EffectusSerializatus {
pub effectus_id: u64,
pub data: Vec<u8>,
pub affinitas: AffinitasNodi,
pub correlatio_id: u64,
}
#[derive(Debug, Clone)]
pub struct ErrorSerializationis {
pub nuntius: String,
}
impl fmt::Display for ErrorSerializationis {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "Serialization error: {}", self.nuntius)
}
}
pub trait TractatorDistributus<E: EffectusDistributus>: Send + Sync {
type Output;
fn handle(
&self,
operation: E,
) -> Pin<Box<dyn Future<Output = Result<Self::Output, ErrorDistributionis>> + Send + '_>>;
fn can_handle_locally(&self) -> bool {
true
}
}
#[derive(Debug, Clone)]
pub enum ErrorDistributionis {
Rete(String),
NodusNonInventus(NodusIdentitas),
EffectusNonSuffultus(u64),
Mora(core::time::Duration),
Serializatio(ErrorSerializationis),
Tractator(String),
NodusInaccessibilis(NodusIdentitas),
}
impl fmt::Display for ErrorDistributionis {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
ErrorDistributionis::Rete(msg) => write!(f, "Network error: {msg}"),
ErrorDistributionis::NodusNonInventus(id) => {
write!(f, "Node not found: {id:?}")
}
ErrorDistributionis::EffectusNonSuffultus(id) => {
write!(f, "Effect {id} not supported")
}
ErrorDistributionis::Mora(d) => write!(f, "Timeout after {d:?}"),
ErrorDistributionis::Serializatio(e) => write!(f, "{e}"),
ErrorDistributionis::Tractator(msg) => write!(f, "Handler error: {msg}"),
ErrorDistributionis::NodusInaccessibilis(id) => {
write!(f, "Node unavailable: {id:?}")
}
}
}
}
#[derive(Clone)]
pub struct ProcuratorRemotus<E> {
pub nodus: NodusIdentitas,
pub connexio: Arc<dyn ConnexioRemota>,
_effectus: core::marker::PhantomData<E>,
}
impl<E: EffectusDistributus> ProcuratorRemotus<E> {
#[inline]
pub fn new(nodus: NodusIdentitas, connexio: Arc<dyn ConnexioRemota>) -> Self {
ProcuratorRemotus {
nodus,
connexio,
_effectus: core::marker::PhantomData,
}
}
pub fn invoke(
&self,
operation: E,
) -> Pin<Box<dyn Future<Output = Result<Vec<u8>, ErrorDistributionis>> + Send + '_>> {
Box::pin(async move {
let serialized = operation
.serialize()
.map_err(ErrorDistributionis::Serializatio)?;
self.connexio.send(self.nodus, serialized).await
})
}
}
impl<E> fmt::Debug for ProcuratorRemotus<E> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ProcuratorRemotus")
.field("nodus", &self.nodus)
.finish_non_exhaustive()
}
}
pub trait ConnexioRemota: Send + Sync {
fn send(
&self,
nodus: NodusIdentitas,
effectus: EffectusSerializatus,
) -> Pin<Box<dyn Future<Output = Result<Vec<u8>, ErrorDistributionis>> + Send + '_>>;
fn is_reachable(&self, nodus: NodusIdentitas) -> bool;
}
pub struct TractatorDirigens<E: EffectusDistributus> {
local: Option<Box<dyn TractatorDistributus<E, Output = Vec<u8>>>>,
remoti: Vec<(NodusIdentitas, ProcuratorRemotus<E>)>,
nodus_localis: NodusIdentitas,
strategia: StrategiaDistributionis,
}
impl<E: EffectusDistributus> TractatorDirigens<E> {
#[inline]
pub fn new(nodus_localis: NodusIdentitas) -> Self {
TractatorDirigens {
local: None,
remoti: Vec::with_capacity(8),
nodus_localis,
strategia: StrategiaDistributionis::RoundRobin { index: 0 },
}
}
pub fn with_local(
mut self,
handler: impl TractatorDistributus<E, Output = Vec<u8>> + 'static,
) -> Self {
self.local = Some(Box::new(handler));
self
}
pub fn with_remote(mut self, nodus: NodusIdentitas, proxy: ProcuratorRemotus<E>) -> Self {
self.remoti.push((nodus, proxy));
self
}
pub fn with_strategy(mut self, strategia: StrategiaDistributionis) -> Self {
self.strategia = strategia;
self
}
pub fn select_node(
&mut self,
operation: &E,
nodes: &[InformationesNodi],
) -> Option<NodusIdentitas> {
let affinity = operation.affinitas();
if matches!(affinity, AffinitasNodi::Localis) {
return Some(self.nodus_localis);
}
if let AffinitasNodi::Nodus(id) = affinity
&& nodes.iter().any(|n| n.identitas == id && n.can_execute())
{
return Some(id);
}
let candidates: Vec<_> = nodes
.iter()
.filter(|n| n.can_execute() && affinity.matches(n))
.collect();
if candidates.is_empty() {
return None;
}
match &mut self.strategia {
StrategiaDistributionis::RoundRobin { index } => {
let selected = &candidates[*index % candidates.len()];
*index = (*index + 1) % candidates.len();
Some(selected.identitas)
}
StrategiaDistributionis::LeastLoaded => candidates
.iter()
.min_by(|a, b| {
a.facultates
.load_factor()
.partial_cmp(&b.facultates.load_factor())
.unwrap_or(core::cmp::Ordering::Equal)
})
.map(|n| n.identitas),
StrategiaDistributionis::AffinityScore => candidates
.iter()
.max_by_key(|n| affinity.score(n))
.map(|n| n.identitas),
StrategiaDistributionis::Random => {
#[allow(clippy::cast_possible_truncation)]
let idx = (candidates[0].identitas.ima as usize) % candidates.len();
Some(candidates[idx].identitas)
}
StrategiaDistributionis::LocalFirst => {
if candidates.iter().any(|n| n.identitas == self.nodus_localis) {
Some(self.nodus_localis)
} else {
Some(candidates[0].identitas)
}
}
}
}
}
#[derive(Debug, Clone)]
pub enum StrategiaDistributionis {
RoundRobin {
index: usize,
},
LeastLoaded,
AffinityScore,
Random,
LocalFirst,
}
#[derive(Debug, Clone, Default)]
pub struct TabulaDirigendi {
pub viae: Vec<(u64, Vec<NodusIdentitas>)>,
}
impl TabulaDirigendi {
#[inline]
pub fn new() -> Self {
TabulaDirigendi {
viae: Vec::with_capacity(16),
}
}
pub fn add_route(&mut self, effectus_id: u64, nodus: NodusIdentitas) {
if let Some((_, nodes)) = self.viae.iter_mut().find(|(id, _)| *id == effectus_id) {
if !nodes.contains(&nodus) {
nodes.push(nodus);
}
} else {
self.viae.push((effectus_id, alloc::vec![nodus]));
}
}
pub fn remove_node(&mut self, nodus: NodusIdentitas) {
for (_, nodes) in &mut self.viae {
nodes.retain(|n| *n != nodus);
}
}
pub fn get_handlers(&self, effectus_id: u64) -> Option<&[NodusIdentitas]> {
self.viae
.iter()
.find(|(id, _)| *id == effectus_id)
.map(|(_, nodes)| nodes.as_slice())
}
}
#[derive(Debug, Clone)]
pub struct ComputatioDistributa {
pub id: u64,
pub nodus: NodusSerializabilis,
pub effectus_requiriti: Vec<u64>,
pub affinitas: AffinitasNodi,
pub prioritas: u32,
pub mora_maxima: Option<core::time::Duration>,
}
impl ComputatioDistributa {
#[inline]
pub fn new(id: u64, nodus: NodusSerializabilis) -> Self {
ComputatioDistributa {
id,
nodus,
effectus_requiriti: Vec::with_capacity(4),
affinitas: AffinitasNodi::Quodlibet,
prioritas: 100,
mora_maxima: None,
}
}
pub fn require_effect(mut self, effectus_id: u64) -> Self {
self.effectus_requiriti.push(effectus_id);
self
}
pub fn with_affinity(mut self, affinitas: AffinitasNodi) -> Self {
self.affinitas = affinitas;
self
}
pub fn with_priority(mut self, prioritas: u32) -> Self {
self.prioritas = prioritas;
self
}
pub fn with_timeout(mut self, timeout: core::time::Duration) -> Self {
self.mora_maxima = Some(timeout);
self
}
}
#[derive(Debug, Clone)]
pub struct ResultatumDistributum {
pub computatio_id: u64,
pub executor: NodusIdentitas,
pub data: Vec<u8>,
pub tempus_executionis: core::time::Duration,
pub effectus_effecti: Vec<u64>,
}
#[cfg(test)]
mod tests {
use super::super::node::{FacultatesNodi, InscriptioNodi, MunusNodi, StatusNodi};
use super::*;
use alloc::vec;
struct EffectusProbationis {
affinitas: AffinitasNodi,
}
impl EffectusDistributus for EffectusProbationis {
const EFFECTUS_ID: u64 = 1;
fn nomen() -> &'static str {
"probatio"
}
fn affinitas(&self) -> AffinitasNodi {
self.affinitas.clone()
}
fn serialize(&self) -> Result<EffectusSerializatus, ErrorSerializationis> {
Ok(EffectusSerializatus {
effectus_id: Self::EFFECTUS_ID,
data: Vec::new(),
affinitas: self.affinitas.clone(),
correlatio_id: 0,
})
}
}
fn probatio_node(id: u64, status: StatusNodi) -> InformationesNodi {
InformationesNodi {
identitas: NodusIdentitas::new(0, id),
inscriptio: InscriptioNodi::new("localhost", 8080 + id as u16),
munus: MunusNodi::Executor,
status,
facultates: FacultatesNodi::default(),
tituli: Vec::new(),
ultima_pulsatio: None,
}
}
#[test]
fn select_node_pinned_but_unhealthy_falls_through() {
let dead_id = NodusIdentitas::new(0, 1);
let nodes = vec![probatio_node(1, StatusNodi::Aegrotus)];
let mut router = TractatorDirigens::<EffectusProbationis>::new(NodusIdentitas::new(0, 0));
let op = EffectusProbationis {
affinitas: AffinitasNodi::Nodus(dead_id),
};
assert_ne!(router.select_node(&op, &nodes), Some(dead_id));
assert_eq!(router.select_node(&op, &nodes), None);
}
#[test]
fn select_node_pinned_and_healthy_is_returned() {
let healthy_id = NodusIdentitas::new(0, 1);
let nodes = vec![probatio_node(1, StatusNodi::Sanus)];
let mut router = TractatorDirigens::<EffectusProbationis>::new(NodusIdentitas::new(0, 0));
let op = EffectusProbationis {
affinitas: AffinitasNodi::Nodus(healthy_id),
};
assert_eq!(router.select_node(&op, &nodes), Some(healthy_id));
}
#[test]
fn test_tabula_dirigendi() {
let mut table = TabulaDirigendi::new();
let node1 = NodusIdentitas::new(1, 1);
let node2 = NodusIdentitas::new(2, 2);
table.add_route(1, node1);
table.add_route(1, node2);
table.add_route(2, node1);
let handlers = table
.get_handlers(1)
.expect("route 1 should have handlers after two add_route calls");
assert_eq!(handlers.len(), 2);
table.remove_node(node1);
let handlers = table
.get_handlers(1)
.expect("route 1 should still have one handler after removing node1");
assert_eq!(handlers.len(), 1);
assert_eq!(handlers[0], node2);
}
#[test]
fn test_computatio_distributa() {
let comp = ComputatioDistributa::new(1, NodusSerializabilis::Vacuus)
.require_effect(10)
.require_effect(20)
.with_priority(50);
assert_eq!(comp.effectus_requiriti, vec![10, 20]);
assert_eq!(comp.prioritas, 50);
}
#[test]
fn test_error_display() {
let err = ErrorDistributionis::EffectusNonSuffultus(42);
let msg = alloc::format!("{err}");
assert!(msg.contains("42"));
}
}