behavior/supervision/adapter/
proxy.rs1use 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
15pub 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
79pub 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}