Skip to main content

behavior/
stashing.rs

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