Skip to main content

behavior_actors/routing/
correlator.rs

1//! Keyed request/result lifecycle correlation.
2
3use 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/// Result delivered to the reply recipient retained by [`Correlator`].
11#[derive(Debug, Clone, PartialEq, Eq)]
12pub enum CorrelationResult<K, V> {
13    /// A matching pending key accepted one value.
14    Resolved { key: K, value: V },
15    /// A matching pending key was explicitly cancelled.
16    Cancelled { key: K },
17}
18
19/// Complete retained lifecycle of one correlation key.
20pub enum CorrelationState<K, Route> {
21    /// The key owns one reply recipient and may accept resolution or cancel.
22    Pending {
23        /// Correlation key.
24        key: K,
25        /// Typed recipient for the terminal result.
26        reply_to: Route,
27    },
28    /// The key resolved and later replies are stale.
29    Completed { key: K },
30    /// The key was cancelled and later replies are stale.
31    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
42/// Commands accepted by [`Correlator`].
43pub enum CorrelatorMessage<K, V, Route> {
44    /// Establish a unique pending key and its terminal reply recipient.
45    Begin { key: K, reply_to: Route },
46    /// Resolve one pending key with an owned value.
47    Resolve { key: K, value: V },
48    /// Cancel one pending key.
49    Cancel { key: K },
50}
51
52/// A rejected correlation transition.
53#[derive(Debug, Error, Clone, PartialEq, Eq)]
54pub enum CorrelatorError<K, V, Route> {
55    /// The key is already pending.
56    #[error("correlation key is already pending")]
57    AlreadyPending { key: K, reply_to: Route },
58    /// A completed key cannot be opened again without an explicit retention policy.
59    #[error("correlation key is already completed")]
60    ReopenCompleted { key: K, reply_to: Route },
61    /// A cancelled key cannot be opened again without an explicit retention policy.
62    #[error("correlation key is already cancelled")]
63    ReopenCancelled { key: K, reply_to: Route },
64    /// An operation targeted an already completed lifecycle.
65    #[error("correlation key is already completed")]
66    AlreadyCompleted(K),
67    /// An operation targeted an already cancelled lifecycle.
68    #[error("correlation key is already cancelled")]
69    AlreadyCancelled(K),
70    /// No lifecycle fact exists for the supplied key.
71    #[error("correlation key is unknown")]
72    Unknown(K),
73    /// No lifecycle fact exists for a reply; the owned value is returned.
74    #[error("correlation reply key is unknown")]
75    UnknownReply { key: K, value: V },
76    /// A reply arrived after successful completion; the owned value is returned.
77    #[error("correlation reply is stale because the key completed")]
78    StaleCompleted { key: K, value: V },
79    /// A reply arrived after cancellation; the owned value is returned.
80    #[error("correlation reply is stale because the key was cancelled")]
81    StaleCancelled { key: K, value: V },
82}
83
84/// Deterministic keyed request/result lifecycle behavior.
85///
86/// Each key is exactly pending, completed, or cancelled. `Begin` owns the
87/// reply recipient. `Resolve` atomically marks a pending key completed and
88/// emits one typed result; `Cancel` marks it cancelled and emits one typed
89/// cancellation. Terminal keys are retained, making duplicate and stale input
90/// explicit rather than indistinguishable from never-seen input. Reopening and
91/// retention expiry require a separate explicit policy and are not inferred
92/// from time or key reuse. Initialization is empty and the behavior does not
93/// terminate itself. These lifecycle and retention choices are Bombay policy;
94/// delivery remains an Address/Communication capability. No method has a
95/// semantic panic condition.
96pub 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    /// Construct an empty correlator definition.
112    #[must_use]
113    pub const fn new() -> Self {
114        Self {
115            states: Vec::new(),
116            address: core::marker::PhantomData,
117        }
118    }
119
120    /// Borrow every retained lifecycle in first-begin order.
121    #[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}