Skip to main content

behavior/supervision/adapter/
proxy.rs

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