1use std::collections::VecDeque;
4
5use crate::behavior::{Actions, Address, Behavior, BirthMode, SendAlgebra, User, UserEvent};
6use crate::verdict::{Never, Step};
7
8#[derive(Debug, Clone, Copy, PartialEq, Eq)]
9pub enum StashRoute {
10 Stash,
11 Deliver,
12 Release,
13}
14
15pub struct Stashing<B: Behavior> {
16 inner: B,
17 route: fn(&B::Msg) -> StashRoute,
18 held: VecDeque<User<B::Addr, B::Msg>>,
19}
20
21impl<B: Behavior<Ph = Never>> Stashing<B> {
22 #[must_use]
23 pub fn new(inner: B, route: fn(&B::Msg) -> StashRoute) -> Self {
24 Self {
25 inner,
26 route,
27 held: VecDeque::new(),
28 }
29 }
30
31 #[must_use]
32 pub fn inner(&self) -> &B {
33 &self.inner
34 }
35
36 #[must_use]
37 pub fn held(&self) -> usize {
38 self.held.len()
39 }
40}
41
42impl<B, A, Sends, Br> Stashing<B>
43where
44 A: Address,
45 Sends: SendAlgebra,
46 Br: BirthMode,
47 B: Behavior<Addr = A, Ph = Never, Sends = Sends, Birth = Br>,
48{
49 async fn drain_into(
50 &mut self,
51 acc: &mut Actions<B::Addr, Never, B::Sends, B::Birth>,
52 ) -> Result<(), B::Error> {
53 let mut batch: VecDeque<_> = self.held.drain(..).collect();
54 while let Some(user) = batch.pop_front() {
55 match (self.route)(&user.message) {
56 StashRoute::Stash => self.held.push_back(user),
57 StashRoute::Deliver | StashRoute::Release => {
58 let actions = self
59 .inner
60 .step(B::Event::user(user.from, user.message))
61 .await?;
62 acc.sends.append(actions.sends);
63 acc.creates.extend(actions.creates);
64 if let Step::Stop(exit) = actions.become_ {
65 self.held.extend(batch);
66 acc.become_ = Step::Stop(exit);
67 return Ok(());
68 }
69 }
70 }
71 }
72 Ok(())
73 }
74}
75
76impl<B, A, Sends, Br> Behavior for Stashing<B>
77where
78 A: Address + Send,
79 Sends: SendAlgebra + Send,
80 Br: BirthMode,
81 B: Behavior<Addr = A, Ph = Never, Sends = Sends, Birth = Br> + Send,
82 A::Nonce: Send,
83 B::Msg: Send,
84 B::Event: Send,
85 Br::Child: Send,
86{
87 type Addr = A;
88 type Msg = B::Msg;
89 type Event = B::Event;
90 type Sends = Sends;
91 type Ph = Never;
92 type Error = B::Error;
93 type Birth = Br;
94
95 async fn init(&mut self) -> Result<Actions<A, Never, Sends, Br>, B::Error> {
96 self.inner.init().await
97 }
98
99 async fn step(&mut self, event: B::Event) -> Result<Actions<A, Never, Sends, Br>, B::Error> {
100 let user = match event.into_user() {
101 Ok(user) => user,
102 Err(other) => return self.inner.step(other).await,
103 };
104 match (self.route)(&user.message) {
105 StashRoute::Stash => {
106 self.held.push_back(user);
107 Ok(Actions::cont())
108 }
109 StashRoute::Deliver => {
110 self.inner
111 .step(B::Event::user(user.from, user.message))
112 .await
113 }
114 StashRoute::Release => {
115 let mut actions = self
116 .inner
117 .step(B::Event::user(user.from, user.message))
118 .await?;
119 if !matches!(actions.become_, Step::Stop(_)) {
120 self.drain_into(&mut actions).await?;
121 }
122 Ok(actions)
123 }
124 }
125 }
126}