Skip to main content

behavior/
stashing.rs

1//! Pure message holding and replay composition.
2
3use 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}