behavior/supervision/adapter/
proxy.rs1use 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
17pub 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#[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
99pub 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}