use alloc::string::String;
use alloc::vec::Vec;
use core::fmt;
use core::time::Duration;
use super::node::{InformationesNodi, NodusIdentitas, StatusNodi};
use super::serializable::NodusSerializabilis;
const DEFAULT_COMPUTATION_PRIORITY: u32 = 100;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct VersioProtocolli {
pub maior: u16,
pub minor: u16,
pub emendatio: u16,
}
impl VersioProtocolli {
pub const CURRENS: Self = VersioProtocolli {
maior: 1,
minor: 0,
emendatio: 0,
};
#[inline]
pub const fn new(maior: u16, minor: u16, emendatio: u16) -> Self {
VersioProtocolli {
maior,
minor,
emendatio,
}
}
pub fn is_compatible(&self, other: &Self) -> bool {
self.maior == other.maior && self.minor >= other.minor
}
}
impl fmt::Display for VersioProtocolli {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}.{}.{}", self.maior, self.minor, self.emendatio)
}
}
#[derive(Debug, Clone)]
pub struct CaputNuntii {
pub versio: VersioProtocolli,
pub id: u64,
pub mittens: NodusIdentitas,
pub recipiens: Option<NodusIdentitas>,
pub tempus: Duration,
pub genus: GenusNuntii,
pub correlatio: Option<u64>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum GenusNuntii {
Pulsatio,
Coniunctio,
Abscessio,
StatusMutatio,
InquisitioNodorum,
ResponsumNodorum,
SubmissioComputationis,
ResultatumComputationis,
CancellatioComputationis,
OperatioEffectus,
ResponsumEffectus,
Suffragium,
AnnuntiatioDucis,
ReplicatioActorum,
Error,
}
#[derive(Debug, Clone)]
pub struct Nuntius {
pub caput: CaputNuntii,
pub corpus: CorpusNuntii,
}
impl Nuntius {
pub fn new(mittens: NodusIdentitas, genus: GenusNuntii, corpus: CorpusNuntii) -> Self {
static COUNTER: core::sync::atomic::AtomicU64 = core::sync::atomic::AtomicU64::new(1);
Nuntius {
caput: CaputNuntii {
versio: VersioProtocolli::CURRENS,
id: COUNTER.fetch_add(1, core::sync::atomic::Ordering::SeqCst),
mittens,
recipiens: None,
tempus: Duration::from_secs(0), genus,
correlatio: None,
},
corpus,
}
}
pub fn to(mut self, recipiens: NodusIdentitas) -> Self {
self.caput.recipiens = Some(recipiens);
self
}
pub fn correlating(mut self, id: u64) -> Self {
self.caput.correlatio = Some(id);
self
}
pub fn respond(&self, genus: GenusNuntii, corpus: CorpusNuntii) -> Self {
Nuntius::new(
self.caput.recipiens.unwrap_or(self.caput.mittens),
genus,
corpus,
)
.to(self.caput.mittens)
.correlating(self.caput.id)
}
}
#[derive(Debug, Clone)]
pub enum CorpusNuntii {
Vacuum,
Pulsatio(PulsatioCorpus),
Coniunctio(ConiunctioCorpus),
Nodi(Vec<InformationesNodi>),
Status(StatusNodi),
Computatio(ComputatioCorpus),
Resultatum(ResultatumCorpus),
Effectus(EffectusCorpus),
ResponsumEffectus(ResponsumEffectusCorpus),
Suffragium(SuffragiumCorpus),
Acta(ActaCorpus),
Error(ErrorCorpus),
}
#[derive(Debug, Clone)]
pub struct PulsatioCorpus {
pub status: StatusNodi,
pub onus: f32,
pub munera_activa: u32,
pub generatio: u64,
}
#[derive(Debug, Clone)]
pub struct ConiunctioCorpus {
pub informationes: InformationesNodi,
pub munus_petitus: super::node::MunusNodi,
}
#[derive(Debug, Clone)]
pub struct ComputatioCorpus {
pub id: u64,
pub nodus: NodusSerializabilis,
pub effectus_requiriti: Vec<u64>,
pub prioritas: u32,
pub mora: Option<Duration>,
}
#[derive(Debug, Clone)]
pub struct ResultatumCorpus {
pub computatio_id: u64,
pub exitus: ExitusComputationis,
pub tempus: Duration,
}
#[derive(Debug, Clone)]
pub enum ExitusComputationis {
Successus(Vec<u8>),
Defectio(String),
Cancellatus,
MoraExcessit,
}
#[derive(Debug, Clone)]
pub struct EffectusCorpus {
pub effectus_id: u64,
pub operatio_id: u64,
pub data: Vec<u8>,
}
#[derive(Debug, Clone)]
pub struct ResponsumEffectusCorpus {
pub operatio_id: u64,
pub exitus: ExitusEffectus,
}
#[derive(Debug, Clone)]
pub enum ExitusEffectus {
Successus(Vec<u8>),
Defectio(String),
NonSuffultus,
}
#[derive(Debug, Clone)]
pub struct SuffragiumCorpus {
pub terminus: u64,
pub candidatus: NodusIdentitas,
pub ultimus_index: u64,
pub ultimus_terminus: u64,
pub concessum: bool,
}
#[derive(Debug, Clone)]
pub struct ActaCorpus {
pub terminus: u64,
pub dux: NodusIdentitas,
pub index_prior: u64,
pub terminus_prior: u64,
pub ingressus: Vec<IngressusActorum>,
pub index_commissi: u64,
}
#[derive(Debug, Clone)]
pub struct IngressusActorum {
pub terminus: u64,
pub index: u64,
pub data: Vec<u8>,
}
#[derive(Debug, Clone)]
pub struct ErrorCorpus {
pub codex: u32,
pub nuntius: String,
pub iterabilis: bool,
}
impl ErrorCorpus {
pub const NODE_NOT_FOUND: u32 = 1;
pub const EFFECT_NOT_SUPPORTED: u32 = 2;
pub const TIMEOUT: u32 = 3;
pub const SERIALIZATION_ERROR: u32 = 4;
pub const NOT_LEADER: u32 = 5;
pub const NO_QUORUM: u32 = 6;
pub const COMPUTATION_FAILED: u32 = 7;
pub const CANCELLED: u32 = 8;
#[inline]
pub fn new(codex: u32, nuntius: impl Into<String>) -> Self {
ErrorCorpus {
codex,
nuntius: nuntius.into(),
iterabilis: false,
}
}
pub fn retriable(mut self) -> Self {
self.iterabilis = true;
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum FormaFili {
#[default]
Binarius,
Json,
MessagePack,
Protobuf,
}
pub trait Serializabilis: Sized {
fn serialize(&self, forma: FormaFili) -> Result<Vec<u8>, ErrorSerialization>;
fn deserialize(bytes: &[u8], forma: FormaFili) -> Result<Self, ErrorSerialization>;
}
#[derive(Debug, Clone)]
pub struct ErrorSerialization {
pub nuntius: String,
pub offset: Option<usize>,
}
impl fmt::Display for ErrorSerialization {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self.offset {
Some(off) => write!(f, "Serialization error at byte {}: {}", off, self.nuntius),
None => write!(f, "Serialization error: {}", self.nuntius),
}
}
}
pub struct AedificatorNuntii {
mittens: NodusIdentitas,
}
impl AedificatorNuntii {
#[inline]
pub fn new(mittens: NodusIdentitas) -> Self {
AedificatorNuntii { mittens }
}
#[inline]
pub fn pulsatio(&self, status: StatusNodi, onus: f32, munera: u32, generatio: u64) -> Nuntius {
Nuntius::new(
self.mittens,
GenusNuntii::Pulsatio,
CorpusNuntii::Pulsatio(PulsatioCorpus {
status,
onus,
munera_activa: munera,
generatio,
}),
)
}
#[inline]
pub fn coniunctio(&self, info: InformationesNodi) -> Nuntius {
Nuntius::new(
self.mittens,
GenusNuntii::Coniunctio,
CorpusNuntii::Coniunctio(ConiunctioCorpus {
munus_petitus: info.munus,
informationes: info,
}),
)
}
#[inline]
pub fn computatio(&self, id: u64, nodus: NodusSerializabilis, effectus: Vec<u64>) -> Nuntius {
Nuntius::new(
self.mittens,
GenusNuntii::SubmissioComputationis,
CorpusNuntii::Computatio(ComputatioCorpus {
id,
nodus,
effectus_requiriti: effectus,
prioritas: DEFAULT_COMPUTATION_PRIORITY,
mora: None,
}),
)
}
#[inline]
pub fn effectus(&self, effectus_id: u64, operatio_id: u64, data: Vec<u8>) -> Nuntius {
Nuntius::new(
self.mittens,
GenusNuntii::OperatioEffectus,
CorpusNuntii::Effectus(EffectusCorpus {
effectus_id,
operatio_id,
data,
}),
)
}
#[inline]
pub fn error(&self, codex: u32, nuntius: impl Into<String>) -> Nuntius {
Nuntius::new(
self.mittens,
GenusNuntii::Error,
CorpusNuntii::Error(ErrorCorpus::new(codex, nuntius)),
)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_versio_compatibility() {
let v1 = VersioProtocolli::new(1, 0, 0);
let v2 = VersioProtocolli::new(1, 1, 0);
let v3 = VersioProtocolli::new(2, 0, 0);
assert!(v2.is_compatible(&v1));
assert!(!v1.is_compatible(&v2));
assert!(!v3.is_compatible(&v1));
}
#[test]
fn test_nuntius_creation() {
let sender = NodusIdentitas::new(1, 1);
let msg = Nuntius::new(sender, GenusNuntii::Pulsatio, CorpusNuntii::Vacuum);
assert_eq!(msg.caput.mittens, sender);
assert_eq!(msg.caput.genus, GenusNuntii::Pulsatio);
}
#[test]
fn test_nuntius_response() {
let sender = NodusIdentitas::new(1, 1);
let recipient = NodusIdentitas::new(2, 2);
let request = Nuntius::new(sender, GenusNuntii::InquisitioNodorum, CorpusNuntii::Vacuum)
.to(recipient);
let response = request.respond(
GenusNuntii::ResponsumNodorum,
CorpusNuntii::Nodi(Vec::new()),
);
assert_eq!(response.caput.recipiens, Some(sender));
assert_eq!(response.caput.correlatio, Some(request.caput.id));
}
#[test]
fn test_aedificator_nuntii() {
let sender = NodusIdentitas::new(1, 1);
let builder = AedificatorNuntii::new(sender);
let msg = builder.pulsatio(StatusNodi::Sanus, 0.5, 10, 1);
assert_eq!(msg.caput.genus, GenusNuntii::Pulsatio);
if let CorpusNuntii::Pulsatio(body) = msg.corpus {
assert_eq!(body.status, StatusNodi::Sanus);
assert!((body.onus - 0.5).abs() < f32::EPSILON);
} else {
panic!("Wrong message body type");
}
}
#[test]
fn test_error_corpus() {
let err = ErrorCorpus::new(ErrorCorpus::TIMEOUT, "Operation timed out").retriable();
assert_eq!(err.codex, ErrorCorpus::TIMEOUT);
assert!(err.iterabilis);
}
}