Skip to main content

behavior_actors/routing/
deduplicator.rs

1//! Explicitly bounded first-seen delivery policy.
2
3use std::collections::VecDeque;
4
5use super::DeliveryOutcomes;
6
7use behavior::{
8    Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
9    SendEffects, User,
10};
11use thiserror::Error;
12
13use crate::DeliveryRoute;
14
15/// Complete observable retention state of a [`Deduplicator`].
16#[derive(Debug, Clone, PartialEq, Eq)]
17pub struct DeduplicatorState<K> {
18    /// Positive maximum number of retained keys.
19    pub capacity: usize,
20    retained: Vec<K>,
21}
22
23impl<K> DeduplicatorState<K> {
24    /// Retained keys from oldest to newest.
25    #[must_use]
26    pub fn retained(&self) -> &[K] {
27        &self.retained
28    }
29}
30
31/// Factual outcome of one keyed delivery attempt.
32#[derive(Debug, PartialEq, Eq)]
33pub enum DeduplicatorOutcome<K, T> {
34    /// The key was new, the value was delivered, and retention was committed.
35    Delivered {
36        /// Accepted key.
37        key: K,
38        /// Oldest key removed to remain within capacity, if any.
39        evicted: Option<K>,
40    },
41    /// The key is retained, so ownership of the undelivered value is returned.
42    Duplicate {
43        /// Duplicate key.
44        key: K,
45        /// Undelivered value.
46        value: T,
47    },
48}
49
50/// Commands accepted by [`Deduplicator`].
51pub enum DeduplicatorMessage<K, T, TargetRoute, ReplyRoute> {
52    /// Deliver `value` only if `key` is absent from bounded retention.
53    Deliver {
54        /// Application-defined idempotency key.
55        key: K,
56        /// Owned value.
57        value: T,
58        /// Typed destination for a first-seen value.
59        to: TargetRoute,
60        /// Typed recipient for the complete result.
61        reply_to: ReplyRoute,
62    },
63}
64
65/// Invalid deduplication-retention definition.
66#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
67pub enum DeduplicatorConfigError {
68    /// Zero retention would make every accepted key immediately forgotten.
69    #[error("deduplicator capacity must be positive")]
70    ZeroCapacity,
71}
72
73/// Bounded, deterministic first-seen delivery behavior.
74///
75/// A retained key rejects delivery and returns the owned payload. A new key is
76/// committed before its delivery is emitted; when retention is full, the
77/// oldest key is returned in the result before the new key becomes newest.
78/// Duplicate attempts do not refresh retention. Initialization is empty, the
79/// fold creates no actors, and it never terminates by policy. The finite FIFO
80/// retention window is deliberate Bombay policy and makes the limitation of
81/// deduplication explicit: an evicted key can be accepted again. Physical
82/// mailbox delivery remains a runtime responsibility. Construction rejects
83/// zero capacity and transitions have no panic path.
84pub struct Deduplicator<
85    A: Address,
86    K,
87    T,
88    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
89    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = DeduplicatorOutcome<K, T>>>,
90> {
91    capacity: usize,
92    retained: VecDeque<K>,
93    marker: core::marker::PhantomData<fn() -> (A, T, TargetRoute, ReplyRoute)>,
94}
95
96impl<A, K, T, TargetRoute, ReplyRoute> Deduplicator<A, K, T, TargetRoute, ReplyRoute>
97where
98    A: Address,
99    K: Clone + Eq,
100    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
101    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = DeduplicatorOutcome<K, T>>>,
102{
103    /// Construct empty positive bounded retention.
104    ///
105    /// # Errors
106    ///
107    /// Returns [`DeduplicatorConfigError::ZeroCapacity`] for zero capacity.
108    pub fn new(capacity: usize) -> Result<Self, DeduplicatorConfigError> {
109        if capacity == 0 {
110            return Err(DeduplicatorConfigError::ZeroCapacity);
111        }
112        Ok(Self {
113            capacity,
114            retained: VecDeque::with_capacity(capacity),
115            marker: core::marker::PhantomData,
116        })
117    }
118
119    /// Return a snapshot of the complete retention policy state.
120    #[must_use]
121    pub fn state(&self) -> DeduplicatorState<K> {
122        DeduplicatorState {
123            capacity: self.capacity,
124            retained: self.retained.iter().cloned().collect(),
125        }
126    }
127}
128
129impl<A, K, T, TargetRoute, ReplyRoute> BehaviorBase
130    for Deduplicator<A, K, T, TargetRoute, ReplyRoute>
131where
132    A: Address,
133    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
134    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = DeduplicatorOutcome<K, T>>>,
135{
136    type Base = Self;
137
138    fn base(&self) -> &Self {
139        self
140    }
141}
142
143impl<A, K, T, TargetRoute, ReplyRoute> behavior::Protocol
144    for Deduplicator<A, K, T, TargetRoute, ReplyRoute>
145where
146    A: Address,
147    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
148    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = DeduplicatorOutcome<K, T>>>,
149{
150    type Addr = A;
151    type Msg = DeduplicatorMessage<K, T, TargetRoute, ReplyRoute>;
152}
153
154impl<A, K, T, TargetRoute, ReplyRoute> Behavior for Deduplicator<A, K, T, TargetRoute, ReplyRoute>
155where
156    A: Address,
157    K: Clone + Eq,
158    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
159    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = DeduplicatorOutcome<K, T>>>,
160    TargetRoute::Sends:
161        behavior::SendsFor<User<A, DeduplicatorMessage<K, T, TargetRoute, ReplyRoute>>>,
162    ReplyRoute::Sends:
163        behavior::SendsFor<User<A, DeduplicatorMessage<K, T, TargetRoute, ReplyRoute>>>,
164{
165    type Protocol = Self;
166    type Event = User<A, behavior::BehaviorMessage<Self>>;
167    type Sends = DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>;
168    type Ph = Never;
169    type Error = Never;
170    type Birth = NoBirths;
171
172    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
173        let DeduplicatorMessage::Deliver {
174            key,
175            value,
176            to,
177            reply_to,
178        } = event.message;
179        if self.retained.contains(&key) {
180            return Ok(Actions::send(DeliveryOutcomes {
181                deliveries: TargetRoute::Sends::empty(),
182                outcomes: reply_to.deliver(DeduplicatorOutcome::Duplicate { key, value }),
183            }));
184        }
185
186        let evicted = if self.retained.len() == self.capacity {
187            self.retained.pop_front()
188        } else {
189            None
190        };
191        self.retained.push_back(key.clone());
192        Ok(Actions::send(DeliveryOutcomes {
193            deliveries: to.deliver(value),
194            outcomes: reply_to.deliver(DeduplicatorOutcome::Delivered { key, evicted }),
195        }))
196    }
197}
198
199#[cfg(test)]
200mod tests {
201    use crate::Activate as _;
202    use behavior::{Delivery, MailAddr, Recipient};
203
204    use super::*;
205
206    struct Target;
207    struct Reply;
208
209    macro_rules! leaf {
210        ($name:ident, $msg:ty) => {
211            impl behavior::Protocol for $name {
212                type Addr = MailAddr;
213                type Msg = $msg;
214            }
215
216            impl Behavior for $name {
217                type Protocol = Self;
218                type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
219                type Sends = Vec<Never>;
220                type Ph = Never;
221                type Error = Never;
222                type Birth = NoBirths;
223                fn transition(
224                    &mut self,
225                    _: behavior::ActiveTurn,
226                    _: Self::Event,
227                ) -> BehaviorActed<Self> {
228                    Ok(Actions::cont())
229                }
230            }
231        };
232    }
233
234    leaf!(Target, u8);
235    leaf!(Reply, DeduplicatorOutcome<u8, u8>);
236
237    type Subject = Deduplicator<MailAddr, u8, u8, Recipient<Target>, Recipient<Reply>>;
238
239    fn deliver(
240        subject: &mut crate::Active<Subject>,
241        key: u8,
242        value: u8,
243    ) -> Actions<
244        MailAddr,
245        Never,
246        DeliveryOutcomes<Vec<Delivery<Target>>, Vec<Delivery<Reply>>>,
247        NoBirths,
248    > {
249        subject
250            .receive(
251                MailAddr(9),
252                DeduplicatorMessage::Deliver {
253                    key,
254                    value,
255                    to: Recipient::global(MailAddr(1)),
256                    reply_to: Recipient::global(MailAddr(2)),
257                },
258            )
259            .unwrap()
260    }
261
262    #[test]
263    fn duplicate_returns_value_without_refreshing_retention() {
264        let mut subject = (Subject::new(2).unwrap()).initialize().unwrap().behavior;
265        let first = deliver(&mut subject, 1, 10);
266        assert_eq!(first.sends.deliveries.len(), 1);
267        assert!(matches!(
268            first.sends.outcomes[0].message,
269            DeduplicatorOutcome::Delivered {
270                key: 1,
271                evicted: None
272            }
273        ));
274        let duplicate = deliver(&mut subject, 1, 11);
275        assert!(duplicate.sends.deliveries.is_empty());
276        assert!(matches!(
277            duplicate.sends.outcomes[0].message,
278            DeduplicatorOutcome::Duplicate { key: 1, value: 11 }
279        ));
280        assert_eq!(subject.state().retained().to_vec(), vec![1]);
281    }
282
283    #[test]
284    fn eviction_is_explicit_and_allows_later_readmission() {
285        let mut subject = (Subject::new(2).unwrap()).initialize().unwrap().behavior;
286        let first = deliver(&mut subject, 1, 10);
287        assert!(matches!(
288            first.sends.outcomes[0].message,
289            DeduplicatorOutcome::Delivered {
290                key: 1,
291                evicted: None
292            }
293        ));
294        let second = deliver(&mut subject, 2, 20);
295        assert!(matches!(
296            second.sends.outcomes[0].message,
297            DeduplicatorOutcome::Delivered {
298                key: 2,
299                evicted: None
300            }
301        ));
302        let third = deliver(&mut subject, 3, 30);
303        assert!(matches!(
304            third.sends.outcomes[0].message,
305            DeduplicatorOutcome::Delivered {
306                key: 3,
307                evicted: Some(1)
308            }
309        ));
310        assert_eq!(subject.state().retained().to_vec(), vec![2, 3]);
311        let readmitted = deliver(&mut subject, 1, 12);
312        assert_eq!(readmitted.sends.deliveries.len(), 1);
313    }
314
315    #[test]
316    fn zero_capacity_is_rejected() {
317        assert!(matches!(
318            Subject::new(0),
319            Err(DeduplicatorConfigError::ZeroCapacity)
320        ));
321    }
322}