use futures::channel::oneshot;
use std::borrow::Cow;
use std::future::Future;
pub use puppet_derive::puppet_actor;
#[doc(hidden)]
pub mod __private {
pub use flume;
pub use futures;
#[cfg(feature = "helper-methods")]
pub use tokio;
}
pub trait Actor {
type Messages;
}
pub trait Message {
type Output;
}
pub trait Executor {
fn spawn(&self, fut: impl Future<Output = ()> + Send + 'static);
}
#[doc(hidden)]
pub trait MessageHandler<T: Message> {
fn create(msg: T) -> (Self, oneshot::Receiver<T::Output>)
where
Self: Sized;
}
pub struct ActorMailbox<A: Actor> {
tx: flume::Sender<A::Messages>,
name: Cow<'static, str>,
}
impl<A: Actor> Clone for ActorMailbox<A> {
fn clone(&self) -> Self {
Self {
tx: self.tx.clone(),
name: self.name.clone(),
}
}
}
impl<A: Actor> ActorMailbox<A> {
#[doc(hidden)]
pub fn new(tx: flume::Sender<A::Messages>, name: Cow<'static, str>) -> Self {
Self { tx, name }
}
#[inline]
pub fn name(&self) -> &str {
self.name.as_ref()
}
pub async fn send<T>(&self, msg: T) -> T::Output
where
T: Message,
A::Messages: MessageHandler<T>,
{
let (msg, rx) = A::Messages::create(msg);
self.tx.send_async(msg).await.expect("Contact actor");
rx.await.expect("Actor response")
}
pub async fn deferred_send<T>(&self, msg: T) -> DeferredResponse<T::Output>
where
T: Message,
A::Messages: MessageHandler<T>,
{
let (msg, rx) = A::Messages::create(msg);
self.tx.send_async(msg).await.expect("Contact actor");
DeferredResponse { rx }
}
pub fn send_sync<T>(&self, msg: T) -> T::Output
where
T: Message,
A::Messages: MessageHandler<T>,
{
let (msg, rx) = A::Messages::create(msg);
self.tx.send(msg).expect("Contact actor");
futures::executor::block_on(rx).expect("Actor response")
}
}
pub struct DeferredResponse<T> {
rx: oneshot::Receiver<T>,
}
impl<T> DeferredResponse<T> {
pub fn try_recv(&mut self) -> Option<T> {
self.rx.try_recv().expect("Get actor response")
}
pub async fn recv(self) -> T {
self.rx.await.expect("Get actor response")
}
}
pub struct Reply<T: Default> {
tx: Option<oneshot::Sender<T>>,
}
impl<T: Default> From<oneshot::Sender<T>> for Reply<T> {
fn from(tx: oneshot::Sender<T>) -> Self {
Self { tx: Some(tx) }
}
}
impl<T: Default> Reply<T> {
pub fn reply(mut self, msg: T) {
if let Some(tx) = self.tx.take() {
let _ = tx.send(msg);
}
}
}
impl<T: Default> Drop for Reply<T> {
fn drop(&mut self) {
if let Some(tx) = self.tx.take() {
let _ = tx.send(T::default());
}
}
}
#[macro_export]
macro_rules! derive_message {
($msg:ident, $output:ty) => {
impl $crate::Message for $msg {
type Output = $output;
}
};
}