1use std::collections::VecDeque;
4
5#[cfg(test)]
6use behavior::MessageProtocol;
7use behavior::{
8 Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
9 SendEffects, User,
10};
11use thiserror::Error;
12
13use super::DeliveryOutcomes;
14use crate::DeliveryRoute;
15
16#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
18pub enum OverflowPolicy {
19 Reject,
21 DropOldest,
23 DropNewest,
25}
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
29pub enum BufferRejection {
30 Full,
32 DroppedNewest,
34}
35
36#[derive(Debug, PartialEq, Eq)]
38pub enum BufferOutcome<T> {
39 Accepted {
41 depth: usize,
43 },
44 Rejected {
46 value: T,
48 reason: BufferRejection,
50 },
51 Evicted {
53 value: T,
55 },
56 Released {
58 remaining: usize,
60 },
61 Empty,
63}
64
65pub struct Buffered<T, Route> {
67 pub value: T,
69 pub reply_to: Route,
71}
72
73pub struct BufferState<T, Route> {
75 pub capacity: usize,
77 pub overflow: OverflowPolicy,
79 queued: VecDeque<Buffered<T, Route>>,
80}
81
82impl<T, Route> BufferState<T, Route> {
83 #[must_use]
85 pub fn len(&self) -> usize {
86 self.queued.len()
87 }
88
89 #[must_use]
91 pub fn is_empty(&self) -> bool {
92 self.queued.is_empty()
93 }
94
95 #[must_use]
97 pub fn queued(&self) -> impl ExactSizeIterator<Item = &Buffered<T, Route>> {
98 self.queued.iter()
99 }
100}
101
102pub enum BufferMessage<T, TargetRoute, ReplyRoute> {
104 Offer {
106 value: T,
108 reply_to: ReplyRoute,
110 },
111 Release {
113 to: TargetRoute,
115 reply_to: ReplyRoute,
117 },
118}
119
120#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
122pub enum BufferConfigError {
123 #[error("buffer capacity must be positive")]
125 ZeroCapacity,
126}
127
128#[derive(Debug, Clone, Copy, PartialEq, Eq)]
130pub struct BufferConfiguration {
131 capacity: usize,
132 overflow: OverflowPolicy,
133}
134
135impl BufferConfiguration {
136 pub fn new(capacity: usize, overflow: OverflowPolicy) -> Result<Self, BufferConfigError> {
142 if capacity == 0 {
143 return Err(BufferConfigError::ZeroCapacity);
144 }
145 Ok(Self { capacity, overflow })
146 }
147}
148
149pub struct Buffer<A, T, TargetRoute, ReplyRoute>
165where
166 A: Address,
167 TargetRoute: DeliveryRoute,
168 TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
169 ReplyRoute: DeliveryRoute,
170 ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
171{
172 state: BufferState<T, ReplyRoute>,
173 marker: core::marker::PhantomData<fn() -> (A, TargetRoute)>,
174}
175
176impl<A, T, TargetRoute, ReplyRoute> Buffer<A, T, TargetRoute, ReplyRoute>
177where
178 A: Address,
179 TargetRoute: DeliveryRoute,
180 TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
181 ReplyRoute: DeliveryRoute,
182 ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
183{
184 #[must_use]
186 pub fn new(configuration: BufferConfiguration) -> Self {
187 Self {
188 state: BufferState {
189 capacity: configuration.capacity,
190 overflow: configuration.overflow,
191 queued: VecDeque::with_capacity(configuration.capacity),
192 },
193 marker: core::marker::PhantomData,
194 }
195 }
196
197 #[must_use]
199 pub const fn state(&self) -> &BufferState<T, ReplyRoute> {
200 &self.state
201 }
202
203 fn actions(
204 deliveries: TargetRoute::Sends,
205 outcomes: ReplyRoute::Sends,
206 ) -> Actions<A, Never, DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>, NoBirths> {
207 Actions::send(DeliveryOutcomes {
208 deliveries,
209 outcomes,
210 })
211 }
212}
213
214impl<A, T, TargetRoute, ReplyRoute> BehaviorBase for Buffer<A, T, TargetRoute, ReplyRoute>
215where
216 A: Address,
217 TargetRoute: DeliveryRoute,
218 TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
219 ReplyRoute: DeliveryRoute,
220 ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
221{
222 type Base = Self;
223
224 fn base(&self) -> &Self {
225 self
226 }
227}
228
229impl<A, T, TargetRoute, ReplyRoute> behavior::Protocol for Buffer<A, T, TargetRoute, ReplyRoute>
230where
231 A: Address,
232 TargetRoute: DeliveryRoute,
233 TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
234 ReplyRoute: DeliveryRoute,
235 ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
236{
237 type Addr = A;
238 type Msg = BufferMessage<T, TargetRoute, ReplyRoute>;
239}
240
241impl<A, T, TargetRoute, ReplyRoute> Behavior for Buffer<A, T, TargetRoute, ReplyRoute>
242where
243 A: Address,
244 TargetRoute: DeliveryRoute,
245 TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
246 ReplyRoute: DeliveryRoute + Clone,
247 ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
248 TargetRoute::Sends: behavior::SendsFor<User<A, BufferMessage<T, TargetRoute, ReplyRoute>>>,
249 ReplyRoute::Sends: behavior::SendsFor<User<A, BufferMessage<T, TargetRoute, ReplyRoute>>>,
250{
251 type Protocol = Self;
252 type Event = User<A, behavior::BehaviorMessage<Self>>;
253 type Sends = DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>;
254 type Ph = Never;
255 type Error = Never;
256 type Birth = NoBirths;
257
258 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
259 match event.message {
260 BufferMessage::Offer { value, reply_to }
261 if self.state.queued.len() < self.state.capacity =>
262 {
263 self.state.queued.push_back(Buffered {
264 value,
265 reply_to: reply_to.clone(),
266 });
267 Ok(Self::actions(
268 TargetRoute::Sends::empty(),
269 reply_to.deliver(BufferOutcome::Accepted {
270 depth: self.state.queued.len(),
271 }),
272 ))
273 }
274 BufferMessage::Offer { value, reply_to } => match self.state.overflow {
275 OverflowPolicy::Reject => Ok(Self::actions(
276 TargetRoute::Sends::empty(),
277 reply_to.deliver(BufferOutcome::Rejected {
278 value,
279 reason: BufferRejection::Full,
280 }),
281 )),
282 OverflowPolicy::DropNewest => Ok(Self::actions(
283 TargetRoute::Sends::empty(),
284 reply_to.deliver(BufferOutcome::Rejected {
285 value,
286 reason: BufferRejection::DroppedNewest,
287 }),
288 )),
289 OverflowPolicy::DropOldest => {
290 let evicted = self
293 .state
294 .queued
295 .pop_front()
296 .expect("positive full buffer contains an oldest value");
297 let mut outcomes = evicted.reply_to.deliver(BufferOutcome::Evicted {
298 value: evicted.value,
299 });
300 self.state.queued.push_back(Buffered {
301 value,
302 reply_to: reply_to.clone(),
303 });
304 outcomes.append(reply_to.deliver(BufferOutcome::Accepted {
305 depth: self.state.queued.len(),
306 }));
307 Ok(Self::actions(TargetRoute::Sends::empty(), outcomes))
308 }
309 },
310 BufferMessage::Release { to, reply_to } => {
311 let Some(buffered) = self.state.queued.pop_front() else {
312 return Ok(Self::actions(
313 TargetRoute::Sends::empty(),
314 reply_to.deliver(BufferOutcome::Empty),
315 ));
316 };
317 Ok(Self::actions(
318 to.deliver(buffered.value),
319 reply_to.deliver(BufferOutcome::Released {
320 remaining: self.state.queued.len(),
321 }),
322 ))
323 }
324 }
325 }
326}
327
328#[cfg(test)]
329mod tests {
330 use super::*;
331 use crate::Activate as _;
332 use behavior::{Delivery, MailAddr, Recipient};
333
334 fn active(
335 policy: OverflowPolicy,
336 ) -> crate::Active<
337 Buffer<
338 MailAddr,
339 u8,
340 Recipient<MessageProtocol<MailAddr, u8>>,
341 Recipient<MessageProtocol<MailAddr, BufferOutcome<u8>>>,
342 >,
343 > {
344 Buffer::new(BufferConfiguration::new(2, policy).unwrap())
345 .initialize()
346 .unwrap()
347 .behavior
348 }
349
350 #[test]
351 fn buffer_uses_the_shared_delivery_outcome_product() {
352 type Target = Recipient<MessageProtocol<MailAddr, u8>>;
353 type Reply = Recipient<MessageProtocol<MailAddr, BufferOutcome<u8>>>;
354 type Expected = crate::DeliveryOutcomes<
355 <Target as DeliveryRoute>::Sends,
356 <Reply as DeliveryRoute>::Sends,
357 >;
358 fn exact<B: Behavior<Sends = Expected>>() {}
359 exact::<Buffer<MailAddr, u8, Target, Reply>>();
360 }
361
362 #[test]
363 fn zero_capacity_is_rejected_before_ownership_is_possible() {
364 assert!(matches!(
365 BufferConfiguration::new(0, OverflowPolicy::Reject),
366 Err(BufferConfigError::ZeroCapacity)
367 ));
368 }
369
370 #[test]
371 fn fifo_release_and_empty_result_use_disjoint_named_lanes() {
372 let reply = Recipient::from(MailAddr(8));
373 let target = Recipient::from(MailAddr(7));
374 let mut buffer = active(OverflowPolicy::Reject);
375 for value in [1, 2] {
376 let accepted = buffer
377 .receive(
378 MailAddr(9),
379 BufferMessage::Offer {
380 value,
381 reply_to: reply,
382 },
383 )
384 .unwrap();
385 assert!(accepted.sends.deliveries.is_empty());
386 assert!(matches!(
387 accepted.sends.outcomes.as_slice(),
388 [delivery]
389 if delivery.message == (BufferOutcome::Accepted {
390 depth: usize::from(value),
391 })
392 ));
393 }
394 for (value, remaining) in [(1, 1), (2, 0)] {
395 let released = buffer
396 .receive(
397 MailAddr(9),
398 BufferMessage::Release {
399 to: target,
400 reply_to: reply,
401 },
402 )
403 .unwrap();
404 assert!(released.sends.deliveries == vec![Delivery::new(target, value)]);
405 assert!(matches!(
406 released.sends.outcomes.as_slice(),
407 [delivery]
408 if delivery.message == (BufferOutcome::Released { remaining })
409 ));
410 }
411 let empty = buffer
412 .receive(
413 MailAddr(9),
414 BufferMessage::Release {
415 to: target,
416 reply_to: reply,
417 },
418 )
419 .unwrap();
420 assert!(empty.sends.deliveries.is_empty());
421 assert!(matches!(
422 empty.sends.outcomes.as_slice(),
423 [delivery] if delivery.message == BufferOutcome::Empty
424 ));
425 }
426
427 #[test]
428 fn every_overflow_policy_preserves_or_returns_all_owned_values() {
429 let first_reply = Recipient::from(MailAddr(1));
430 let newest_reply = Recipient::from(MailAddr(2));
431 for policy in [
432 OverflowPolicy::Reject,
433 OverflowPolicy::DropNewest,
434 OverflowPolicy::DropOldest,
435 ] {
436 let mut buffer = active(policy);
437 for value in [10, 11] {
438 let offered = buffer
439 .receive(
440 MailAddr(9),
441 BufferMessage::Offer {
442 value,
443 reply_to: first_reply,
444 },
445 )
446 .unwrap();
447 assert!(offered.sends.deliveries.is_empty());
448 assert!(matches!(
449 offered.sends.outcomes.as_slice(),
450 [delivery]
451 if delivery.message
452 == (BufferOutcome::Accepted {
453 depth: usize::from(value - 9),
454 })
455 ));
456 assert!(offered.creates.is_empty());
457 assert_eq!(offered.become_, behavior::Step::Continue);
458 }
459 let overflow = buffer
460 .receive(
461 MailAddr(9),
462 BufferMessage::Offer {
463 value: 12,
464 reply_to: newest_reply,
465 },
466 )
467 .unwrap();
468 match policy {
469 OverflowPolicy::Reject => assert!(matches!(
470 overflow.sends.outcomes.as_slice(),
471 [delivery]
472 if delivery.message == (BufferOutcome::Rejected {
473 value: 12,
474 reason: BufferRejection::Full,
475 })
476 )),
477 OverflowPolicy::DropNewest => assert!(matches!(
478 overflow.sends.outcomes.as_slice(),
479 [delivery]
480 if delivery.message == (BufferOutcome::Rejected {
481 value: 12,
482 reason: BufferRejection::DroppedNewest,
483 })
484 )),
485 OverflowPolicy::DropOldest => {
486 assert!(matches!(
487 overflow.sends.outcomes.as_slice(),
488 [evicted, accepted]
489 if evicted.message == (BufferOutcome::Evicted { value: 10 })
490 && accepted.message == (BufferOutcome::Accepted { depth: 2 })
491 ));
492 assert!(buffer.state().queued().map(|item| item.value).eq([11, 12]));
493 }
494 }
495 assert_eq!(buffer.state().len(), 2);
496 }
497 }
498}