feat: refactor: split out components into modules

This commit is contained in:
Zachery Aaron Shores-Chmielewski 2026-01-23 14:09:31 +07:00
parent a9f7abba25
commit 2450442566
7 changed files with 261 additions and 245 deletions

View file

@ -1,4 +1,4 @@
use swactor::{ActorAddress, ActorInterface, Message, Runtime, RuntimeFlavor}; use swactor::{actor::{ActorAddress, ActorInterface}, runtime::{Runtime, RuntimeFlavor}};
#[derive(Debug, Default)] #[derive(Debug, Default)]
struct Greeter { struct Greeter {
@ -31,7 +31,6 @@ impl ActorInterface for Greeter {
} }
} }
fn main() { fn main() {
let mut rt = Runtime::new(100, Some(RuntimeFlavor::SingleThreaded)); let mut rt = Runtime::new(100, Some(RuntimeFlavor::SingleThreaded));
let addr = rt let addr = rt

View file

@ -1,7 +1,7 @@
use swactor::{ActorAddress, ActorInterface, Inbox, Runtime, RuntimeFlavor}; use swactor::{actor::{ActorAddress, ActorInterface}, runtime::{Inbox, Runtime, RuntimeFlavor}};
#[derive(Debug, Default, Clone)] #[derive(Debug, Default, Clone)]
struct RingMessage { pub struct RingMessage {
count: usize, count: usize,
} }

60
src/actor.rs Normal file
View file

@ -0,0 +1,60 @@
use crate::{runtime::Runtime, WATERLEVEL, ring_buffer::Receiver};
pub trait Message: 'static + Sized + Clone + Send {}
impl<T: 'static + Sized + Clone + Send> 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<A>
where
A: ActorInterface,
{
_addr: ActorAddress,
inbox: Receiver<A::Incoming>,
inner: A,
}
impl<A: ActorInterface> Actor<A> {
pub(crate) fn new(addr: ActorAddress, inbox: Receiver<A::Incoming>, 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<A> AnyActor for Actor<A>
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"
),
}
}
}
}

View file

@ -1,6 +1,5 @@
#[derive(Debug)] #[derive(Debug)]
pub struct Error(Box<dyn std::error::Error + Send + Sync + 'static>); pub struct Error(Box<dyn std::error::Error + Send + Sync + 'static>);
pub type Result<T> = std::result::Result<T, Error>;
pub(crate) fn convert_err<E: std::fmt::Debug>(e: E) -> Error { pub(crate) fn convert_err<E: std::fmt::Debug>(e: E) -> Error {
Error(format!("{e:?}").into()) Error(format!("{e:?}").into())
} }

View file

@ -1,15 +1,14 @@
pub mod actor;
pub(crate) mod error;
pub use error::Error;
mod ring_buffer; mod ring_buffer;
mod router;
use std::collections::HashMap; pub mod runtime;
use crossbeam_queue::ArrayQueue;
use ring_buffer::{Receiver, Sender};
pub mod error;
use error::Error;
#[cfg(feature = "getrandom")] #[cfg(feature = "getrandom")]
pub fn get_random(buf: &mut [u8]) { pub(crate) fn get_random(buf: &mut [u8]) {
getrandom::getrandom(buf).unwrap() getrandom::getrandom(buf).unwrap()
} }
@ -19,235 +18,4 @@ pub fn get_random(buf: &mut [u8]) {
/// else /// else
/// process total_messages // 2 /// process total_messages // 2
const WATERLEVEL: usize = 10; const WATERLEVEL: usize = 10;
const DEFAULT_INBOX_CAPACITY: usize = 1_000; const DEFAULT_INBOX_CAPACITY: usize = 1_000;
pub trait Message: 'static + Sized + Clone + Send {}
impl<T: 'static + Sized + Clone + Send> Message for T {}
pub type Envelope = Box<dyn std::any::Any + Send>;
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<A>
where
A: ActorInterface,
{
_addr: ActorAddress,
inbox: Receiver<A::Incoming>,
inner: A,
}
/// Trait for type-erased actors
trait AnyActor: Send {
fn tick(&mut self, ctx: &Runtime);
}
impl<A> AnyActor for Actor<A>
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<M: Message> {
addr: ActorAddress,
inner: Receiver<M>,
}
impl<M: Message> Inbox<M> {
pub fn addr(&self) -> &ActorAddress {
&self.addr
}
pub fn try_recv(&self) -> Option<M> {
self.inner.try_recv()
}
}
#[derive(Debug, Default)]
pub enum RuntimeFlavor {
#[default]
SingleThreaded,
Multithreaded(usize),
}
pub struct Runtime {
flavor: RuntimeFlavor,
router: Router,
router_inbox: Sender<RouterMessage>,
actor_queue: ArrayQueue<Box<dyn AnyActor>>,
}
impl Runtime {
pub fn new(capacity: usize, flavor: Option<RuntimeFlavor>) -> 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<A: ActorInterface>(&self, actor: A) -> Result<ActorAddress, Error> {
let addr = {
let mut bytes = u64::to_le_bytes(0);
get_random(&mut bytes);
u64::from_le_bytes(bytes)
};
let inbox = Receiver::<A::Incoming>::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<M: Message>(&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<M: Message>(&self) -> Inbox<M> {
let addr = {
let mut bytes = u64::to_le_bytes(0);
get_random(&mut bytes);
u64::from_le_bytes(bytes)
};
let receiver = Receiver::<M>::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<M: Message> SenderT for Sender<M> {
fn try_send(&self, envelope: Envelope) {
if let Ok(msg) = envelope.downcast::<M>() {
let _ = Sender::try_send(self, *msg);
}
}
}
/// Internal messages for the Router's own inbox
pub enum RouterMessage {
/// register addrs <addr> with sender <sender>
AddAddr(ActorAddress, Box<dyn SenderT>),
/// remove an actor from the address book
RemoveAddr(ActorAddress),
/// send <msg> to <addr>
SendToAddr { addr: ActorAddress, msg: Envelope },
}
struct Router {
directory: HashMap<ActorAddress, Box<dyn SenderT>>,
inbox: Receiver<RouterMessage>,
}
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<RouterMessage> {
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);
}
}
}
}
}

77
src/router.rs Normal file
View file

@ -0,0 +1,77 @@
use std::collections::HashMap;
use crate::{WATERLEVEL, actor::{ActorAddress, Message}, ring_buffer::{Receiver, Sender}};
pub(crate) type Envelope = Box<dyn std::any::Any + Send>;
pub(crate) trait SenderT: Send {
fn try_send(&self, envelope: Envelope);
}
impl<M: Message> SenderT for Sender<M> {
fn try_send(&self, envelope: Envelope) {
if let Ok(msg) = envelope.downcast::<M>() {
let _ = Sender::try_send(self, *msg);
}
}
}
/// Internal messages for the Router's own inbox
pub(crate) enum RouterMessage {
/// register addrs <addr> with sender <sender>
AddAddr(ActorAddress, Box<dyn SenderT>),
/// remove an actor from the address book
RemoveAddr(ActorAddress),
/// send <msg> to <addr>
SendToAddr { addr: ActorAddress, msg: Envelope },
}
pub(crate) struct Router {
directory: HashMap<ActorAddress, Box<dyn SenderT>>,
inbox: Receiver<RouterMessage>,
}
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<RouterMessage> {
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);
}
}
}
}
}

113
src/runtime.rs Normal file
View file

@ -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<M: Message> {
addr: ActorAddress,
inner: Receiver<M>,
}
impl<M: Message> Inbox<M> {
pub fn addr(&self) -> &ActorAddress {
&self.addr
}
pub fn try_recv(&self) -> Option<M> {
self.inner.try_recv()
}
}
pub struct Runtime {
flavor: RuntimeFlavor,
router: Router,
router_inbox: Sender<RouterMessage>,
actor_queue: ArrayQueue<Box<dyn AnyActor>>,
}
impl Runtime {
pub fn new(capacity: usize, flavor: Option<RuntimeFlavor>) -> 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<A: ActorInterface>(&self, actor: A) -> Result<ActorAddress, Error> {
let addr = {
let mut bytes = u64::to_le_bytes(0);
get_random(&mut bytes);
u64::from_le_bytes(bytes)
};
let inbox = Receiver::<A::Incoming>::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<M: Message>(&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<M: Message>(&self) -> Inbox<M> {
let addr = {
let mut bytes = u64::to_le_bytes(0);
get_random(&mut bytes);
u64::from_le_bytes(bytes)
};
let receiver = Receiver::<M>::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,
}
}
}