Compare commits
No commits in common. "d504377ba984c5a6dc4c5d0e29e09eecd2c93544" and "55077dbf57d23c60d6c9f1d817f99713457329d4" have entirely different histories.
d504377ba9
...
55077dbf57
9 changed files with 248 additions and 723 deletions
|
|
@ -1,7 +1,4 @@
|
||||||
use swactor::{
|
use swactor::{ActorAddress, ActorInterface, Message, Runtime, RuntimeFlavor};
|
||||||
actor::{ActorAddress, ActorInterface},
|
|
||||||
runtime::{Context, Runtime, RuntimeFlavor},
|
|
||||||
};
|
|
||||||
|
|
||||||
#[derive(Debug, Default)]
|
#[derive(Debug, Default)]
|
||||||
struct Greeter {
|
struct Greeter {
|
||||||
|
|
@ -16,15 +13,13 @@ struct GreetMessage {
|
||||||
/// who do we send out greeting back to?
|
/// who do we send out greeting back to?
|
||||||
return_addr: ActorAddress,
|
return_addr: ActorAddress,
|
||||||
}
|
}
|
||||||
|
impl Message for GreetMessage {}
|
||||||
#[derive(Debug, Default, Clone)]
|
|
||||||
struct GreetResponse(String);
|
|
||||||
|
|
||||||
impl ActorInterface for Greeter {
|
impl ActorInterface for Greeter {
|
||||||
type Incoming = GreetMessage;
|
type Incoming = GreetMessage;
|
||||||
type Response = GreetResponse;
|
type Response = GreetResponse;
|
||||||
|
|
||||||
fn handle(&mut self, ctx: &Context, msg: GreetMessage) {
|
fn handle(&mut self, ctx: &Runtime, msg: GreetMessage) {
|
||||||
let res = GreetResponse(format!("Hello, {}!", msg.who));
|
let res = GreetResponse(format!("Hello, {}!", msg.who));
|
||||||
self.num_greeted += 1;
|
self.num_greeted += 1;
|
||||||
if let Err(_) = ctx.send_to(msg.return_addr, res) {
|
if let Err(_) = ctx.send_to(msg.return_addr, res) {
|
||||||
|
|
@ -34,8 +29,12 @@ impl ActorInterface for Greeter {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Default, Clone)]
|
||||||
|
struct GreetResponse(String);
|
||||||
|
impl Message for GreetResponse {}
|
||||||
|
|
||||||
fn main() {
|
fn main() {
|
||||||
let rt = Runtime::new(100, RuntimeFlavor::SingleThreaded);
|
let mut rt = Runtime::new(100, Some(RuntimeFlavor::SingleThreaded));
|
||||||
let addr = rt
|
let addr = rt
|
||||||
.spawn(Greeter::default())
|
.spawn(Greeter::default())
|
||||||
.expect("failed to spawn greeter");
|
.expect("failed to spawn greeter");
|
||||||
|
|
|
||||||
|
|
@ -1,68 +0,0 @@
|
||||||
use swactor::{
|
|
||||||
actor::{ActorAddress, ActorInterface},
|
|
||||||
runtime::{Context, Inbox, Runtime, RuntimeFlavor},
|
|
||||||
};
|
|
||||||
|
|
||||||
#[derive(Debug, Default, Clone)]
|
|
||||||
pub struct RingMessage {
|
|
||||||
count: usize,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl RingMessage {
|
|
||||||
pub fn next(self) -> Self {
|
|
||||||
Self {
|
|
||||||
count: self.count + 1,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Debug, Default)]
|
|
||||||
struct RingActor {
|
|
||||||
next: ActorAddress,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl RingActor {
|
|
||||||
pub fn new(next: ActorAddress) -> Self {
|
|
||||||
Self { next }
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl ActorInterface for RingActor {
|
|
||||||
type Incoming = RingMessage;
|
|
||||||
type Response = ();
|
|
||||||
fn handle(&mut self, ctx: &Context, msg: Self::Incoming) {
|
|
||||||
if let Err(_) = ctx.send_to(self.next, msg.next()) {
|
|
||||||
// do nothing
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn main() {
|
|
||||||
let rt = Runtime::new(10_000, RuntimeFlavor::SingleThreaded);
|
|
||||||
let inbox: Inbox<RingMessage> = rt.new_inbox();
|
|
||||||
|
|
||||||
let mut next = rt
|
|
||||||
.spawn(RingActor::new(*inbox.addr()))
|
|
||||||
.expect("failed to spawn");
|
|
||||||
for _ in 0..500 {
|
|
||||||
let new = rt.spawn(RingActor::new(next)).expect("failed to spawn");
|
|
||||||
next = new;
|
|
||||||
}
|
|
||||||
rt.send_to(next, RingMessage { count: 0 })
|
|
||||||
.expect("failed to start message ring");
|
|
||||||
|
|
||||||
let msg: RingMessage;
|
|
||||||
loop {
|
|
||||||
match inbox.try_recv() {
|
|
||||||
Some(m) => {
|
|
||||||
msg = m;
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
None => {
|
|
||||||
rt.tick();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
println!("{msg:?}")
|
|
||||||
}
|
|
||||||
62
src/actor.rs
62
src/actor.rs
|
|
@ -1,62 +0,0 @@
|
||||||
use crate::{ring_buffer::Receiver, runtime::Context, WATERLEVEL};
|
|
||||||
|
|
||||||
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: &Context, msg: Self::Incoming);
|
|
||||||
}
|
|
||||||
|
|
||||||
pub type ActorAddress = u64;
|
|
||||||
|
|
||||||
/// FIXME: If we never have the actor struct reference its own address, should
|
|
||||||
/// we even include it as a variable here? We could instead grab this information
|
|
||||||
/// from the runtime or router.
|
|
||||||
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: &Context);
|
|
||||||
}
|
|
||||||
|
|
||||||
impl<A> AnyActor for Actor<A>
|
|
||||||
where
|
|
||||||
A: ActorInterface,
|
|
||||||
{
|
|
||||||
fn tick(&mut self, ctx: &Context) {
|
|
||||||
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"
|
|
||||||
),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -1,5 +1,6 @@
|
||||||
#[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())
|
||||||
}
|
}
|
||||||
|
|
|
||||||
255
src/lib.rs
255
src/lib.rs
|
|
@ -1,28 +1,251 @@
|
||||||
pub mod actor;
|
|
||||||
|
|
||||||
pub(crate) mod error;
|
|
||||||
pub use error::Error;
|
|
||||||
|
|
||||||
mod ring_buffer;
|
mod ring_buffer;
|
||||||
mod router;
|
|
||||||
pub mod runtime;
|
|
||||||
|
|
||||||
// Re-export commonly used types
|
use std::collections::HashMap;
|
||||||
pub use actor::{ActorAddress, ActorInterface, Message};
|
|
||||||
pub use runtime::{Context, Inbox, Runtime, RuntimeFlavor};
|
use crossbeam_queue::ArrayQueue;
|
||||||
|
use ring_buffer::{Receiver, Sender};
|
||||||
|
|
||||||
|
pub mod error;
|
||||||
|
use error::Error;
|
||||||
|
|
||||||
#[cfg(feature = "getrandom")]
|
#[cfg(feature = "getrandom")]
|
||||||
pub(crate) fn get_random(buf: &mut [u8]) {
|
pub fn get_random(buf: &mut [u8]) {
|
||||||
getrandom::getrandom(buf).unwrap()
|
getrandom::getrandom(buf).unwrap()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The strategy for message processing is such:
|
/// The strategy for message processing is such:
|
||||||
///
|
|
||||||
/// ```ignore
|
|
||||||
/// if total_messages < WATERLEVEL:
|
/// if total_messages < WATERLEVEL:
|
||||||
/// process all
|
/// process all
|
||||||
/// else
|
/// else
|
||||||
/// process total_messages >> 1
|
/// process total_messages // 2
|
||||||
/// ```
|
|
||||||
const WATERLEVEL: usize = 10;
|
const WATERLEVEL: usize = 10;
|
||||||
const DEFAULT_INBOX_CAPACITY: usize = 1_000;
|
|
||||||
|
const DEFAULT_INBOX_CAPACITY: usize = 100;
|
||||||
|
|
||||||
|
pub trait Message: 'static + Sized + Clone + Send {}
|
||||||
|
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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -48,17 +48,6 @@ pub(crate) struct Sender<T> {
|
||||||
queue: Arc<ArrayQueue<T>>,
|
queue: Arc<ArrayQueue<T>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// FIXME: I don't like this. Why do we need to clone the Sender
|
|
||||||
/// Because there are no guarentees on the existence of the Receiver
|
|
||||||
/// we need to be very careful about passing around access to the buffer.
|
|
||||||
impl<T> Clone for Sender<T> {
|
|
||||||
fn clone(&self) -> Self {
|
|
||||||
Self {
|
|
||||||
queue: Arc::clone(&self.queue),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl<T> Sender<T> {
|
impl<T> Sender<T> {
|
||||||
/// Attempt to push a value to the queue. Returns Err(value) if the queue is full.
|
/// Attempt to push a value to the queue. Returns Err(value) if the queue is full.
|
||||||
pub fn try_send(&self, value: T) -> Result<(), T> {
|
pub fn try_send(&self, value: T) -> Result<(), T> {
|
||||||
|
|
|
||||||
|
|
@ -1,86 +0,0 @@
|
||||||
use std::collections::HashMap;
|
|
||||||
|
|
||||||
use crate::{
|
|
||||||
actor::{ActorAddress, Message},
|
|
||||||
ring_buffer::{Receiver, Sender},
|
|
||||||
WATERLEVEL,
|
|
||||||
};
|
|
||||||
|
|
||||||
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>),
|
|
||||||
|
|
||||||
/// FIXME: this will be active when we allow actors to shut themselves
|
|
||||||
/// down. For now, disable the warning.
|
|
||||||
#[allow(dead_code)]
|
|
||||||
/// 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);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
316
src/runtime.rs
316
src/runtime.rs
|
|
@ -1,316 +0,0 @@
|
||||||
use std::sync::atomic::{AtomicBool, Ordering};
|
|
||||||
use std::sync::{Arc, Mutex};
|
|
||||||
use std::thread::{self, JoinHandle};
|
|
||||||
|
|
||||||
use crossbeam_queue::ArrayQueue;
|
|
||||||
|
|
||||||
use crate::{
|
|
||||||
actor::{Actor, ActorAddress, ActorInterface, AnyActor, Message},
|
|
||||||
get_random,
|
|
||||||
ring_buffer::{Receiver, Sender},
|
|
||||||
router::{Envelope, Router, RouterMessage},
|
|
||||||
Error, DEFAULT_INBOX_CAPACITY,
|
|
||||||
};
|
|
||||||
|
|
||||||
#[derive(Debug, Clone, Default)]
|
|
||||||
pub enum RuntimeFlavor {
|
|
||||||
#[default]
|
|
||||||
SingleThreaded,
|
|
||||||
Multithreaded {
|
|
||||||
workers: 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()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A lightweight handle for sending messages to actors
|
|
||||||
/// This is what actors receive in their handle() method
|
|
||||||
#[derive(Clone)]
|
|
||||||
pub struct Context {
|
|
||||||
router_inbox: Sender<RouterMessage>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Context {
|
|
||||||
/// Send a message to an actor address
|
|
||||||
pub fn send_to<M: Message>(&self, addr: ActorAddress, msg: M) -> Result<(), Error> {
|
|
||||||
let envelope: Envelope = Box::new(msg);
|
|
||||||
self.router_inbox
|
|
||||||
.try_send(RouterMessage::SendToAddr {
|
|
||||||
addr,
|
|
||||||
msg: envelope,
|
|
||||||
})
|
|
||||||
.map_err(|_| Error::from("Failed to send message: router inbox full"))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Shared runtime state, wrapped in Arc for thread sharing
|
|
||||||
/// FIXME: We check the `running` variable too often. Should
|
|
||||||
/// push it somewhere where it is infreqently checked, and instead
|
|
||||||
/// have shutdown logic for everything else.
|
|
||||||
struct RuntimeInner {
|
|
||||||
running: AtomicBool,
|
|
||||||
router_inbox: Sender<RouterMessage>,
|
|
||||||
actor_queue: ArrayQueue<Box<dyn AnyActor>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Thread handles for the multithreaded runtime
|
|
||||||
struct RuntimeHandles {
|
|
||||||
workers: Vec<JoinHandle<()>>,
|
|
||||||
router_thread: Option<JoinHandle<()>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The main runtime for executing actors
|
|
||||||
pub struct Runtime {
|
|
||||||
inner: Arc<RuntimeInner>,
|
|
||||||
flavor: RuntimeFlavor,
|
|
||||||
/// Router is only accessed from a single thread (either main or dedicated router thread)
|
|
||||||
///
|
|
||||||
/// FIXME: If single threaded, why do we have a mutex
|
|
||||||
router: Mutex<Router>,
|
|
||||||
/// Thread handles, created lazily when run() is called
|
|
||||||
handles: Mutex<Option<RuntimeHandles>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Runtime {
|
|
||||||
/// Create a new runtime with given actor queue capacity and flavor
|
|
||||||
///
|
|
||||||
/// FIXME: I don't like this interface. Maybe a builder or config pattern. I shouldnt
|
|
||||||
/// have to read code to understand what these variable names are.
|
|
||||||
pub fn new(capacity: usize, flavor: RuntimeFlavor) -> Self {
|
|
||||||
// FIXME: Avoid hard coded defaults, or at least put them all in one place
|
|
||||||
let router = Router::new(DEFAULT_INBOX_CAPACITY);
|
|
||||||
let router_inbox = router.new_sender();
|
|
||||||
Self {
|
|
||||||
inner: Arc::new(RuntimeInner {
|
|
||||||
running: AtomicBool::new(false),
|
|
||||||
router_inbox,
|
|
||||||
actor_queue: ArrayQueue::new(capacity),
|
|
||||||
}),
|
|
||||||
flavor,
|
|
||||||
router: Mutex::new(router),
|
|
||||||
handles: Mutex::new(None),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Get a context handle for sending messages
|
|
||||||
///
|
|
||||||
/// FIXME: Do we need all this indirection?
|
|
||||||
pub fn context(&self) -> Context {
|
|
||||||
Context {
|
|
||||||
router_inbox: self.inner.router_inbox.clone(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Spawn an actor, returns its address
|
|
||||||
pub fn spawn<A: ActorInterface>(&self, actor: A) -> Result<ActorAddress, Error> {
|
|
||||||
// FIXME: Figure out what to do with this.
|
|
||||||
// Its distracting to include here, but not used elsewhere for now.
|
|
||||||
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
|
|
||||||
.inner
|
|
||||||
.router_inbox
|
|
||||||
.try_send(RouterMessage::AddAddr(addr, Box::new(sender)));
|
|
||||||
|
|
||||||
self.inner
|
|
||||||
.actor_queue
|
|
||||||
.push(Box::new(Actor::new(addr, inbox, actor)))
|
|
||||||
.map_err(|_| Error::from("Runtime error: Failed to spawn actor."))?;
|
|
||||||
|
|
||||||
Ok(addr)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Send a message to an actor address
|
|
||||||
pub fn send_to<M: Message>(&self, addr: ActorAddress, msg: M) -> Result<(), Error> {
|
|
||||||
self.context().send_to(addr, msg)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Create an external inbox for receiving messages outside actors
|
|
||||||
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
|
|
||||||
.inner
|
|
||||||
.router_inbox
|
|
||||||
.try_send(RouterMessage::AddAddr(addr, Box::new(sender)));
|
|
||||||
Inbox {
|
|
||||||
addr,
|
|
||||||
inner: receiver,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Process one actor tick + router messages
|
|
||||||
/// Works in both single/multi mode (useful for testing and fine-grained control)
|
|
||||||
///
|
|
||||||
/// FIXME: `tick` does not make sense in multithreaded context. Add a check to ensure single threaded
|
|
||||||
pub fn tick(&self) {
|
|
||||||
let ctx = self.context();
|
|
||||||
|
|
||||||
if let Some(mut actor) = self.inner.actor_queue.pop() {
|
|
||||||
actor.tick(&ctx);
|
|
||||||
let _ = self.inner.actor_queue.push(actor);
|
|
||||||
}
|
|
||||||
|
|
||||||
// In single-threaded mode, also tick the router
|
|
||||||
if matches!(self.flavor, RuntimeFlavor::SingleThreaded)
|
|
||||||
&& let Ok(mut router) = self.router.lock()
|
|
||||||
{
|
|
||||||
router.tick();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Run the runtime (blocking)
|
|
||||||
/// - SingleThreaded: runs in current thread until shutdown
|
|
||||||
/// - Multithreaded: spawns workers + router thread, blocks until shutdown
|
|
||||||
pub fn run(&self) {
|
|
||||||
self.inner.running.store(true, Ordering::Release);
|
|
||||||
|
|
||||||
match &self.flavor {
|
|
||||||
RuntimeFlavor::SingleThreaded => {
|
|
||||||
self.run_single_threaded();
|
|
||||||
}
|
|
||||||
RuntimeFlavor::Multithreaded { workers } => {
|
|
||||||
self.run_multi_threaded(*workers);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// FIXME: I don't like this loop. We can send shutdown signals from the process
|
|
||||||
/// that calls the actor runtime instead. It also does not make sense to have this
|
|
||||||
/// around for single threaded runtimes (people can instead loop over `runtime.tick()`).
|
|
||||||
fn run_single_threaded(&self) {
|
|
||||||
let ctx = self.context();
|
|
||||||
|
|
||||||
while self.inner.running.load(Ordering::Relaxed) {
|
|
||||||
// Pop actor, tick it, push it back
|
|
||||||
if let Some(mut actor) = self.inner.actor_queue.pop() {
|
|
||||||
actor.tick(&ctx);
|
|
||||||
let _ = self.inner.actor_queue.push(actor);
|
|
||||||
}
|
|
||||||
|
|
||||||
// Tick the router
|
|
||||||
if let Ok(mut router) = self.router.lock() {
|
|
||||||
router.tick();
|
|
||||||
}
|
|
||||||
|
|
||||||
thread::yield_now();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn run_multi_threaded(&self, num_workers: usize) {
|
|
||||||
// Take ownership of router for the dedicated router thread
|
|
||||||
// FIXME: this is weird and concerning. Introduces a class of logic errors
|
|
||||||
// wherein we try and call a dummy router.
|
|
||||||
let router = {
|
|
||||||
let mut guard = self.router.lock().unwrap();
|
|
||||||
std::mem::replace(&mut *guard, Router::new(1)) // placeholder
|
|
||||||
};
|
|
||||||
|
|
||||||
// Spawn dedicated router thread
|
|
||||||
// FIXME: what is this, these variables are named terribly
|
|
||||||
let router_running = Arc::clone(&self.inner);
|
|
||||||
let router_handle = {
|
|
||||||
let running = router_running;
|
|
||||||
thread::spawn(move || {
|
|
||||||
router_loop(router, running);
|
|
||||||
})
|
|
||||||
};
|
|
||||||
|
|
||||||
// Spawn worker threads
|
|
||||||
let worker_handles: Vec<_> = (0..num_workers)
|
|
||||||
.map(|_| {
|
|
||||||
let inner = Arc::clone(&self.inner);
|
|
||||||
thread::spawn(move || {
|
|
||||||
worker_loop(inner);
|
|
||||||
})
|
|
||||||
})
|
|
||||||
.collect();
|
|
||||||
|
|
||||||
// Store handles
|
|
||||||
*self.handles.lock().unwrap() = Some(RuntimeHandles {
|
|
||||||
workers: worker_handles,
|
|
||||||
router_thread: Some(router_handle),
|
|
||||||
});
|
|
||||||
|
|
||||||
// Block until shutdown - wait for all threads to complete
|
|
||||||
self.wait_for_shutdown();
|
|
||||||
}
|
|
||||||
|
|
||||||
fn wait_for_shutdown(&self) {
|
|
||||||
// Wait for the running flag to be set to false, then join threads
|
|
||||||
while self.inner.running.load(Ordering::Relaxed) {
|
|
||||||
thread::yield_now();
|
|
||||||
}
|
|
||||||
|
|
||||||
// Join all threads
|
|
||||||
let handles = self.handles.lock().unwrap().take();
|
|
||||||
if let Some(h) = handles {
|
|
||||||
for worker in h.workers {
|
|
||||||
let _ = worker.join();
|
|
||||||
}
|
|
||||||
if let Some(rt) = h.router_thread {
|
|
||||||
let _ = rt.join();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Signal all workers to stop
|
|
||||||
pub fn shutdown(&self) {
|
|
||||||
self.inner.running.store(false, Ordering::Release);
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Check if runtime is still active
|
|
||||||
pub fn is_running(&self) -> bool {
|
|
||||||
self.inner.running.load(Ordering::Acquire)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Worker thread loop - processes actors from the shared queue
|
|
||||||
fn worker_loop(inner: Arc<RuntimeInner>) {
|
|
||||||
let ctx = Context {
|
|
||||||
router_inbox: inner.router_inbox.clone(),
|
|
||||||
};
|
|
||||||
|
|
||||||
while inner.running.load(Ordering::Relaxed) {
|
|
||||||
if let Some(mut actor) = inner.actor_queue.pop() {
|
|
||||||
actor.tick(&ctx);
|
|
||||||
let _ = inner.actor_queue.push(actor);
|
|
||||||
} else {
|
|
||||||
thread::yield_now();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Router thread loop - processes router messages
|
|
||||||
fn router_loop(mut router: Router, inner: Arc<RuntimeInner>) {
|
|
||||||
while inner.running.load(Ordering::Relaxed) {
|
|
||||||
router.tick();
|
|
||||||
thread::yield_now();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -1,155 +0,0 @@
|
||||||
use std::sync::Arc;
|
|
||||||
use std::thread;
|
|
||||||
use std::time::Duration;
|
|
||||||
|
|
||||||
use swactor::{
|
|
||||||
actor::{ActorAddress, ActorInterface},
|
|
||||||
runtime::{Context, Inbox, Runtime, RuntimeFlavor},
|
|
||||||
};
|
|
||||||
|
|
||||||
// ============================================================================
|
|
||||||
// Test Helpers
|
|
||||||
// ============================================================================
|
|
||||||
|
|
||||||
#[derive(Clone)]
|
|
||||||
struct PingMessage {
|
|
||||||
reply_to: ActorAddress,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Clone)]
|
|
||||||
struct PongMessage;
|
|
||||||
|
|
||||||
struct PongActor;
|
|
||||||
|
|
||||||
impl ActorInterface for PongActor {
|
|
||||||
type Incoming = PingMessage;
|
|
||||||
type Response = PongMessage;
|
|
||||||
|
|
||||||
fn handle(&mut self, ctx: &Context, msg: PingMessage) {
|
|
||||||
let _ = ctx.send_to(msg.reply_to, PongMessage);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// An actor that forwards messages to another address
|
|
||||||
struct ForwarderActor {
|
|
||||||
target: ActorAddress,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Clone)]
|
|
||||||
struct ForwardMessage(usize);
|
|
||||||
|
|
||||||
impl ActorInterface for ForwarderActor {
|
|
||||||
type Incoming = ForwardMessage;
|
|
||||||
type Response = ();
|
|
||||||
|
|
||||||
fn handle(&mut self, ctx: &Context, msg: ForwardMessage) {
|
|
||||||
let _ = ctx.send_to(self.target, msg);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_single_threaded_ping_pong() {
|
|
||||||
let rt = Runtime::new(100, RuntimeFlavor::SingleThreaded);
|
|
||||||
let inbox: Inbox<PongMessage> = rt.new_inbox();
|
|
||||||
|
|
||||||
let pong_addr = rt.spawn(PongActor).expect("spawn pong");
|
|
||||||
|
|
||||||
// Send ping
|
|
||||||
rt.send_to(
|
|
||||||
pong_addr,
|
|
||||||
PingMessage {
|
|
||||||
reply_to: *inbox.addr(),
|
|
||||||
},
|
|
||||||
)
|
|
||||||
.unwrap();
|
|
||||||
|
|
||||||
// Tick until we get a response
|
|
||||||
for _ in 0..10 {
|
|
||||||
rt.tick();
|
|
||||||
if inbox.try_recv().is_some() {
|
|
||||||
return; // Success!
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
panic!("Did not receive pong response");
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_single_threaded_message_chain() {
|
|
||||||
let rt = Runtime::new(100, RuntimeFlavor::SingleThreaded);
|
|
||||||
let inbox: Inbox<ForwardMessage> = rt.new_inbox();
|
|
||||||
|
|
||||||
// Create a chain: A -> B -> C -> inbox
|
|
||||||
let c_addr = rt
|
|
||||||
.spawn(ForwarderActor {
|
|
||||||
target: *inbox.addr(),
|
|
||||||
})
|
|
||||||
.unwrap();
|
|
||||||
let b_addr = rt.spawn(ForwarderActor { target: c_addr }).unwrap();
|
|
||||||
let a_addr = rt.spawn(ForwarderActor { target: b_addr }).unwrap();
|
|
||||||
|
|
||||||
// Send message to start of chain
|
|
||||||
rt.send_to(a_addr, ForwardMessage(42)).unwrap();
|
|
||||||
|
|
||||||
// Tick until message arrives
|
|
||||||
for _ in 0..20 {
|
|
||||||
rt.tick();
|
|
||||||
if let Some(ForwardMessage(val)) = inbox.try_recv() {
|
|
||||||
assert_eq!(val, 42);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
panic!("Message did not traverse the chain");
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_multithreaded_message_passing() {
|
|
||||||
let rt = Arc::new(Runtime::new(
|
|
||||||
1000,
|
|
||||||
RuntimeFlavor::Multithreaded { workers: 4 },
|
|
||||||
));
|
|
||||||
let inbox: Inbox<ForwardMessage> = rt.new_inbox();
|
|
||||||
|
|
||||||
// Create a longer chain to exercise multi-threading
|
|
||||||
let mut target = *inbox.addr();
|
|
||||||
for _ in 0..20 {
|
|
||||||
target = rt.spawn(ForwarderActor { target }).unwrap();
|
|
||||||
}
|
|
||||||
|
|
||||||
let start_addr = target;
|
|
||||||
|
|
||||||
// Send message
|
|
||||||
rt.send_to(start_addr, ForwardMessage(999)).unwrap();
|
|
||||||
|
|
||||||
// Spawn thread to check for result and shutdown
|
|
||||||
let rt_clone = Arc::clone(&rt);
|
|
||||||
let inbox_check = thread::spawn(move || {
|
|
||||||
for _ in 0..100 {
|
|
||||||
thread::sleep(Duration::from_millis(10));
|
|
||||||
if let Some(ForwardMessage(val)) = inbox.try_recv() {
|
|
||||||
rt_clone.shutdown();
|
|
||||||
return Some(val);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
rt_clone.shutdown();
|
|
||||||
None
|
|
||||||
});
|
|
||||||
|
|
||||||
rt.run();
|
|
||||||
|
|
||||||
let result = inbox_check.join().unwrap();
|
|
||||||
assert_eq!(result, Some(999));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_is_running_flag() {
|
|
||||||
let rt = Runtime::new(100, RuntimeFlavor::SingleThreaded);
|
|
||||||
|
|
||||||
// Before run(), is_running should be false
|
|
||||||
assert!(!rt.is_running());
|
|
||||||
|
|
||||||
// After shutdown before run, still false
|
|
||||||
rt.shutdown();
|
|
||||||
assert!(!rt.is_running());
|
|
||||||
}
|
|
||||||
Loading…
Reference in a new issue