use super::super::domain::{
Incarnation, IncarnationEffects, IncarnationError, IncarnationPhase, IncarnationStopEffects,
};
use super::super::protocol::{ProxyCommand, ProxyEvent};
use crate::behavior::{
Actions, Address, Behavior, Births, Create, Delivery, Recipient, SendAlgebra, ServiceSends,
User,
};
use crate::next::{Never, Step};
use crate::protocol::{
ObserveChild, ObserveCreation, ReportWorkerCreationResolved, ReportWorkerStopped,
};
use crate::{Own, SendInput};
pub struct ProxySends<C: Behavior> {
pub deliveries: Vec<Delivery<C>>,
pub child_observations: ServiceSends<ObserveChild<<C::Addr as Address>::Nonce>>,
pub creation_observations: ServiceSends<ObserveCreation<<C::Addr as Address>::Nonce>>,
pub stopped_reports: ServiceSends<ReportWorkerStopped<C::Addr>>,
pub creation_reports: ServiceSends<ReportWorkerCreationResolved<<C::Addr as Address>::Nonce>>,
}
pub(crate) type ProxyActions<C> = Actions<<C as Behavior>::Addr, Never, ProxySends<C>, Births<C>>;
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
pub enum ProxyError {
#[error("the proxy worker lifecycle was already initialized")]
AlreadyInitialized,
#[error("the proxy worker creation-attempt sequence is exhausted")]
AttemptSequenceExhausted,
}
impl From<IncarnationError> for ProxyError {
fn from(error: IncarnationError) -> Self {
match error {
IncarnationError::AlreadyInitialized => Self::AlreadyInitialized,
IncarnationError::AttemptSequenceExhausted => Self::AttemptSequenceExhausted,
}
}
}
impl<C: Behavior> SendAlgebra for ProxySends<C> {
fn empty() -> Self {
Self {
deliveries: Vec::new(),
child_observations: ServiceSends::empty(),
creation_observations: ServiceSends::empty(),
stopped_reports: ServiceSends::empty(),
creation_reports: ServiceSends::empty(),
}
}
fn append(&mut self, mut other: Self) {
self.deliveries.append(&mut other.deliveries);
self.child_observations.append(other.child_observations);
self.creation_observations
.append(other.creation_observations);
self.stopped_reports.append(other.stopped_reports);
self.creation_reports.append(other.creation_reports);
}
}
impl<C: Behavior> SendInput<Delivery<C>, Own> for ProxySends<C> {
fn emit(&mut self, input: Delivery<C>) {
self.deliveries.push(input);
}
}
impl<C: Behavior> SendInput<ObserveChild<<C::Addr as Address>::Nonce>, Own> for ProxySends<C> {
fn emit(&mut self, input: ObserveChild<<C::Addr as Address>::Nonce>) {
self.child_observations.send(input);
}
}
impl<C: Behavior> SendInput<ObserveCreation<<C::Addr as Address>::Nonce>, Own> for ProxySends<C> {
fn emit(&mut self, input: ObserveCreation<<C::Addr as Address>::Nonce>) {
self.creation_observations.send(input);
}
}
impl<C: Behavior> SendInput<ReportWorkerStopped<C::Addr>, Own> for ProxySends<C> {
fn emit(&mut self, input: ReportWorkerStopped<C::Addr>) {
self.stopped_reports.send(input);
}
}
impl<C: Behavior> SendInput<ReportWorkerCreationResolved<<C::Addr as Address>::Nonce>, Own>
for ProxySends<C>
{
fn emit(&mut self, input: ReportWorkerCreationResolved<<C::Addr as Address>::Nonce>) {
self.creation_reports.send(input);
}
}
pub struct Proxy<C: Behavior<Ph = Never>> {
incarnation: Incarnation<<C::Addr as Address>::Nonce, C>,
}
impl<C: Behavior<Ph = Never>> Proxy<C> {
#[must_use]
pub fn new(worker: C) -> Self {
Self {
incarnation: Incarnation::new(worker),
}
}
#[must_use]
pub const fn phase(&self) -> IncarnationPhase<<C::Addr as Address>::Nonce> {
self.incarnation.phase()
}
}
impl<C> Proxy<C>
where
C: Behavior<Ph = Never>,
C::Addr: Address,
<C::Addr as Address>::Nonce: From<u64>,
{
fn actions(
effects: IncarnationEffects<<C::Addr as Address>::Nonce, C, C::Msg>,
) -> ProxyActions<C> {
let mut sends = ProxySends::empty();
if let Some((incarnation, message)) = effects.delivery {
sends
.deliveries
.push(Delivery::new(Recipient::child(incarnation), message));
}
if let Some(resolved) = effects.creation_report {
sends.creation_reports.extend([resolved.into()]);
}
let creates = effects.creation.map_or_else(Vec::new, |creation| {
sends
.child_observations
.extend([ObserveChild::new(creation.attempt)]);
sends
.creation_observations
.extend([ObserveCreation::new(creation.attempt)]);
vec![Create::new(creation.attempt, creation.child, creation.kind)]
});
Actions::new(sends, creates, Step::Continue)
}
fn stopped_actions(
effects: IncarnationStopEffects<<C::Addr as Address>::Nonce, C>,
event: crate::ChildStopped<C::Addr>,
) -> ProxyActions<C> {
let mut actions = Self::actions(IncarnationEffects::new(effects.creation, None, None));
if let Some(incarnation) = effects.stopped {
actions
.sends
.stopped_reports
.extend([ReportWorkerStopped::from(crate::ChildStopped::new(
incarnation,
event.outcome,
event.at,
))]);
}
actions
}
}
impl<C> Behavior for Proxy<C>
where
C: Behavior<Ph = Never>,
<C::Addr as Address>::Nonce: From<u64>,
{
type Addr = C::Addr;
type Msg = ProxyCommand<C>;
type Event = ProxyEvent<User<C::Addr, ProxyCommand<C>>>;
type Sends = ProxySends<C>;
type Ph = Never;
type Error = ProxyError;
type Birth = Births<C>;
fn init(
&mut self,
_: crate::InitializationTurn,
) -> Result<Actions<C::Addr, Never, Self::Sends, Births<C>>, ProxyError> {
let effects = self.incarnation.initialize()?;
Ok(Self::actions(effects))
}
fn transition(
&mut self,
_: crate::ActiveTurn,
event: Self::Event,
) -> Result<Actions<C::Addr, Never, Self::Sends, Births<C>>, ProxyError> {
Ok(match event {
ProxyEvent::CreationResolved(resolved) => Self::actions(
self.incarnation
.creation_resolved(resolved.nonce, resolved.kind, resolved.result),
),
ProxyEvent::ChildStopped(event) => {
let effects = self.incarnation.child_stopped(event.nonce)?;
Self::stopped_actions(effects, event)
}
ProxyEvent::Command(event) => match event.message {
ProxyCommand::Forward(message) => {
let effects = self.incarnation.forward(message);
Self::actions(effects)
}
ProxyCommand::Replace(child) => {
let effects = self.incarnation.replace(child)?;
Self::actions(effects)
}
},
})
}
}