use crate::*;
use futures::future::BoxFuture;
use std::{any::TypeId, future::Future, marker::PhantomData, sync::Arc};
use type_sets::SubsetOf;
pub trait PolyBox: DynPolyBox + Clone {
type Set: Members;
fn into_dyn_unchecked<T>(self) -> DynInbox<T>;
}
pub trait PolyboxExt: PolyBox {
fn into_dyn_subset<T>(self) -> DynInbox<T>
where
T: SubsetOf<Self::Set>,
{
self.into_dyn_unchecked()
}
fn into_dyn(self) -> DynInbox<Self::Set> {
self.into_dyn_unchecked()
}
fn into_dyn_checked<T: Members>(self) -> Result<DynInbox<T>, Self> {
if self.accepts_msgs(&T::members()) {
Ok(self.into_dyn_unchecked())
} else {
Err(self)
}
}
#[must_use]
fn accepts_msg(&self, id: TypeId) -> bool {
<Self::Set as Members>::members().contains(&id)
}
#[must_use]
fn accepts_msgs(&self, ids: &[TypeId]) -> bool {
ids.iter()
.all(|id| <Self::Set as Members>::members().contains(id))
}
fn send_checked<T: Message>(
&self,
msg: T,
) -> impl Future<Output = Result<Output<T>, SendCheckedError<T>>> + Send {
async {
let (payload, output) = T::build_payload(msg);
let payload = BoxedPayload::new::<T>(payload);
match self._send_boxed_payload_checked(payload).await {
Ok(()) => Ok(output),
Err(SendCheckedError::Closed(payload)) => {
let payload = payload
.downcast::<T>()
.expect("Failed to convert payload back");
Err(SendCheckedError::Closed(T::destroy_payload(payload)))
}
Err(SendCheckedError::NotAccepted(payload)) => {
Err(SendCheckedError::NotAccepted(T::destroy_payload(
payload
.downcast::<T>()
.expect("Failed to convert payload back"),
)))
}
}
}
}
fn send_checked_blocking<T: Message>(&self, msg: T) -> Result<Output<T>, SendCheckedError<T>> {
let (payload, output) = T::build_payload(msg);
let payload = BoxedPayload::new::<T>(payload);
match self._send_boxed_payload_checked_blocking(payload) {
Ok(()) => Ok(output),
Err(SendCheckedError::Closed(payload)) => {
let payload = payload
.downcast::<T>()
.expect("Failed to convert payload back");
Err(SendCheckedError::Closed(T::destroy_payload(payload)))
}
Err(SendCheckedError::NotAccepted(payload)) => {
Err(SendCheckedError::NotAccepted(T::destroy_payload(
payload
.downcast::<T>()
.expect("Failed to convert payload back"),
)))
}
}
}
}
impl<T: PolyBox> PolyboxExt for T {}
pub trait DynPolyBox: Send + Sync {
fn _send_boxed_payload_checked(
&self,
msg: BoxedPayload,
) -> BoxFuture<'_, Result<(), SendCheckedError<BoxedPayload>>>;
fn _send_boxed_payload_checked_blocking(
&self,
msg: BoxedPayload,
) -> Result<(), SendCheckedError<BoxedPayload>> {
futures::executor::block_on(self._send_boxed_payload_checked(msg))
}
}
pub struct DynInbox<T> {
inbox: Arc<dyn DynPolyBox>,
_t: PhantomData<fn() -> T>,
}
impl<T> Clone for DynInbox<T> {
fn clone(&self) -> Self {
Self {
inbox: self.inbox.clone(),
_t: PhantomData,
}
}
}
impl<T> DynInbox<T> {
pub fn new_unchecked(inbox: Arc<dyn DynPolyBox>) -> Self {
Self {
inbox,
_t: PhantomData,
}
}
pub fn new<R>(inbox: R) -> Self
where
R: DynPolyBox + PolyBox + 'static,
T: SubsetOf<R::Set>,
{
Self {
inbox: Arc::new(inbox),
_t: PhantomData,
}
}
}
impl<T: Members> PolyBox for DynInbox<T> {
type Set = T;
fn into_dyn_unchecked<R>(self) -> DynInbox<R> {
DynInbox::new_unchecked(self.inbox)
}
}
impl<T, R> Sends<T> for DynInbox<R>
where
T: Message<Kind: MessageSpecifier<T, Output: Send, Payload: Send>>,
R: Members + Contains<T>,
{
async fn send(&self, msg: T) -> Result<Output<T>, SendError<T>> {
self.send_checked(msg).await.map_err(|e| match e {
SendCheckedError::Closed(msg) => SendError(msg),
SendCheckedError::NotAccepted(_msg) => {
panic!(
"Payload was not accepted, this should not happen if the type system is used correctly"
)
}
})
}
}
impl<T> DynPolyBox for DynInbox<T> {
fn _send_boxed_payload_checked(
&self,
msg: BoxedPayload,
) -> BoxFuture<'_, Result<(), SendCheckedError<BoxedPayload>>> {
self.inbox._send_boxed_payload_checked(msg)
}
}