use std::marker::PhantomData;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use crate::host::{self, node_id, process_id};
use crate::protocol::ProtocolCapture;
use crate::serializer::{Bincode, Serializer};
use crate::timer::TimerRef;
use crate::{ProcessConfig, Tag};
pub trait IntoProcess<M, S> {
type Process;
fn spawn<C>(
capture: C,
entry: fn(C, Self),
link: Option<Tag>,
config: Option<&ProcessConfig>,
node: Option<u64>,
) -> Self::Process
where
S: Serializer<C> + Serializer<ProtocolCapture<C>>;
}
pub trait NoLink {}
#[derive(Serialize, Deserialize)]
pub struct Process<M, S = Bincode> {
node_id: u64,
id: u64,
#[serde(skip_serializing, default)]
serializer_type: PhantomData<(M, S)>,
}
impl<M, S> Process<M, S> {
pub(crate) fn new(node_id: u64, process_id: u64) -> Self {
Self {
node_id,
id: process_id,
serializer_type: PhantomData,
}
}
pub fn this() -> Self {
Self::new(node_id(), process_id())
}
pub fn spawn<C, T>(capture: C, entry: fn(C, T)) -> T::Process
where
S: Serializer<C> + Serializer<ProtocolCapture<C>>,
T: IntoProcess<M, S>,
T: NoLink,
{
T::spawn(capture, entry, None, None, None)
}
pub fn spawn_node<C, T>(node_id: u64, capture: C, entry: fn(C, T)) -> T::Process
where
S: Serializer<C> + Serializer<ProtocolCapture<C>>,
T: IntoProcess<M, S>,
T: NoLink,
{
T::spawn(capture, entry, None, None, Some(node_id))
}
pub fn spawn_node_config<C, T>(
node_id: u64,
config: &ProcessConfig,
capture: C,
entry: fn(C, T),
) -> T::Process
where
S: Serializer<C> + Serializer<ProtocolCapture<C>>,
T: IntoProcess<M, S>,
T: NoLink,
{
T::spawn(capture, entry, None, Some(config), Some(node_id))
}
pub fn spawn_link<C, T>(capture: C, entry: fn(C, T)) -> T::Process
where
S: Serializer<C> + Serializer<ProtocolCapture<C>>,
T: IntoProcess<M, S>,
{
T::spawn(capture, entry, Some(Tag::new()), None, None)
}
pub fn spawn_link_tag<C, T>(capture: C, tag: Tag, entry: fn(C, T)) -> T::Process
where
S: Serializer<C> + Serializer<ProtocolCapture<C>>,
T: IntoProcess<M, S>,
{
T::spawn(capture, entry, Some(tag), None, None)
}
pub fn spawn_config<C, T>(config: &ProcessConfig, capture: C, entry: fn(C, T)) -> T::Process
where
S: Serializer<C> + Serializer<ProtocolCapture<C>>,
T: IntoProcess<M, S>,
T: NoLink,
{
T::spawn(capture, entry, None, Some(config), None)
}
pub fn spawn_link_config<C, T>(
config: &ProcessConfig,
capture: C,
entry: fn(C, T),
) -> T::Process
where
S: Serializer<C> + Serializer<ProtocolCapture<C>>,
T: IntoProcess<M, S>,
{
T::spawn(capture, entry, Some(Tag::new()), Some(config), None)
}
pub fn spawn_link_config_tag<C, T>(
config: &ProcessConfig,
capture: C,
tag: Tag,
entry: fn(C, T),
) -> T::Process
where
S: Serializer<C> + Serializer<ProtocolCapture<C>>,
T: IntoProcess<M, S>,
{
T::spawn(capture, entry, Some(tag), Some(config), None)
}
pub fn id(&self) -> u64 {
self.id
}
pub fn node_id(&self) -> u64 {
self.node_id
}
pub fn link(&self) {
unsafe { host::api::process::link(0, self.id) };
}
pub fn unlink(&self) {
unsafe { host::api::process::unlink(self.id) };
}
pub fn kill(&self) {
unsafe { host::api::process::kill(self.id) };
}
pub fn register(&self, name: &str) {
let name = format!(
"{} + Process + {}/{}",
name,
std::any::type_name::<M>(),
std::any::type_name::<S>()
);
unsafe { host::api::registry::put(name.as_ptr(), name.len(), self.node_id, self.id) };
}
pub fn lookup(name: &str) -> Option<Self> {
let name = format!(
"{} + Process + {}/{}",
name,
std::any::type_name::<M>(),
std::any::type_name::<S>()
);
let mut id = 0;
let mut node_id = 0;
let result =
unsafe { host::api::registry::get(name.as_ptr(), name.len(), &mut node_id, &mut id) };
if result == 0 {
Some(Self {
node_id,
id,
serializer_type: PhantomData,
})
} else {
None
}
}
}
impl<M, S> Process<M, S>
where
S: Serializer<M>,
{
pub fn send(&self, message: M) {
unsafe { host::api::message::create_data(Tag::none().id(), 0) };
S::encode(&message).unwrap();
host::send(self.node_id, self.id);
}
pub fn send_after(&self, message: M, duration: Duration) -> TimerRef {
unsafe { host::api::message::create_data(Tag::none().id(), 0) };
S::encode(&message).unwrap();
let timer_id =
unsafe { host::api::timer::send_after(self.id, duration.as_millis() as u64) };
TimerRef::new(timer_id)
}
pub fn tag_send(&self, tag: Tag, message: M) {
unsafe { host::api::message::create_data(tag.id(), 0) };
S::encode(&message).unwrap();
host::send(self.node_id, self.id);
}
}
impl<M, S> PartialEq for Process<M, S> {
fn eq(&self, other: &Self) -> bool {
self.id() == other.id() && self.node_id() == other.node_id()
}
}
impl<M, S> Eq for Process<M, S> {}
impl<M, S> std::hash::Hash for Process<M, S> {
fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
self.node_id.hash(state);
self.id.hash(state);
}
}
impl<M, S> std::fmt::Debug for Process<M, S> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Process")
.field("id", &self.id())
.field("node_id", &self.node_id())
.finish()
}
}
impl<M, S> Clone for Process<M, S> {
fn clone(&self) -> Self {
Self {
node_id: self.node_id,
id: self.id,
serializer_type: self.serializer_type,
}
}
}
impl<M, S> Copy for Process<M, S> {}