diff --git a/Cargo.lock b/Cargo.lock index fd2fd7a..318f003 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2,6 +2,26 @@ # It is not intended for manual editing. version = 4 +[[package]] +name = "bytemuck" +version = "1.24.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fbdf580320f38b612e485521afda1ee26d10cc9884efaaa750d383e13e3c5f4" +dependencies = [ + "bytemuck_derive", +] + +[[package]] +name = "bytemuck_derive" +version = "1.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f9abbd1bc6865053c427f7198e6af43bfdedc55ab791faed4fbd361d789575ff" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "bytes" version = "1.11.0" @@ -48,6 +68,7 @@ dependencies = [ name = "swactor" version = "0.1.0" dependencies = [ + "bytemuck", "tokio", "tokio-util", ] diff --git a/Cargo.toml b/Cargo.toml index ba24c26..3ce83f1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -7,5 +7,6 @@ edition = "2024" crate-type = ["cdylib", "rlib"] [dependencies] +bytemuck = { version = "1.24.0", features = ["derive"] } tokio = { version = "1.48.0", features = ["rt", "macros"] } tokio-util = "0.7.17" diff --git a/DESIGN.md b/DESIGN.md new file mode 100644 index 0000000..bb03be1 --- /dev/null +++ b/DESIGN.md @@ -0,0 +1,35 @@ +# Design goals +Get as much usability and speed as possible while keeping line count low. Aim for no footguns, ability to plug in +logic easily, and run near anywhere. + +We are not building a new erlang/BEAM. Minimal feature set means spawning actor processes, not having supervisiors, lots of process +monitoring tools, prempting, etc. + +## Actor model + +An actor has: + + - An inbox: + this is a mpsc channel that the runtime/router dumps messages into and the actor consumes when the runtime loads it + + - an outbox channel connection: + this is a mpmc channel that is implemented by the runtime and router. Actors on this specific channel put responses and outgoing messages into this channel, to be routed to the given address. + + - a growable and mutable state: + An actor owns some, from the runtime perspective, type erased bytes. The actor when processing messages can access its own state, but no other task can. This includes viewing. + + - a set of functions for processing messages: + When the runtime loads the actor, it locks the inbox and attempts to process the messages therein. + +## Runtime and Router + +In order for an actor to consume and send messages, it is processed by a runtime. The runtime, in order to negotiate messages between +actors, possesses a router. + +A runtime has: + + - An actor processing thread(s): + the processor will mark an actor as busy, load its state and inbox, and begin consuming messages from the inbox. The number of messages consumed is determined by the runtime + + - A message router: + the router is responsible for ensuring messages posted by actors get delivered to the appropriate inbox. \ No newline at end of file diff --git a/examples/hello.rs b/examples/hello.rs index 40440e1..6fa5b1b 100644 --- a/examples/hello.rs +++ b/examples/hello.rs @@ -1,46 +1,63 @@ -use swactor::Actor; -use tokio::sync::oneshot; -pub enum GreeterMessage { - Name(String), -} +// use std::sync::{Arc, atomic::AtomicBool}; -pub enum GreeterResponse { - Hello(String), -} +// use swactor::error::*; +// use tokio::task::JoinHandle; -pub struct Greeter; +// struct ActorInbox { +// _guard: AtomicBool +// } -impl Actor for Greeter { - type Message = GreeterMessage; - type Response = GreeterResponse; +// struct GenericGreeter { +// _guard: Arc, +// inbox: Vec, +// outbox: Vec, +// } - fn handle_message(&self, msg: Self::Message, tx: oneshot::Sender) { - let rep = match msg { - GreeterMessage::Name(name) => GreeterResponse::Hello(format!("Hello, {name}!")), - }; +// impl GenericGreeter { +// pub fn new() -> Self { +// Self { +// _guard: Arc::new(AtomicBool::new(false)), +// inbox: Vec::new(), +// outbox: Vec::new(), +// } +// } +// } - if let Err(_) = tx.send(rep) { - // Greeter is not responsible for a dropped Receiver - } - } -} -fn main() { - let rt = tokio::runtime::Builder::new_current_thread() - .build() - .expect("failed to build runtime"); - let greeter = Greeter.spawn(&rt); +// pub enum GreeterMessage { +// Name(String), +// } - let response = rt - .block_on(async move { - greeter - .send(GreeterMessage::Name("world".to_string())) - .await - }) - .expect("failed to get respose"); +// pub enum GreeterResponse { +// Hello(String), +// } - match response { - GreeterResponse::Hello(hello) => println!("{hello}"), - } -} +// pub struct Greeter; + + +// fn main() { +// let rt = tokio::runtime::Builder::new_current_thread() +// .build() +// .expect("failed to build runtime"); + +// let greet = GenericGreeter::spawn(&rt); + +// let res = rt.block_on(async {greet.await}).expect("runtime error").expect("greeter error"); + +// println!("Success!"); + +// // let greeter = Greeter.spawn(&rt); + +// // let response = rt +// // .block_on(async move { +// // greeter +// // .send(GreeterMessage::Name("world".to_string())) +// // .await +// // }) +// // .expect("failed to get respose"); + +// // match response { +// // GreeterResponse::Hello(hello) => println!("{hello}"), +// // } +// } diff --git a/src/kimi.rs b/src/kimi.rs new file mode 100644 index 0000000..e69de29 diff --git a/src/lib.bak.rs b/src/lib.bak.rs new file mode 100644 index 0000000..1641435 --- /dev/null +++ b/src/lib.bak.rs @@ -0,0 +1,118 @@ +pub mod error; + +/// Public export as the oneshot channel is in the `Actor` trait signature +pub use tokio::sync::oneshot; + +use tokio::{sync::mpsc, task::JoinHandle}; +use tokio_util::sync::CancellationToken; + +use crate::error::{Result, convert_err}; + +const DEFAULT_CHANNEL_SIZE: usize = 100; + +type ActorRequest = ( + ::Message, + oneshot::Sender<::Response>, +); + +/// Wrapper defining the transmission end of a Request/Response channel with an `Actor` +pub struct ActorRequestSender(mpsc::Sender>); + +impl ActorRequestSender { + pub async fn send(&self, request: A::Message) -> Result { + let (tx, rx) = oneshot::channel::(); + self.0.send((request, tx)).await.map_err(convert_err)?; + + rx.await.map_err(convert_err) + } +} + +impl From>> for ActorRequestSender { + fn from(value: mpsc::Sender>) -> Self { + Self(value) + } +} + +impl Clone for ActorRequestSender { + fn clone(&self) -> Self { + Self(self.0.clone()) + } +} + +/// Combined 'JoinHandle' to await the actor process and 'Sender' for communication +pub struct Handle +where + A: Actor, +{ + cancel_token: CancellationToken, + tx: ActorRequestSender, + /// the task drops when the `JoinHandle` does, so be careful with the `Handle` + _handle: JoinHandle>, + // to prevent accidental swaps, strongly type the handle + _type: std::marker::PhantomData, +} + +impl Handle { + /// Send a message to the spawned `Actor` task and get a response corresponding to the `Actor::Response` type + pub async fn send(&self, msg: A::Message) -> Result { + self.tx.send(msg).await + } + + /// Get a cloned sender for messaging the `Actor` this handle is for + pub fn get_connection(&self) -> ActorRequestSender { + self.tx.clone() + } +} + +impl Drop for Handle { + fn drop(&mut self) { + self.cancel_token.cancel(); + } +} + +/// Primary trait defining an `Actor` capable of receiving, processing, and transmitting messages +pub trait Actor: Send + Sized + 'static { + /// The type for messages received by this `Actor` + type Message: Send; + /// The type for responses given by this actor when called from `Handle::send(..)` + type Response: Send; + + /// Inner method that defines actor behavior + fn handle_message(&self, msg: Self::Message, tx: oneshot::Sender); + + /// Spawns the `Actor` utilizing the given runtime context + /// Only `tokio` runtime is accepted for now + fn spawn(self, ctx: &tokio::runtime::Runtime) -> Handle { + let cancel_token = CancellationToken::new(); + let cancel = cancel_token.clone(); + + let (tx, mut rx) = + mpsc::channel::<(Self::Message, oneshot::Sender)>(DEFAULT_CHANNEL_SIZE); + let handle = ctx.spawn(async move { + let mut res = Ok(()); + loop { + tokio::select! { + _ = cancel.cancelled() => { + break; + }, + + msg = rx.recv() => { + match msg { + Some(m) => { self.handle_message(m.0, m.1); }, + None => {res = Err(format!("Sender handle was dropped without calling cancel!").into()); break; }, + } + } + }; + } + + res + }); + + Handle { + cancel_token, + _handle: handle, + tx: tx.into(), + _type: std::marker::PhantomData::, + } + } +} diff --git a/src/lib.rs b/src/lib.rs index 1641435..dbe780e 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,118 +1,164 @@ pub mod error; -/// Public export as the oneshot channel is in the `Actor` trait signature -pub use tokio::sync::oneshot; +use std::{ + marker::PhantomData, + sync::{Mutex, mpsc::TryRecvError}, +}; -use tokio::{sync::mpsc, task::JoinHandle}; -use tokio_util::sync::CancellationToken; +use crate::error::{Error, Result, convert_err}; -use crate::error::{Result, convert_err}; +use std::sync::mpsc; -const DEFAULT_CHANNEL_SIZE: usize = 100; +use bytemuck::{Pod, Zeroable}; -type ActorRequest = ( - ::Message, - oneshot::Sender<::Response>, -); - -/// Wrapper defining the transmission end of a Request/Response channel with an `Actor` -pub struct ActorRequestSender(mpsc::Sender>); - -impl ActorRequestSender { - pub async fn send(&self, request: A::Message) -> Result { - let (tx, rx) = oneshot::channel::(); - self.0.send((request, tx)).await.map_err(convert_err)?; - - rx.await.map_err(convert_err) - } +#[repr(C)] +#[derive(Copy, Clone, Pod, Zeroable)] +struct GreeterState { + pub num_greeted: usize, } -impl From>> for ActorRequestSender { - fn from(value: mpsc::Sender>) -> Self { - Self(value) - } +enum GreetMessage { + Name(String), } -impl Clone for ActorRequestSender { - fn clone(&self) -> Self { - Self(self.0.clone()) - } +enum GreetResponse { + Greeting(String), } -/// Combined 'JoinHandle' to await the actor process and 'Sender' for communication -pub struct Handle -where - A: Actor, -{ - cancel_token: CancellationToken, - tx: ActorRequestSender, - /// the task drops when the `JoinHandle` does, so be careful with the `Handle` - _handle: JoinHandle>, - // to prevent accidental swaps, strongly type the handle - _type: std::marker::PhantomData, +type GreeterId = u64; + +struct Greeter { + id: GreeterId, + inbox: mpsc::Receiver, + outbox: mpsc::Sender, + state: GreeterState, } -impl Handle { - /// Send a message to the spawned `Actor` task and get a response corresponding to the `Actor::Response` type - pub async fn send(&self, msg: A::Message) -> Result { - self.tx.send(msg).await +impl Greeter { + pub fn new( + id: GreeterId, + inbox: mpsc::Receiver, + outbox: mpsc::Sender, + ) -> Self { + Self { + id, + inbox, + outbox, + state: GreeterState { num_greeted: 0 }, + } } - /// Get a cloned sender for messaging the `Actor` this handle is for - pub fn get_connection(&self) -> ActorRequestSender { - self.tx.clone() + pub fn id(&self) -> GreeterId { + self.id } -} -impl Drop for Handle { - fn drop(&mut self) { - self.cancel_token.cancel(); - } -} - -/// Primary trait defining an `Actor` capable of receiving, processing, and transmitting messages -pub trait Actor: Send + Sized + 'static { - /// The type for messages received by this `Actor` - type Message: Send; - /// The type for responses given by this actor when called from `Handle::send(..)` - type Response: Send; - - /// Inner method that defines actor behavior - fn handle_message(&self, msg: Self::Message, tx: oneshot::Sender); - - /// Spawns the `Actor` utilizing the given runtime context - /// Only `tokio` runtime is accepted for now - fn spawn(self, ctx: &tokio::runtime::Runtime) -> Handle { - let cancel_token = CancellationToken::new(); - let cancel = cancel_token.clone(); - - let (tx, mut rx) = - mpsc::channel::<(Self::Message, oneshot::Sender)>(DEFAULT_CHANNEL_SIZE); - let handle = ctx.spawn(async move { - let mut res = Ok(()); - loop { - tokio::select! { - _ = cancel.cancelled() => { - break; - }, - - msg = rx.recv() => { - match msg { - Some(m) => { self.handle_message(m.0, m.1); }, - None => {res = Err(format!("Sender handle was dropped without calling cancel!").into()); break; }, - } - } - }; + pub fn process_message(&mut self) -> Result<()> { + match self.inbox.try_recv() { + Ok(m) => { + let GreetMessage::Name(n) = m; + self.outbox + .send(GreetResponse::Greeting(format!("Hello, {n}!"))) + .map_err(convert_err)?; + self.state.num_greeted += 1; } + Err(e) => match e { + TryRecvError::Empty => return Ok(()), + TryRecvError::Disconnected => { + return Err("Outbox has been disconnected, actor in an improper state".into()); + } + }, + } - res - }); + Ok(()) + } +} - Handle { - cancel_token, - _handle: handle, - tx: tx.into(), - _type: std::marker::PhantomData::, +use std::collections::HashMap; + +struct Router { + address_book: HashMap>, + next_id: GreeterId, +} + +impl Router { + pub fn new() -> Self { + Self { + address_book: HashMap::new(), + next_id: 0, + } + } + + pub fn register(&mut self, sender: mpsc::Sender) -> GreeterId { + let id = self.next_id; + self.next_id += 1; + self.address_book.insert(id, sender); + id + } + + pub fn unregister(&mut self, id: GreeterId) -> Option> { + self.address_book.remove(&id) + } + + pub fn get_sender(&self, id: GreeterId) -> Option<&mpsc::Sender> { + self.address_book.get(&id) + } + + pub fn send(&self, id: GreeterId, message: GreetMessage) -> Result<()> { + match self.address_book.get(&id) { + Some(sender) => sender.send(message).map_err(convert_err), + None => Err(format!("No sender found for id {}", id).into()), } } } + +struct Runtime { + router: Router, + greeters: Vec, + response_rx: mpsc::Receiver, + response_tx: mpsc::Sender, +} + +impl Runtime { + pub fn new() -> Self { + let (response_tx, response_rx) = mpsc::channel(); + Self { + router: Router::new(), + greeters: Vec::new(), + response_rx, + response_tx, + } + } + + pub fn spawn_greeter(&mut self) -> GreeterId { + let (inbox_tx, inbox_rx) = mpsc::channel(); + let id = self.router.register(inbox_tx); + let greeter = Greeter::new(id, inbox_rx, self.response_tx.clone()); + self.greeters.push(greeter); + id + } + + pub fn send_message(&self, id: GreeterId, message: GreetMessage) -> Result<()> { + self.router.send(id, message) + } + + pub fn tick(&mut self) -> Result<()> { + for greeter in &mut self.greeters { + greeter.process_message()?; + } + Ok(()) + } + + pub fn try_recv_response(&self) -> Option { + self.response_rx.try_recv().ok() + } + + pub fn run_until_idle(&mut self) -> Result<()> { + loop { + self.tick()?; + if self.response_rx.try_recv().is_err() { + break; + } + } + Ok(()) + } +} \ No newline at end of file