diff --git a/examples/hello.rs b/examples/hello.rs index d83fbc0..7005d24 100644 --- a/examples/hello.rs +++ b/examples/hello.rs @@ -1,4 +1,4 @@ -use swactor::{ActorAddress, ActorInterface, Message, Runtime, RuntimeFlavor}; +use swactor::{actor::{ActorAddress, ActorInterface}, runtime::{Runtime, RuntimeFlavor}}; #[derive(Debug, Default)] struct Greeter { @@ -31,7 +31,6 @@ impl ActorInterface for Greeter { } } - fn main() { let mut rt = Runtime::new(100, Some(RuntimeFlavor::SingleThreaded)); let addr = rt diff --git a/examples/ring.rs b/examples/ring.rs index 4e854f5..bcfc07e 100644 --- a/examples/ring.rs +++ b/examples/ring.rs @@ -1,7 +1,7 @@ -use swactor::{ActorAddress, ActorInterface, Inbox, Runtime, RuntimeFlavor}; +use swactor::{actor::{ActorAddress, ActorInterface}, runtime::{Inbox, Runtime, RuntimeFlavor}}; #[derive(Debug, Default, Clone)] -struct RingMessage { +pub struct RingMessage { count: usize, } diff --git a/src/actor.rs b/src/actor.rs new file mode 100644 index 0000000..4f896c1 --- /dev/null +++ b/src/actor.rs @@ -0,0 +1,60 @@ +use crate::{runtime::Runtime, WATERLEVEL, ring_buffer::Receiver}; + + +pub trait Message: 'static + Sized + Clone + Send {} +impl Message for T {} + +pub trait ActorInterface: 'static + Send { + type Incoming: Message; + type Response: Message; + fn handle(&mut self, ctx: &Runtime, msg: Self::Incoming); +} + +pub type ActorAddress = u64; + +pub struct Actor +where + A: ActorInterface, +{ + _addr: ActorAddress, + inbox: Receiver, + inner: A, +} + +impl Actor { + pub(crate) fn new(addr: ActorAddress, inbox: Receiver, inner: A) -> Self { + Self { + _addr: addr, + inbox, + inner + } + } +} + +/// Trait for type-erased actors +pub(crate) trait AnyActor: Send { + fn tick(&mut self, ctx: &Runtime); +} + +impl AnyActor for Actor +where + A: ActorInterface, +{ + fn tick(&mut self, ctx: &Runtime) { + let total_messages = self.inbox.len(); + let messages_to_process = if total_messages < WATERLEVEL { + total_messages + } else { + total_messages >> 1 + }; + + for _ in 0..messages_to_process { + match self.inbox.try_recv() { + Some(msg) => self.inner.handle(ctx, msg), + None => unreachable!( + "We checked number of unprocessed messages in the queue ahead of processing" + ), + } + } + } +} \ No newline at end of file diff --git a/src/error.rs b/src/error.rs index e00c61c..91a260a 100644 --- a/src/error.rs +++ b/src/error.rs @@ -1,6 +1,5 @@ #[derive(Debug)] pub struct Error(Box); -pub type Result = std::result::Result; pub(crate) fn convert_err(e: E) -> Error { Error(format!("{e:?}").into()) } diff --git a/src/lib.rs b/src/lib.rs index 7da4647..85d51ae 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,15 +1,14 @@ +pub mod actor; + +pub(crate) mod error; +pub use error::Error; + mod ring_buffer; - -use std::collections::HashMap; - -use crossbeam_queue::ArrayQueue; -use ring_buffer::{Receiver, Sender}; - -pub mod error; -use error::Error; +mod router; +pub mod runtime; #[cfg(feature = "getrandom")] -pub fn get_random(buf: &mut [u8]) { +pub(crate) fn get_random(buf: &mut [u8]) { getrandom::getrandom(buf).unwrap() } @@ -19,235 +18,4 @@ pub fn get_random(buf: &mut [u8]) { /// else /// process total_messages // 2 const WATERLEVEL: usize = 10; - const DEFAULT_INBOX_CAPACITY: usize = 1_000; - -pub trait Message: 'static + Sized + Clone + Send {} -impl Message for T {} - -pub type Envelope = Box; - -pub trait ActorInterface: 'static + Send { - type Incoming: Message; - type Response: Message; - fn handle(&mut self, ctx: &Runtime, msg: Self::Incoming); -} - -pub type ActorAddress = u64; - -pub struct Actor -where - A: ActorInterface, -{ - _addr: ActorAddress, - inbox: Receiver, - inner: A, -} - -/// Trait for type-erased actors -trait AnyActor: Send { - fn tick(&mut self, ctx: &Runtime); -} - -impl AnyActor for Actor -where - A: ActorInterface, -{ - fn tick(&mut self, ctx: &Runtime) { - let total_messages = self.inbox.len(); - let messages_to_process = if total_messages < WATERLEVEL { - total_messages - } else { - total_messages >> 1 - }; - - for _ in 0..messages_to_process { - match self.inbox.try_recv() { - Some(msg) => self.inner.handle(ctx, msg), - None => unreachable!( - "We checked number of unprocessed messages in the queue ahead of processing" - ), - } - } - } -} - -pub struct Inbox { - addr: ActorAddress, - inner: Receiver, -} - -impl Inbox { - pub fn addr(&self) -> &ActorAddress { - &self.addr - } - - pub fn try_recv(&self) -> Option { - self.inner.try_recv() - } -} - -#[derive(Debug, Default)] -pub enum RuntimeFlavor { - #[default] - SingleThreaded, - Multithreaded(usize), -} - -pub struct Runtime { - flavor: RuntimeFlavor, - router: Router, - router_inbox: Sender, - actor_queue: ArrayQueue>, -} - -impl Runtime { - pub fn new(capacity: usize, flavor: Option) -> Self { - let router = Router::new(DEFAULT_INBOX_CAPACITY); - let router_inbox = router.new_sender(); - Self { - flavor: flavor.unwrap_or_default(), - router, - router_inbox, - actor_queue: ArrayQueue::new(capacity), - } - } - - pub fn spawn(&self, actor: A) -> Result { - let addr = { - let mut bytes = u64::to_le_bytes(0); - get_random(&mut bytes); - u64::from_le_bytes(bytes) - }; - let inbox = Receiver::::new(DEFAULT_INBOX_CAPACITY); - let sender = inbox.new_sender(); - - // Register the sender with the router - let _ = self - .router_inbox - .try_send(RouterMessage::AddAddr(addr, Box::new(sender))); - - self.actor_queue - .push(Box::new(Actor { - _addr: addr, - inbox, - inner: actor, - })) - .map_err(|_| Error::from("Runtime error: Failed to spawn actor."))?; - - Ok(addr) - } - - pub fn send_to(&self, addr: ActorAddress, msg: M) -> Result<(), ()> { - let envelope: Envelope = Box::new(msg); - self.router_inbox - .try_send(RouterMessage::SendToAddr { - addr, - msg: envelope, - }) - .map_err(|_| ()) - } - - pub fn tick(&mut self) { - // Pop actor, tick it, push it back - if let Some(mut actor) = self.actor_queue.pop() { - actor.tick(self); - let _ = self.actor_queue.push(actor); - } - - match self.flavor { - RuntimeFlavor::Multithreaded(_) => (), // router has its own thread - RuntimeFlavor::SingleThreaded => self.router.tick(), - } - } - - pub fn new_inbox(&self) -> Inbox { - let addr = { - let mut bytes = u64::to_le_bytes(0); - get_random(&mut bytes); - u64::from_le_bytes(bytes) - }; - let receiver = Receiver::::new(DEFAULT_INBOX_CAPACITY); - let sender = receiver.new_sender(); - // Register the sender with the router - let _ = self - .router_inbox - .try_send(RouterMessage::AddAddr(addr, Box::new(sender))); - Inbox { - addr, - inner: receiver, - } - } -} - -pub trait SenderT: Send { - fn try_send(&self, envelope: Envelope); -} - -impl SenderT for Sender { - fn try_send(&self, envelope: Envelope) { - if let Ok(msg) = envelope.downcast::() { - let _ = Sender::try_send(self, *msg); - } - } -} - -/// Internal messages for the Router's own inbox -pub enum RouterMessage { - /// register addrs with sender - AddAddr(ActorAddress, Box), - /// remove an actor from the address book - RemoveAddr(ActorAddress), - /// send to - SendToAddr { addr: ActorAddress, msg: Envelope }, -} - -struct Router { - directory: HashMap>, - inbox: Receiver, -} - -impl Router { - pub fn new(cap: usize) -> Self { - Self { - directory: HashMap::new(), - inbox: Receiver::new(cap), - } - } - - pub fn tick(&mut self) { - let total_messages = self.inbox.len(); - let messages_to_process = if total_messages < WATERLEVEL { - total_messages - } else { - total_messages >> 1 - }; - - for _ in 0..messages_to_process { - match self.inbox.try_recv() { - Some(msg) => self.handle(msg), - None => unreachable!("We ran checks on total messages before processing."), - } - } - } - - pub fn new_sender(&self) -> Sender { - self.inbox.new_sender() - } - - fn handle(&mut self, msg: RouterMessage) { - match msg { - RouterMessage::AddAddr(addr, sender) => { - self.directory.insert(addr, sender); - } - RouterMessage::RemoveAddr(addr) => { - self.directory.remove(&addr); - } - RouterMessage::SendToAddr { addr, msg } => { - if let Some(sender) = self.directory.get(&addr) { - sender.try_send(msg); - } - } - } - } -} diff --git a/src/router.rs b/src/router.rs new file mode 100644 index 0000000..77b73d3 --- /dev/null +++ b/src/router.rs @@ -0,0 +1,77 @@ +use std::collections::HashMap; + +use crate::{WATERLEVEL, actor::{ActorAddress, Message}, ring_buffer::{Receiver, Sender}}; + +pub(crate) type Envelope = Box; + +pub(crate) trait SenderT: Send { + fn try_send(&self, envelope: Envelope); +} + +impl SenderT for Sender { + fn try_send(&self, envelope: Envelope) { + if let Ok(msg) = envelope.downcast::() { + let _ = Sender::try_send(self, *msg); + } + } +} + +/// Internal messages for the Router's own inbox +pub(crate) enum RouterMessage { + /// register addrs with sender + AddAddr(ActorAddress, Box), + /// remove an actor from the address book + RemoveAddr(ActorAddress), + /// send to + SendToAddr { addr: ActorAddress, msg: Envelope }, +} + +pub(crate) struct Router { + directory: HashMap>, + inbox: Receiver, +} + +impl Router { + pub fn new(cap: usize) -> Self { + Self { + directory: HashMap::new(), + inbox: Receiver::new(cap), + } + } + + pub fn tick(&mut self) { + let total_messages = self.inbox.len(); + let messages_to_process = if total_messages < WATERLEVEL { + total_messages + } else { + total_messages >> 1 + }; + + for _ in 0..messages_to_process { + match self.inbox.try_recv() { + Some(msg) => self.handle(msg), + None => unreachable!("We ran checks on total messages before processing."), + } + } + } + + pub fn new_sender(&self) -> Sender { + self.inbox.new_sender() + } + + fn handle(&mut self, msg: RouterMessage) { + match msg { + RouterMessage::AddAddr(addr, sender) => { + self.directory.insert(addr, sender); + } + RouterMessage::RemoveAddr(addr) => { + self.directory.remove(&addr); + } + RouterMessage::SendToAddr { addr, msg } => { + if let Some(sender) = self.directory.get(&addr) { + sender.try_send(msg); + } + } + } + } +} diff --git a/src/runtime.rs b/src/runtime.rs new file mode 100644 index 0000000..28f8e4e --- /dev/null +++ b/src/runtime.rs @@ -0,0 +1,113 @@ +use crossbeam_queue::ArrayQueue; + +use crate::{ + DEFAULT_INBOX_CAPACITY, Error, + actor::{Actor, ActorAddress, ActorInterface, AnyActor, Message}, + get_random, + ring_buffer::{Receiver, Sender}, + router::{Envelope, Router, RouterMessage}, +}; + +#[derive(Debug, Default)] +pub enum RuntimeFlavor { + #[default] + SingleThreaded, + Multithreaded(usize), +} + +pub struct Inbox { + addr: ActorAddress, + inner: Receiver, +} + +impl Inbox { + pub fn addr(&self) -> &ActorAddress { + &self.addr + } + + pub fn try_recv(&self) -> Option { + self.inner.try_recv() + } +} + +pub struct Runtime { + flavor: RuntimeFlavor, + router: Router, + router_inbox: Sender, + actor_queue: ArrayQueue>, +} + +impl Runtime { + pub fn new(capacity: usize, flavor: Option) -> Self { + let router = Router::new(DEFAULT_INBOX_CAPACITY); + let router_inbox = router.new_sender(); + Self { + flavor: flavor.unwrap_or_default(), + router, + router_inbox, + actor_queue: ArrayQueue::new(capacity), + } + } + + pub fn spawn(&self, actor: A) -> Result { + let addr = { + let mut bytes = u64::to_le_bytes(0); + get_random(&mut bytes); + u64::from_le_bytes(bytes) + }; + let inbox = Receiver::::new(DEFAULT_INBOX_CAPACITY); + let sender = inbox.new_sender(); + + // Register the sender with the router + let _ = self + .router_inbox + .try_send(RouterMessage::AddAddr(addr, Box::new(sender))); + + self.actor_queue + .push(Box::new(Actor::new(addr, inbox, actor))) + .map_err(|_| Error::from("Runtime error: Failed to spawn actor."))?; + + Ok(addr) + } + + pub fn send_to(&self, addr: ActorAddress, msg: M) -> Result<(), ()> { + let envelope: Envelope = Box::new(msg); + self.router_inbox + .try_send(RouterMessage::SendToAddr { + addr, + msg: envelope, + }) + .map_err(|_| ()) + } + + pub fn tick(&mut self) { + // Pop actor, tick it, push it back + if let Some(mut actor) = self.actor_queue.pop() { + actor.tick(self); + let _ = self.actor_queue.push(actor); + } + + match self.flavor { + RuntimeFlavor::Multithreaded(_) => (), // router has its own thread + RuntimeFlavor::SingleThreaded => self.router.tick(), + } + } + + pub fn new_inbox(&self) -> Inbox { + let addr = { + let mut bytes = u64::to_le_bytes(0); + get_random(&mut bytes); + u64::from_le_bytes(bytes) + }; + let receiver = Receiver::::new(DEFAULT_INBOX_CAPACITY); + let sender = receiver.new_sender(); + // Register the sender with the router + let _ = self + .router_inbox + .try_send(RouterMessage::AddAddr(addr, Box::new(sender))); + Inbox { + addr, + inner: receiver, + } + } +}