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