use bytemuck::{Pod, Zeroable};
use mesocarp::logging::journal::Journal;
use crate::{
env::Environment,
objects::{AntiMsg, Msg, SchedulingTask, Transfer},
AikaError,
};
#[derive(Debug)]
pub struct Context<MessageType: Pod + Zeroable + Clone> {
pub env: Box<dyn Environment>,
pub time: u64,
pub terminal: u64,
pub cluster_id: usize,
pub connected: bool,
pub anti_msgs: Journal,
pub(crate) outbox: Vec<Transfer<MessageType>>,
pub(crate) counter: u16,
}
impl<MessageType: Pod + Zeroable + Clone> Context<MessageType> {
pub fn new(
env: impl Environment + 'static,
connected: bool,
cluster_id: usize,
terminal: u64,
) -> Self {
let env = Box::new(env);
Self {
env,
time: 0,
cluster_id,
connected,
anti_msgs: Journal::init(8 * 1024),
outbox: Vec::new(),
terminal,
counter: 0,
}
}
pub fn send_mail(
&mut self,
mut msg: Msg<MessageType>,
to_cluster: usize,
) -> Result<(), AikaError> {
if msg.recv > self.terminal {
return Ok(());
}
if msg.recv < self.time {
return Err(AikaError::TimeTravel);
}
let to_cluster = if self.connected {
to_cluster
} else {
self.cluster_id
};
msg.to.0 = to_cluster;
msg.from.0 = self.cluster_id;
if self.connected {
let anti = AntiMsg::new(msg.sent, msg.recv, msg.from, msg.to);
self.anti_msgs.write(anti, self.time, None);
}
let outgoing = Transfer::Msg(msg);
if self.counter < 1000 {
self.counter += 1;
}
self.outbox.push(outgoing);
Ok(())
}
}
pub trait Actor<MessageType: Clone + Pod + Zeroable> {
fn step(
&mut self,
env: &mut Context<MessageType>,
actor_id: usize,
) -> Result<SchedulingTask, AikaError>;
}
pub trait ConnectedActor<MessageType: Pod + Zeroable + Clone>: Actor<MessageType> {
fn read_message(
&mut self,
env: &mut Context<MessageType>,
msg: Msg<MessageType>,
actor_id: usize,
) -> Result<(), AikaError>;
}
pub enum ActorType<MessageType: Pod + Zeroable + Clone> {
Basic(Box<dyn Actor<MessageType>>),
Connected(Box<dyn ConnectedActor<MessageType>>),
}
impl<MessageType: Pod + Zeroable + Clone> ActorType<MessageType> {
pub fn step(
&mut self,
env: &mut Context<MessageType>,
actor_id: usize,
) -> Result<SchedulingTask, AikaError> {
match self {
ActorType::Basic(actor) => actor.step(env, actor_id),
ActorType::Connected(connected_actor) => connected_actor.step(env, actor_id),
}
}
}