Skip to main content

behavior/supervision/adapter/
proxy.rs

1//! Stable proxy lifecycle and fresh worker incarnation replacement.
2
3use super::super::domain::{
4    Incarnation, IncarnationEffects, IncarnationError, IncarnationPhase, IncarnationStopEffects,
5};
6use super::super::protocol::{ProxyCommand, ProxyEvent};
7use crate::behavior::{
8    Actions, Address, Behavior, Births, Create, Delivery, Recipient, SendAlgebra, ServiceSends,
9    User,
10};
11use crate::next::{Never, Step};
12use crate::protocol::{
13    ObserveChild, ObserveCreation, ReportWorkerCreationResolved, ReportWorkerStopped,
14};
15use crate::{Own, SendInput};
16
17/// The concrete, statically dispatched effect lanes emitted by a [`Proxy`].
18pub struct ProxySends<C: Behavior> {
19    pub deliveries: Vec<Delivery<C>>,
20    pub child_observations: ServiceSends<ObserveChild<<C::Addr as Address>::Nonce>>,
21    pub creation_observations: ServiceSends<ObserveCreation<<C::Addr as Address>::Nonce>>,
22    pub stopped_reports: ServiceSends<ReportWorkerStopped<C::Addr>>,
23    pub creation_reports: ServiceSends<ReportWorkerCreationResolved<<C::Addr as Address>::Nonce>>,
24}
25
26pub(crate) type ProxyActions<C> = Actions<<C as Behavior>::Addr, Never, ProxySends<C>, Births<C>>;
27
28/// A failure of the proxy's own worker-incarnation lifecycle.
29#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
30pub enum ProxyError {
31    #[error("the proxy worker lifecycle was already initialized")]
32    AlreadyInitialized,
33    #[error("the proxy worker creation-attempt sequence is exhausted")]
34    AttemptSequenceExhausted,
35}
36
37impl From<IncarnationError> for ProxyError {
38    fn from(error: IncarnationError) -> Self {
39        match error {
40            IncarnationError::AlreadyInitialized => Self::AlreadyInitialized,
41            IncarnationError::AttemptSequenceExhausted => Self::AttemptSequenceExhausted,
42        }
43    }
44}
45
46impl<C: Behavior> SendAlgebra for ProxySends<C> {
47    fn empty() -> Self {
48        Self {
49            deliveries: Vec::new(),
50            child_observations: ServiceSends::empty(),
51            creation_observations: ServiceSends::empty(),
52            stopped_reports: ServiceSends::empty(),
53            creation_reports: ServiceSends::empty(),
54        }
55    }
56
57    fn append(&mut self, mut other: Self) {
58        self.deliveries.append(&mut other.deliveries);
59        self.child_observations.append(other.child_observations);
60        self.creation_observations
61            .append(other.creation_observations);
62        self.stopped_reports.append(other.stopped_reports);
63        self.creation_reports.append(other.creation_reports);
64    }
65}
66
67impl<C: Behavior> SendInput<Delivery<C>, Own> for ProxySends<C> {
68    fn emit(&mut self, input: Delivery<C>) {
69        self.deliveries.push(input);
70    }
71}
72
73impl<C: Behavior> SendInput<ObserveChild<<C::Addr as Address>::Nonce>, Own> for ProxySends<C> {
74    fn emit(&mut self, input: ObserveChild<<C::Addr as Address>::Nonce>) {
75        self.child_observations.send(input);
76    }
77}
78
79impl<C: Behavior> SendInput<ObserveCreation<<C::Addr as Address>::Nonce>, Own> for ProxySends<C> {
80    fn emit(&mut self, input: ObserveCreation<<C::Addr as Address>::Nonce>) {
81        self.creation_observations.send(input);
82    }
83}
84
85impl<C: Behavior> SendInput<ReportWorkerStopped<C::Addr>, Own> for ProxySends<C> {
86    fn emit(&mut self, input: ReportWorkerStopped<C::Addr>) {
87        self.stopped_reports.send(input);
88    }
89}
90
91impl<C: Behavior> SendInput<ReportWorkerCreationResolved<<C::Addr as Address>::Nonce>, Own>
92    for ProxySends<C>
93{
94    fn emit(&mut self, input: ReportWorkerCreationResolved<<C::Addr as Address>::Nonce>) {
95        self.creation_reports.send(input);
96    }
97}
98
99/// A stable actor that serializes fresh worker-incarnation installation.
100///
101/// A worker is routable only in `Running`. Deadline most one creation can be
102/// `Installing`; stale or provenance-mismatched results are inert. Rejection
103/// leaves `last_installed` unchanged, so a later attempt still names the last
104/// incarnation that actually existed.
105pub struct Proxy<C: Behavior<Ph = Never>> {
106    incarnation: Incarnation<<C::Addr as Address>::Nonce, C>,
107}
108
109impl<C: Behavior<Ph = Never>> Proxy<C> {
110    #[must_use]
111    pub fn new(worker: C) -> Self {
112        Self {
113            incarnation: Incarnation::new(worker),
114        }
115    }
116
117    #[must_use]
118    pub const fn phase(&self) -> IncarnationPhase<<C::Addr as Address>::Nonce> {
119        self.incarnation.phase()
120    }
121}
122
123impl<C> Proxy<C>
124where
125    C: Behavior<Ph = Never>,
126    C::Addr: Address,
127    <C::Addr as Address>::Nonce: From<u64>,
128{
129    fn actions(
130        effects: IncarnationEffects<<C::Addr as Address>::Nonce, C, C::Msg>,
131    ) -> ProxyActions<C> {
132        let mut sends = ProxySends::empty();
133        if let Some((incarnation, message)) = effects.delivery {
134            sends
135                .deliveries
136                .push(Delivery::new(Recipient::child(incarnation), message));
137        }
138        if let Some(resolved) = effects.creation_report {
139            sends.creation_reports.extend([resolved.into()]);
140        }
141        let creates = effects.creation.map_or_else(Vec::new, |creation| {
142            sends
143                .child_observations
144                .extend([ObserveChild::new(creation.attempt)]);
145            sends
146                .creation_observations
147                .extend([ObserveCreation::new(creation.attempt)]);
148            vec![Create::new(creation.attempt, creation.child, creation.kind)]
149        });
150        Actions::new(sends, creates, Step::Continue)
151    }
152
153    fn stopped_actions(
154        effects: IncarnationStopEffects<<C::Addr as Address>::Nonce, C>,
155        event: crate::ChildStopped<C::Addr>,
156    ) -> ProxyActions<C> {
157        let mut actions = Self::actions(IncarnationEffects::new(effects.creation, None, None));
158        if let Some(incarnation) = effects.stopped {
159            actions
160                .sends
161                .stopped_reports
162                .extend([ReportWorkerStopped::from(crate::ChildStopped::new(
163                    incarnation,
164                    event.outcome,
165                    event.at,
166                ))]);
167        }
168        actions
169    }
170}
171
172impl<C> Behavior for Proxy<C>
173where
174    C: Behavior<Ph = Never>,
175    <C::Addr as Address>::Nonce: From<u64>,
176{
177    type Addr = C::Addr;
178    type Msg = ProxyCommand<C>;
179    type Event = ProxyEvent<User<C::Addr, ProxyCommand<C>>>;
180    type Sends = ProxySends<C>;
181    type Ph = Never;
182    type Error = ProxyError;
183    type Birth = Births<C>;
184
185    fn init(
186        &mut self,
187        _: crate::InitializationTurn,
188    ) -> Result<Actions<C::Addr, Never, Self::Sends, Births<C>>, ProxyError> {
189        let effects = self.incarnation.initialize()?;
190        Ok(Self::actions(effects))
191    }
192
193    fn transition(
194        &mut self,
195        _: crate::ActiveTurn,
196        event: Self::Event,
197    ) -> Result<Actions<C::Addr, Never, Self::Sends, Births<C>>, ProxyError> {
198        Ok(match event {
199            ProxyEvent::CreationResolved(resolved) => Self::actions(
200                self.incarnation
201                    .creation_resolved(resolved.nonce, resolved.kind, resolved.result),
202            ),
203            ProxyEvent::ChildStopped(event) => {
204                let effects = self.incarnation.child_stopped(event.nonce)?;
205                Self::stopped_actions(effects, event)
206            }
207            ProxyEvent::Command(event) => match event.message {
208                ProxyCommand::Forward(message) => {
209                    let effects = self.incarnation.forward(message);
210                    Self::actions(effects)
211                }
212                ProxyCommand::Replace(child) => {
213                    let effects = self.incarnation.replace(child)?;
214                    Self::actions(effects)
215                }
216            },
217        })
218    }
219}