pub const DEFAULT_SHUTDOWN_TIMEOUT_MS: u64 = 5000;
pub const DEFAULT_MAX_RESTARTS: u32 = 3;
pub const DEFAULT_RESTART_WITHIN_SECONDS: u32 = 5;
pub const DEFAULT_EXPONENS_INITIAL_MS: u64 = 100;
pub const DEFAULT_EXPONENS_MAX_MS: u64 = 30000;
pub const DEFAULT_EXPONENS_FACTOR: f64 = 2.0;
use alloc::boxed::Box;
use alloc::string::String;
use alloc::sync::Arc;
use alloc::vec::Vec;
use core::future::Future;
use core::marker::PhantomData;
use core::pin::Pin;
use core::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use super::fibra::{Fibra, FibraError, FibraManubrium, TesseraAbrogationis};
use super::runtime::RuntimeGenerare;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum StrategiaRestart {
#[default]
UnusProUno,
OmnesProUno,
ReliquiProUno,
Nullus,
}
#[derive(Debug, Clone, Copy)]
pub struct IntensitasRestart {
pub max_restarts: u32,
pub within_seconds: u32,
}
impl Default for IntensitasRestart {
fn default() -> Self {
IntensitasRestart {
max_restarts: DEFAULT_MAX_RESTARTS,
within_seconds: DEFAULT_RESTART_WITHIN_SECONDS,
}
}
}
pub struct InfansSpecificatio {
pub nomen: String,
pub fabricare: Arc<dyn Fn() -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync>,
pub restart: bool,
pub shutdown_timeout_ms: u64,
}
impl InfansSpecificatio {
pub fn new<F, Fut>(nomen: impl Into<String>, fabricare: F) -> Self
where
F: Fn() -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
InfansSpecificatio {
nomen: nomen.into(),
fabricare: Arc::new(move || Box::pin(fabricare())),
restart: true,
shutdown_timeout_ms: DEFAULT_SHUTDOWN_TIMEOUT_MS,
}
}
#[inline]
pub fn with_restart(mut self, restart: bool) -> Self {
self.restart = restart;
self
}
#[inline]
pub fn with_shutdown_timeout(mut self, timeout_ms: u64) -> Self {
self.shutdown_timeout_ms = timeout_ms;
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StatusInfans {
Incipit,
Currens,
Perfectus,
Defectus,
Renovatur,
Terminatus,
}
pub struct Praefectus<R: RuntimeGenerare> {
strategia: StrategiaRestart,
intensitas: IntensitasRestart,
specs: Vec<InfansSpecificatio>,
running: Arc<AtomicBool>,
restart_count: Arc<AtomicU32>,
sistendum: Arc<TesseraAbrogationis>,
_runtime: PhantomData<R>,
}
impl<R: RuntimeGenerare> Praefectus<R> {
#[inline]
pub fn new(strategia: StrategiaRestart) -> Self {
Praefectus {
strategia,
intensitas: IntensitasRestart::default(),
specs: Vec::with_capacity(8),
running: Arc::new(AtomicBool::new(false)),
restart_count: Arc::new(AtomicU32::new(0)),
sistendum: Arc::new(TesseraAbrogationis::default()),
_runtime: PhantomData,
}
}
#[inline]
pub fn with_intensitas(mut self, intensitas: IntensitasRestart) -> Self {
self.intensitas = intensitas;
self
}
#[inline]
pub fn add_child(&mut self, spec: InfansSpecificatio) {
self.specs.push(spec);
}
#[inline]
pub fn add_child_fn<F, Fut>(&mut self, nomen: impl Into<String>, f: F)
where
F: Fn() -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
self.add_child(InfansSpecificatio::new(nomen, f));
}
#[inline]
pub fn strategia(&self) -> StrategiaRestart {
self.strategia
}
#[inline]
pub fn child_count(&self) -> usize {
self.specs.len()
}
#[inline]
pub fn is_running(&self) -> bool {
self.running.load(Ordering::SeqCst)
}
#[inline]
pub fn stop(&self) {
self.running.store(false, Ordering::SeqCst);
self.sistendum.abrogare();
}
#[inline]
pub fn stop_handle(&self) -> Arc<TesseraAbrogationis> {
self.sistendum.clone()
}
#[allow(clippy::too_many_lines)]
pub fn start(self) -> Fibra<()> {
let running = self.running.clone();
let restart_count = self.restart_count.clone();
let strategia = self.strategia;
let intensitas = self.intensitas;
let sistendum = self.sistendum.clone();
let specs = self.specs;
running.store(true, Ordering::SeqCst);
Fibra::new(async move {
let mut children: Vec<(
InfansSpecificatio,
Option<FibraManubrium<()>>,
StatusInfans,
u32,
)> = specs
.into_iter()
.map(|spec| (spec, None, StatusInfans::Incipit, 0u32))
.collect();
for (spec, handle, status, _restart_count) in &mut children {
let fut = (spec.fabricare)();
let fibra = Fibra::new(fut);
let id = fibra.id();
let cancelled = fibra.cancellation_token();
let join_handle = R::spawn(fibra);
*handle = Some(FibraManubrium::new(id, join_handle, cancelled));
*status = StatusInfans::Currens;
}
let max_restarts = intensitas.max_restarts as usize;
let window_secs = u64::from(intensitas.within_seconds);
let mut restart_times: Vec<u64> = Vec::with_capacity(max_restarts);
let origo = std::time::Instant::now();
loop {
let event = core::future::poll_fn(|cx| {
use core::task::Poll;
sistendum.registrare(cx);
if sistendum.est_abrogata() || !running.load(Ordering::SeqCst) {
return Poll::Ready(None);
}
for (i, (_spec, handle, status, _rc)) in children.iter_mut().enumerate() {
if *status == StatusInfans::Currens
&& let Some(h) = handle.as_mut()
&& let Poll::Ready(exitus) = h.poll_conjunctio(cx)
{
return Poll::Ready(Some((i, exitus)));
}
}
Poll::Pending
})
.await;
let Some((idx, exitus)) = event else { break };
children[idx].1 = None;
if exitus.is_ok() {
children[idx].2 = StatusInfans::Perfectus;
continue;
}
if !children[idx].0.restart {
children[idx].2 = StatusInfans::Defectus;
continue;
}
let len = children.len();
let indices: Vec<usize> = match strategia {
StrategiaRestart::UnusProUno => alloc::vec![idx],
StrategiaRestart::OmnesProUno => (0..len).collect(),
StrategiaRestart::ReliquiProUno => (idx..len).collect(),
StrategiaRestart::Nullus => Vec::new(),
};
if indices.is_empty() {
children[idx].2 = StatusInfans::Defectus;
continue;
}
let now_secs = origo.elapsed().as_secs();
let cutoff = now_secs.saturating_sub(window_secs);
restart_times.retain(|&t| t >= cutoff);
if restart_times.len() >= max_restarts {
for (_s, h, st, _rc) in &mut children {
if let Some(handle) = h.take() {
handle.cancel();
let _ = handle.conjungere().await; }
*st = StatusInfans::Terminatus;
}
break;
}
restart_times.push(now_secs);
for i in indices {
if let Some(old) = children[i].1.take() {
old.cancel();
let _ = old.conjungere().await;
}
let fut = (children[i].0.fabricare)();
let fibra = Fibra::new(fut);
let id = fibra.id();
let token = fibra.cancellation_token();
let join_handle = R::spawn(fibra);
children[i].1 = Some(FibraManubrium::new(id, join_handle, token));
children[i].2 = StatusInfans::Currens;
children[i].3 += 1;
restart_count.fetch_add(1, Ordering::SeqCst);
}
}
for (_spec, handle, status, _rc) in &mut children {
if let Some(h) = handle.take() {
h.cancel();
let _ = h.conjungere().await;
}
*status = StatusInfans::Terminatus;
}
running.store(false, Ordering::SeqCst);
})
}
}
impl<R: RuntimeGenerare> Default for Praefectus<R> {
fn default() -> Self {
Self::new(StrategiaRestart::UnusProUno)
}
}
#[derive(Debug, Clone)]
pub enum EventusSupervisio {
InfansIncepit {
nomen: String,
},
InfansPerfectus {
nomen: String,
},
InfansDefectus {
nomen: String,
error: String,
},
InfansRenovatur {
nomen: String,
attempt: u32,
},
LimesExcessus {
nomen: String,
},
SupervisorSistit,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum PolitiaDefectus {
#[default]
Renovare,
Sistere,
Escalare,
Resumere,
}
pub type DeciderDefectus = Box<dyn Fn(&FibraError) -> PolitiaDefectus + Send + Sync>;
#[inline]
pub fn default_decider() -> DeciderDefectus {
Box::new(|_| PolitiaDefectus::Renovare)
}
#[derive(Debug, Clone, Copy)]
pub enum StrategiaMora {
Nulla,
Constans {
delay_ms: u64,
},
Exponens {
initial_ms: u64,
max_ms: u64,
factor: f64,
},
Linearis {
initial_ms: u64,
increment_ms: u64,
max_ms: u64,
},
}
impl Default for StrategiaMora {
fn default() -> Self {
StrategiaMora::Exponens {
initial_ms: DEFAULT_EXPONENS_INITIAL_MS,
max_ms: DEFAULT_EXPONENS_MAX_MS,
factor: DEFAULT_EXPONENS_FACTOR,
}
}
}
impl StrategiaMora {
#[inline]
pub fn delay_for_attempt(&self, attempt: u32) -> u64 {
match self {
StrategiaMora::Nulla => 0,
StrategiaMora::Constans { delay_ms } => *delay_ms,
StrategiaMora::Exponens {
initial_ms,
max_ms,
factor,
} => {
let delay = (*initial_ms as f64) * factor.powi(attempt as i32);
(delay as u64).min(*max_ms)
}
StrategiaMora::Linearis {
initial_ms,
increment_ms,
max_ms,
} => {
let delay = initial_ms + (increment_ms * u64::from(attempt));
delay.min(*max_ms)
}
}
}
}
#[inline]
pub fn supervisor_unus_pro_uno<R: RuntimeGenerare>() -> Praefectus<R> {
Praefectus::new(StrategiaRestart::UnusProUno)
}
#[inline]
pub fn supervisor_omnes_pro_uno<R: RuntimeGenerare>() -> Praefectus<R> {
Praefectus::new(StrategiaRestart::OmnesProUno)
}
#[inline]
pub fn supervisor_reliqui_pro_uno<R: RuntimeGenerare>() -> Praefectus<R> {
Praefectus::new(StrategiaRestart::ReliquiProUno)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::async_core::runtime::NullRuntime;
#[test]
fn test_strategia_restart_default() {
assert_eq!(StrategiaRestart::default(), StrategiaRestart::UnusProUno);
}
#[test]
fn test_intensitas_default() {
let i = IntensitasRestart::default();
assert_eq!(i.max_restarts, 3);
assert_eq!(i.within_seconds, 5);
}
#[test]
fn test_praefectus_new() {
let sup: Praefectus<NullRuntime> = Praefectus::new(StrategiaRestart::OmnesProUno);
assert_eq!(sup.strategia(), StrategiaRestart::OmnesProUno);
assert_eq!(sup.child_count(), 0);
assert!(!sup.is_running());
}
#[test]
fn test_praefectus_add_child() {
let mut sup: Praefectus<NullRuntime> = Praefectus::new(StrategiaRestart::UnusProUno);
sup.add_child_fn("test", || async {});
assert_eq!(sup.child_count(), 1);
}
#[test]
fn test_strategia_mora_nulla() {
let s = StrategiaMora::Nulla;
assert_eq!(s.delay_for_attempt(0), 0);
assert_eq!(s.delay_for_attempt(10), 0);
}
#[test]
fn test_strategia_mora_constans() {
let s = StrategiaMora::Constans { delay_ms: 1000 };
assert_eq!(s.delay_for_attempt(0), 1000);
assert_eq!(s.delay_for_attempt(10), 1000);
}
#[test]
fn test_strategia_mora_exponens() {
let s = StrategiaMora::Exponens {
initial_ms: 100,
max_ms: 10000,
factor: 2.0,
};
assert_eq!(s.delay_for_attempt(0), 100);
assert_eq!(s.delay_for_attempt(1), 200);
assert_eq!(s.delay_for_attempt(2), 400);
assert_eq!(s.delay_for_attempt(10), 10000); }
#[test]
fn test_strategia_mora_linearis() {
let s = StrategiaMora::Linearis {
initial_ms: 100,
increment_ms: 100,
max_ms: 500,
};
assert_eq!(s.delay_for_attempt(0), 100);
assert_eq!(s.delay_for_attempt(1), 200);
assert_eq!(s.delay_for_attempt(2), 300);
assert_eq!(s.delay_for_attempt(10), 500); }
#[test]
fn test_politia_defectus_default() {
assert_eq!(PolitiaDefectus::default(), PolitiaDefectus::Renovare);
}
#[test]
fn test_status_infans() {
assert_eq!(StatusInfans::Currens, StatusInfans::Currens);
assert_ne!(StatusInfans::Currens, StatusInfans::Defectus);
}
#[test]
fn is_running_false_after_supervision_loop_terminates() {
use core::future::Future;
use core::pin::Pin;
use core::task::{Context, Poll, Waker};
let sup: Praefectus<NullRuntime> = Praefectus::new(StrategiaRestart::UnusProUno);
let running = sup.running.clone();
sup.stop_handle().abrogare(); assert!(!running.load(Ordering::SeqCst));
let mut fiber = sup.start();
assert!(running.load(Ordering::SeqCst));
let waker = Waker::noop();
let mut cx = Context::from_waker(waker);
let poll = Pin::new(&mut fiber).poll(&mut cx);
assert!(
matches!(poll, Poll::Ready(_)),
"expected the zero-child, pre-stopped supervision loop to resolve on the first poll"
);
assert!(
!running.load(Ordering::SeqCst),
"is_running() must report false once the supervision loop has terminated"
);
}
}