1use 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#[derive(Debug, Clone, PartialEq, Eq)]
17pub struct DeduplicatorState<K> {
18 pub capacity: usize,
20 retained: Vec<K>,
21}
22
23impl<K> DeduplicatorState<K> {
24 #[must_use]
26 pub fn retained(&self) -> &[K] {
27 &self.retained
28 }
29}
30
31#[derive(Debug, PartialEq, Eq)]
33pub enum DeduplicatorOutcome<K, T> {
34 Delivered {
36 key: K,
38 evicted: Option<K>,
40 },
41 Duplicate {
43 key: K,
45 value: T,
47 },
48}
49
50pub enum DeduplicatorMessage<K, T, TargetRoute, ReplyRoute> {
52 Deliver {
54 key: K,
56 value: T,
58 to: TargetRoute,
60 reply_to: ReplyRoute,
62 },
63}
64
65#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
67pub enum DeduplicatorConfigError {
68 #[error("deduplicator capacity must be positive")]
70 ZeroCapacity,
71}
72
73pub 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 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 #[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}