use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use crate::vm::Vm;
use super::mailbox::Mailbox;
use super::oneshot;
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct FlowId(pub(crate) u64);
impl FlowId {
pub fn as_u64(self) -> u64 {
self.0
}
}
impl std::fmt::Display for FlowId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "flow#{}", self.0)
}
}
static NEXT_FLOW_ID: AtomicU64 = AtomicU64::new(1);
pub fn next_flow_id() -> FlowId {
FlowId(NEXT_FLOW_ID.fetch_add(1, Ordering::Relaxed))
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum FlowState {
Ready,
Running,
Waiting,
Sleeping,
Terminated,
Failed,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum RestartPolicy {
Always,
OnFailure,
Never,
}
#[derive(Debug, Default)]
pub struct FlowMetrics {
pub instructions: AtomicU64,
pub messages_sent: AtomicU64,
pub messages_received: AtomicU64,
pub reschedules: AtomicU64,
}
pub struct Flow {
pub id: FlowId,
pub vm: Vm,
pub mailbox: Arc<Mailbox>,
pub metrics: Arc<FlowMetrics>,
pub restart_policy: RestartPolicy,
pub(crate) completion: oneshot::Sender<FlowOutcome>,
pub pending_message: Option<crate::bytecode::Value>,
pub last_receive_dest: Option<u8>,
pub(crate) supervisor: Option<super::supervisor::SupervisorLink>,
}
#[derive(Clone, Debug, PartialEq)]
pub enum FlowOutcome {
Completed(crate::bytecode::Value),
Failed(String),
}
impl Flow {
pub fn new(
id: FlowId,
vm: Vm,
mailbox: Arc<Mailbox>,
restart_policy: RestartPolicy,
completion: oneshot::Sender<FlowOutcome>,
) -> Self {
Self {
id,
vm,
mailbox,
metrics: Arc::new(FlowMetrics::default()),
restart_policy,
completion,
pending_message: None,
last_receive_dest: None,
supervisor: None,
}
}
pub(crate) fn complete(self, outcome: FlowOutcome) {
self.completion.send(outcome);
}
}