1use behavior::{Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, User};
4#[cfg(test)]
5use behavior::{Delivery, Recipient};
6use thiserror::Error;
7
8use crate::DeliveryRoute;
9
10#[derive(Debug, Clone, PartialEq, Eq)]
12pub enum CorrelationResult<K, V> {
13 Resolved { key: K, value: V },
15 Cancelled { key: K },
17}
18
19pub enum CorrelationState<K, Route> {
21 Pending {
23 key: K,
25 reply_to: Route,
27 },
28 Completed { key: K },
30 Cancelled { key: K },
32}
33
34impl<K, Route> CorrelationState<K, Route> {
35 fn key(&self) -> &K {
36 match self {
37 Self::Pending { key, .. } | Self::Completed { key } | Self::Cancelled { key } => key,
38 }
39 }
40}
41
42pub enum CorrelatorMessage<K, V, Route> {
44 Begin { key: K, reply_to: Route },
46 Resolve { key: K, value: V },
48 Cancel { key: K },
50}
51
52#[derive(Debug, Error, Clone, PartialEq, Eq)]
54pub enum CorrelatorError<K, V, Route> {
55 #[error("correlation key is already pending")]
57 AlreadyPending { key: K, reply_to: Route },
58 #[error("correlation key is already completed")]
60 ReopenCompleted { key: K, reply_to: Route },
61 #[error("correlation key is already cancelled")]
63 ReopenCancelled { key: K, reply_to: Route },
64 #[error("correlation key is already completed")]
66 AlreadyCompleted(K),
67 #[error("correlation key is already cancelled")]
69 AlreadyCancelled(K),
70 #[error("correlation key is unknown")]
72 Unknown(K),
73 #[error("correlation reply key is unknown")]
75 UnknownReply { key: K, value: V },
76 #[error("correlation reply is stale because the key completed")]
78 StaleCompleted { key: K, value: V },
79 #[error("correlation reply is stale because the key was cancelled")]
81 StaleCancelled { key: K, value: V },
82}
83
84pub struct Correlator<
97 A: Address,
98 K,
99 V,
100 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = CorrelationResult<K, V>>>,
101> {
102 states: Vec<CorrelationState<K, Route>>,
103 address: core::marker::PhantomData<fn() -> A>,
104}
105
106impl<A, K, V, Route> Correlator<A, K, V, Route>
107where
108 A: Address,
109 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = CorrelationResult<K, V>>>,
110{
111 #[must_use]
113 pub const fn new() -> Self {
114 Self {
115 states: Vec::new(),
116 address: core::marker::PhantomData,
117 }
118 }
119
120 #[must_use]
122 pub fn states(&self) -> &[CorrelationState<K, Route>] {
123 &self.states
124 }
125}
126
127impl<A, K, V, Route> Default for Correlator<A, K, V, Route>
128where
129 A: Address,
130 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = CorrelationResult<K, V>>>,
131{
132 fn default() -> Self {
133 Self::new()
134 }
135}
136
137impl<A, K, V, Route> BehaviorBase for Correlator<A, K, V, Route>
138where
139 A: Address,
140 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = CorrelationResult<K, V>>>,
141{
142 type Base = Self;
143
144 fn base(&self) -> &Self {
145 self
146 }
147}
148
149impl<A, K, V, Route> behavior::Protocol for Correlator<A, K, V, Route>
150where
151 A: Address,
152 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = CorrelationResult<K, V>>>,
153{
154 type Addr = A;
155 type Msg = CorrelatorMessage<K, V, Route>;
156}
157
158impl<A, K, V, Route> Behavior for Correlator<A, K, V, Route>
159where
160 A: Address,
161 K: Clone + Eq,
162 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = CorrelationResult<K, V>>>
163 + Clone,
164 Route::Sends: behavior::SendsFor<User<A, CorrelatorMessage<K, V, Route>>>,
165{
166 type Protocol = Self;
167 type Event = User<A, behavior::BehaviorMessage<Self>>;
168 type Sends = Route::Sends;
169 type Ph = Never;
170 type Error = CorrelatorError<K, V, Route>;
171 type Birth = NoBirths;
172
173 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
174 match event.message {
175 CorrelatorMessage::Begin { key, reply_to } => {
176 if let Some(existing) = self.states.iter().find(|state| state.key() == &key) {
177 return Err(match existing {
178 CorrelationState::Pending { .. } => {
179 CorrelatorError::AlreadyPending { key, reply_to }
180 }
181 CorrelationState::Completed { .. } => {
182 CorrelatorError::ReopenCompleted { key, reply_to }
183 }
184 CorrelationState::Cancelled { .. } => {
185 CorrelatorError::ReopenCancelled { key, reply_to }
186 }
187 });
188 }
189 self.states
190 .push(CorrelationState::Pending { key, reply_to });
191 Ok(Actions::cont())
192 }
193 CorrelatorMessage::Resolve { key, value } => {
194 let Some(index) = self.states.iter().position(|state| state.key() == &key) else {
195 return Err(CorrelatorError::UnknownReply { key, value });
196 };
197 let reply_to = match &self.states[index] {
198 CorrelationState::Pending { reply_to, .. } => reply_to.clone(),
199 CorrelationState::Completed { .. } => {
200 return Err(CorrelatorError::StaleCompleted { key, value });
201 }
202 CorrelationState::Cancelled { .. } => {
203 return Err(CorrelatorError::StaleCancelled { key, value });
204 }
205 };
206 self.states[index] = CorrelationState::Completed { key: key.clone() };
207 Ok(Actions::send(
208 reply_to.deliver(CorrelationResult::Resolved { key, value }),
209 ))
210 }
211 CorrelatorMessage::Cancel { key } => {
212 let Some(index) = self.states.iter().position(|state| state.key() == &key) else {
213 return Err(CorrelatorError::Unknown(key));
214 };
215 let reply_to = match &self.states[index] {
216 CorrelationState::Pending { reply_to, .. } => reply_to.clone(),
217 CorrelationState::Completed { .. } => {
218 return Err(CorrelatorError::AlreadyCompleted(key));
219 }
220 CorrelationState::Cancelled { .. } => {
221 return Err(CorrelatorError::AlreadyCancelled(key));
222 }
223 };
224 self.states[index] = CorrelationState::Cancelled { key: key.clone() };
225 Ok(Actions::send(
226 reply_to.deliver(CorrelationResult::Cancelled { key }),
227 ))
228 }
229 }
230 }
231}
232
233#[cfg(test)]
234mod tests {
235 use super::*;
236 use crate::Activate as _;
237 use behavior::MailAddr;
238
239 struct Reply;
240
241 impl behavior::Protocol for Reply {
242 type Addr = MailAddr;
243 type Msg = CorrelationResult<u8, u16>;
244 }
245
246 impl Behavior for Reply {
247 type Protocol = Self;
248 type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
249 type Sends = Vec<Never>;
250 type Ph = Never;
251 type Error = Never;
252 type Birth = NoBirths;
253
254 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
255 Ok(Actions::cont())
256 }
257 }
258
259 type TestCorrelator = Correlator<MailAddr, u8, u16, Recipient<Reply>>;
260
261 #[test]
262 fn resolve_commits_before_one_terminal_delivery_and_marks_duplicates_stale() {
263 let reply = Recipient::<Reply>::global(MailAddr(8));
264 let mut correlator = (TestCorrelator::new()).initialize().unwrap().behavior;
265 let begun = correlator
266 .receive(
267 MailAddr(9),
268 CorrelatorMessage::Begin {
269 key: 1,
270 reply_to: reply,
271 },
272 )
273 .unwrap();
274 assert!(begun.sends.is_empty());
275 assert!(begun.creates.is_empty());
276 assert_eq!(begun.become_, behavior::Step::Continue);
277 let resolved = correlator
278 .receive(
279 MailAddr(9),
280 CorrelatorMessage::Resolve { key: 1, value: 42 },
281 )
282 .unwrap();
283 assert!(
284 resolved.sends
285 == vec![Delivery::new(
286 reply,
287 CorrelationResult::Resolved { key: 1, value: 42 },
288 )]
289 );
290 assert!(resolved.creates.is_empty());
291 assert!(matches!(
292 correlator.states(),
293 [CorrelationState::Completed { key: 1 }]
294 ));
295 let rejection = correlator.receive(
296 MailAddr(9),
297 CorrelatorMessage::Resolve { key: 1, value: 43 },
298 );
299 assert!(matches!(
300 rejection,
301 Err(CorrelatorError::StaleCompleted { key: 1, value: 43 })
302 ));
303 }
304
305 #[test]
306 fn cancel_is_terminal_and_unknown_reply_preserves_the_value() {
307 let reply = Recipient::<Reply>::global(MailAddr(8));
308 let mut correlator = (TestCorrelator::new()).initialize().unwrap().behavior;
309 let rejection = correlator.receive(
310 MailAddr(9),
311 CorrelatorMessage::Resolve { key: 7, value: 99 },
312 );
313 assert!(matches!(
314 rejection,
315 Err(CorrelatorError::UnknownReply { key: 7, value: 99 })
316 ));
317 let begun = correlator
318 .receive(
319 MailAddr(9),
320 CorrelatorMessage::Begin {
321 key: 2,
322 reply_to: reply,
323 },
324 )
325 .unwrap();
326 assert!(begun.sends.is_empty());
327 assert!(begun.creates.is_empty());
328 assert_eq!(begun.become_, behavior::Step::Continue);
329 let cancelled = correlator
330 .receive(MailAddr(9), CorrelatorMessage::Cancel { key: 2 })
331 .unwrap();
332 assert!(
333 cancelled.sends
334 == vec![Delivery::new(
335 reply,
336 CorrelationResult::Cancelled { key: 2 },
337 )]
338 );
339 let rejection = correlator.receive(
340 MailAddr(9),
341 CorrelatorMessage::Resolve { key: 2, value: 100 },
342 );
343 assert!(matches!(
344 rejection,
345 Err(CorrelatorError::StaleCancelled { key: 2, value: 100 })
346 ));
347 }
348}