Skip to main content

behavior_actors/routing/
acknowledgements.rs

1//! Multi-participant acknowledgement lifecycle correlation.
2
3#[cfg(test)]
4use behavior::Recipient;
5use behavior::{Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, User};
6use thiserror::Error;
7
8use crate::DeliveryRoute;
9
10/// Exhaustive lifecycle phase for one acknowledgement key.
11#[derive(Debug, Clone, PartialEq, Eq)]
12pub enum AcknowledgementState<P> {
13    /// Some required participants have not acknowledged.
14    Pending {
15        /// Participants still required, in declaration order.
16        remaining: Vec<P>,
17        /// Participants already accepted, in acknowledgement order.
18        acknowledged: Vec<P>,
19    },
20    /// Every declared participant acknowledged.
21    Completed,
22    /// The pending lifecycle was explicitly cancelled.
23    Cancelled,
24}
25
26/// One retained acknowledgement lifecycle.
27#[derive(Debug, Clone, PartialEq, Eq)]
28pub struct AcknowledgementRecord<K, P> {
29    /// Application-defined correlation key.
30    pub key: K,
31    /// Complete current phase.
32    pub state: AcknowledgementState<P>,
33}
34
35/// Exact operation submitted to an already terminal or unknown lifecycle.
36#[derive(Debug, Clone, PartialEq, Eq)]
37pub enum AcknowledgementInput<K, P> {
38    /// One participant attempted to acknowledge.
39    Acknowledge { key: K, participant: P },
40    /// The lifecycle was asked to cancel.
41    Cancel { key: K },
42}
43
44/// Typed rejection of an acknowledgement operation.
45#[derive(Debug, Error, Clone, PartialEq, Eq)]
46pub enum AcknowledgementError<K, P> {
47    /// A lifecycle already exists for this key.
48    #[error("acknowledgement key already exists")]
49    Existing {
50        /// Rejected key.
51        key: K,
52        /// Exact participant declaration from the rejected command.
53        participants: Vec<P>,
54    },
55    /// No lifecycle exists for this key.
56    #[error("acknowledgement key is unknown")]
57    Unknown(AcknowledgementInput<K, P>),
58    /// The participant was not declared for the pending lifecycle.
59    #[error("participant is not required by this acknowledgement")]
60    UnexpectedParticipant {
61        /// Correlation key.
62        key: K,
63        /// Rejected participant.
64        participant: P,
65    },
66    /// The participant already acknowledged.
67    #[error("participant acknowledgement is a duplicate")]
68    DuplicateParticipant {
69        /// Correlation key.
70        key: K,
71        /// Rejected participant.
72        participant: P,
73    },
74    /// An operation targeted a completed lifecycle.
75    #[error("acknowledgement lifecycle is already complete")]
76    Completed(AcknowledgementInput<K, P>),
77    /// An operation targeted a cancelled lifecycle.
78    #[error("acknowledgement lifecycle is cancelled")]
79    Cancelled(AcknowledgementInput<K, P>),
80}
81
82/// Complete result of one acknowledgement operation.
83#[derive(Debug, Clone, PartialEq, Eq)]
84pub enum AcknowledgementOutcome<K, P> {
85    /// A new lifecycle was accepted.
86    Started {
87        /// Correlation key.
88        key: K,
89        /// Number of distinct participants still required.
90        remaining: usize,
91    },
92    /// One participant was accepted and requirements remain.
93    Acknowledged {
94        /// Correlation key.
95        key: K,
96        /// Accepted participant.
97        participant: P,
98        /// Number still required.
99        remaining: usize,
100    },
101    /// The lifecycle reached its terminal successful phase.
102    Completed {
103        /// Correlation key.
104        key: K,
105    },
106    /// The lifecycle reached its terminal cancelled phase.
107    Cancelled {
108        /// Correlation key.
109        key: K,
110    },
111    /// The operation was rejected without changing state.
112    Rejected(AcknowledgementError<K, P>),
113}
114
115/// Operations accepted by [`Acknowledgements`].
116pub enum AcknowledgementMessage<K, P, Route> {
117    /// Establish a fresh acknowledgement lifecycle.
118    Begin {
119        /// Correlation key.
120        key: K,
121        /// Required participants; duplicates are normalized by first occurrence.
122        participants: Vec<P>,
123        /// Typed outcome recipient.
124        reply_to: Route,
125    },
126    /// Record one participant acknowledgement.
127    Acknowledge {
128        /// Correlation key.
129        key: K,
130        /// Participant making the acknowledgement.
131        participant: P,
132        /// Typed outcome recipient.
133        reply_to: Route,
134    },
135    /// Cancel a pending lifecycle.
136    Cancel {
137        /// Correlation key.
138        key: K,
139        /// Typed outcome recipient.
140        reply_to: Route,
141    },
142}
143
144/// Typed multi-participant acknowledgement behavior.
145///
146/// `Begin` normalizes participant membership by first occurrence. An empty set
147/// completes immediately. Each declared participant can advance a pending
148/// lifecycle exactly once; the final acknowledgement atomically commits
149/// [`AcknowledgementState::Completed`] and reports completion. Cancellation is
150/// accepted only while pending. Completed and cancelled records are retained,
151/// making stale terminal input distinguishable from an unknown key. Every
152/// rejection is emitted as a concrete [`AcknowledgementError`] and leaves the
153/// state unchanged. Initialization is empty, the template creates no actors,
154/// and it does not terminate its host. Membership normalization, terminal
155/// retention, and acknowledgement ordering are Bombay policy. Delivery remains
156/// a runtime capability, and transitions have no panic path.
157pub struct Acknowledgements<
158    A: Address,
159    K,
160    P,
161    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
162> {
163    records: Vec<AcknowledgementRecord<K, P>>,
164    marker: core::marker::PhantomData<fn() -> (A, Route)>,
165}
166
167impl<A, K, P, Route> Acknowledgements<A, K, P, Route>
168where
169    A: Address,
170    K: Clone + Eq,
171    P: Clone + Eq,
172    Route:
173        DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
174{
175    /// Construct an empty acknowledgement table.
176    #[must_use]
177    pub const fn new() -> Self {
178        Self {
179            records: Vec::new(),
180            marker: core::marker::PhantomData,
181        }
182    }
183
184    /// Borrow every pending or retained terminal lifecycle.
185    #[must_use]
186    pub fn records(&self) -> &[AcknowledgementRecord<K, P>] {
187        &self.records
188    }
189
190    fn result(
191        reply_to: Route,
192        outcome: AcknowledgementOutcome<K, P>,
193    ) -> Actions<A, Never, Route::Sends, NoBirths> {
194        Actions::send(reply_to.deliver(outcome))
195    }
196
197    fn begin(
198        &mut self,
199        key: K,
200        participants: Vec<P>,
201        reply_to: Route,
202    ) -> Actions<A, Never, Route::Sends, NoBirths> {
203        if self.records.iter().any(|record| record.key == key) {
204            return Self::result(
205                reply_to,
206                AcknowledgementOutcome::Rejected(AcknowledgementError::Existing {
207                    key,
208                    participants,
209                }),
210            );
211        }
212        let mut distinct = Vec::new();
213        for participant in participants {
214            if !distinct.contains(&participant) {
215                distinct.push(participant);
216            }
217        }
218        let remaining = distinct.len();
219        let state = if distinct.is_empty() {
220            AcknowledgementState::Completed
221        } else {
222            AcknowledgementState::Pending {
223                remaining: distinct,
224                acknowledged: Vec::new(),
225            }
226        };
227        self.records.push(AcknowledgementRecord {
228            key: key.clone(),
229            state,
230        });
231        let outcome = if remaining == 0 {
232            AcknowledgementOutcome::Completed { key }
233        } else {
234            AcknowledgementOutcome::Started { key, remaining }
235        };
236        Self::result(reply_to, outcome)
237    }
238
239    fn acknowledge(
240        &mut self,
241        key: K,
242        participant: P,
243        reply_to: Route,
244    ) -> Actions<A, Never, Route::Sends, NoBirths> {
245        let Some(record) = self.records.iter_mut().find(|record| record.key == key) else {
246            return Self::result(
247                reply_to,
248                AcknowledgementOutcome::Rejected(AcknowledgementError::Unknown(
249                    AcknowledgementInput::Acknowledge { key, participant },
250                )),
251            );
252        };
253        let (remaining, acknowledged) = match &mut record.state {
254            AcknowledgementState::Pending {
255                remaining,
256                acknowledged,
257            } => (remaining, acknowledged),
258            AcknowledgementState::Completed => {
259                return Self::result(
260                    reply_to,
261                    AcknowledgementOutcome::Rejected(AcknowledgementError::Completed(
262                        AcknowledgementInput::Acknowledge { key, participant },
263                    )),
264                );
265            }
266            AcknowledgementState::Cancelled => {
267                return Self::result(
268                    reply_to,
269                    AcknowledgementOutcome::Rejected(AcknowledgementError::Cancelled(
270                        AcknowledgementInput::Acknowledge { key, participant },
271                    )),
272                );
273            }
274        };
275        if acknowledged.contains(&participant) {
276            return Self::result(
277                reply_to,
278                AcknowledgementOutcome::Rejected(AcknowledgementError::DuplicateParticipant {
279                    key,
280                    participant,
281                }),
282            );
283        }
284        let Some(index) = remaining
285            .iter()
286            .position(|required| required == &participant)
287        else {
288            return Self::result(
289                reply_to,
290                AcknowledgementOutcome::Rejected(AcknowledgementError::UnexpectedParticipant {
291                    key,
292                    participant,
293                }),
294            );
295        };
296        remaining.remove(index);
297        acknowledged.push(participant.clone());
298        if remaining.is_empty() {
299            record.state = AcknowledgementState::Completed;
300            Self::result(reply_to, AcknowledgementOutcome::Completed { key })
301        } else {
302            Self::result(
303                reply_to,
304                AcknowledgementOutcome::Acknowledged {
305                    key,
306                    participant,
307                    remaining: remaining.len(),
308                },
309            )
310        }
311    }
312
313    fn cancel(&mut self, key: K, reply_to: Route) -> Actions<A, Never, Route::Sends, NoBirths> {
314        let Some(record) = self.records.iter_mut().find(|record| record.key == key) else {
315            return Self::result(
316                reply_to,
317                AcknowledgementOutcome::Rejected(AcknowledgementError::Unknown(
318                    AcknowledgementInput::Cancel { key },
319                )),
320            );
321        };
322        match record.state {
323            AcknowledgementState::Pending { .. } => {
324                record.state = AcknowledgementState::Cancelled;
325                Self::result(reply_to, AcknowledgementOutcome::Cancelled { key })
326            }
327            AcknowledgementState::Completed => Self::result(
328                reply_to,
329                AcknowledgementOutcome::Rejected(AcknowledgementError::Completed(
330                    AcknowledgementInput::Cancel { key },
331                )),
332            ),
333            AcknowledgementState::Cancelled => Self::result(
334                reply_to,
335                AcknowledgementOutcome::Rejected(AcknowledgementError::Cancelled(
336                    AcknowledgementInput::Cancel { key },
337                )),
338            ),
339        }
340    }
341}
342
343impl<A, K, P, Route> Default for Acknowledgements<A, K, P, Route>
344where
345    A: Address,
346    K: Clone + Eq,
347    P: Clone + Eq,
348    Route:
349        DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
350{
351    fn default() -> Self {
352        Self::new()
353    }
354}
355
356impl<A, K, P, Route> BehaviorBase for Acknowledgements<A, K, P, Route>
357where
358    A: Address,
359    Route:
360        DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
361{
362    type Base = Self;
363    fn base(&self) -> &Self {
364        self
365    }
366}
367
368impl<A, K, P, Route> behavior::Protocol for Acknowledgements<A, K, P, Route>
369where
370    A: Address,
371    Route:
372        DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
373{
374    type Addr = A;
375    type Msg = AcknowledgementMessage<K, P, Route>;
376}
377
378impl<A, K, P, Route> Behavior for Acknowledgements<A, K, P, Route>
379where
380    A: Address,
381    K: Clone + Eq,
382    P: Clone + Eq,
383    Route:
384        DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
385    Route::Sends: behavior::SendsFor<User<A, AcknowledgementMessage<K, P, Route>>>,
386{
387    type Protocol = Self;
388    type Event = User<A, behavior::BehaviorMessage<Self>>;
389    type Sends = Route::Sends;
390    type Ph = Never;
391    type Error = Never;
392    type Birth = NoBirths;
393
394    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
395        Ok(match event.message {
396            AcknowledgementMessage::Begin {
397                key,
398                participants,
399                reply_to,
400            } => self.begin(key, participants, reply_to),
401            AcknowledgementMessage::Acknowledge {
402                key,
403                participant,
404                reply_to,
405            } => self.acknowledge(key, participant, reply_to),
406            AcknowledgementMessage::Cancel { key, reply_to } => self.cancel(key, reply_to),
407        })
408    }
409}
410
411#[cfg(test)]
412mod tests {
413    use super::*;
414    use crate::Activate as _;
415    use behavior::MailAddr;
416
417    struct Reply;
418    impl behavior::Protocol for Reply {
419        type Addr = MailAddr;
420        type Msg = AcknowledgementOutcome<u8, u8>;
421    }
422
423    impl Behavior for Reply {
424        type Protocol = Self;
425        type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
426        type Sends = Vec<Never>;
427        type Ph = Never;
428        type Error = Never;
429        type Birth = NoBirths;
430        fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
431            Ok(Actions::cont())
432        }
433    }
434
435    type Subject = Acknowledgements<MailAddr, u8, u8, Recipient<Reply>>;
436    fn reply() -> Recipient<Reply> {
437        Recipient::global(MailAddr(1))
438    }
439
440    #[test]
441    fn membership_is_normalized_and_completion_is_terminal() {
442        let mut subject = (Subject::new()).initialize().unwrap().behavior;
443        let started = subject
444            .receive(
445                MailAddr(9),
446                AcknowledgementMessage::Begin {
447                    key: 7,
448                    participants: vec![1, 1, 2],
449                    reply_to: reply(),
450                },
451            )
452            .unwrap();
453        assert!(matches!(
454            started.sends[0].message,
455            AcknowledgementOutcome::Started {
456                key: 7,
457                remaining: 2
458            }
459        ));
460        let first = subject
461            .receive(
462                MailAddr(9),
463                AcknowledgementMessage::Acknowledge {
464                    key: 7,
465                    participant: 1,
466                    reply_to: reply(),
467                },
468            )
469            .unwrap();
470        assert!(matches!(
471            first.sends[0].message,
472            AcknowledgementOutcome::Acknowledged { remaining: 1, .. }
473        ));
474        let completed = subject
475            .receive(
476                MailAddr(9),
477                AcknowledgementMessage::Acknowledge {
478                    key: 7,
479                    participant: 2,
480                    reply_to: reply(),
481                },
482            )
483            .unwrap();
484        assert!(matches!(
485            completed.sends[0].message,
486            AcknowledgementOutcome::Completed { key: 7 }
487        ));
488        let stale = subject
489            .receive(
490                MailAddr(9),
491                AcknowledgementMessage::Acknowledge {
492                    key: 7,
493                    participant: 2,
494                    reply_to: reply(),
495                },
496            )
497            .unwrap();
498        assert!(matches!(
499            stale.sends[0].message,
500            AcknowledgementOutcome::Rejected(AcknowledgementError::Completed(
501                AcknowledgementInput::Acknowledge {
502                    key: 7,
503                    participant: 2,
504                }
505            ))
506        ));
507    }
508
509    #[test]
510    fn rejection_and_cancellation_do_not_conflate_phases() {
511        let mut subject = (Subject::new()).initialize().unwrap().behavior;
512        let started = subject
513            .receive(
514                MailAddr(9),
515                AcknowledgementMessage::Begin {
516                    key: 1,
517                    participants: vec![3],
518                    reply_to: reply(),
519                },
520            )
521            .unwrap();
522        assert!(matches!(
523            started.sends[0].message,
524            AcknowledgementOutcome::Started {
525                key: 1,
526                remaining: 1
527            }
528        ));
529        let unexpected = subject
530            .receive(
531                MailAddr(9),
532                AcknowledgementMessage::Acknowledge {
533                    key: 1,
534                    participant: 4,
535                    reply_to: reply(),
536                },
537            )
538            .unwrap();
539        assert!(matches!(
540            unexpected.sends[0].message,
541            AcknowledgementOutcome::Rejected(AcknowledgementError::UnexpectedParticipant { .. })
542        ));
543        let cancelled = subject
544            .receive(
545                MailAddr(9),
546                AcknowledgementMessage::Cancel {
547                    key: 1,
548                    reply_to: reply(),
549                },
550            )
551            .unwrap();
552        assert!(matches!(
553            cancelled.sends[0].message,
554            AcknowledgementOutcome::Cancelled { key: 1 }
555        ));
556        assert!(matches!(
557            subject.records()[0].state,
558            AcknowledgementState::Cancelled
559        ));
560    }
561}