Skip to main content

behavior_actors/time/
receive_timeout.rs

1//! Pure receive-inactivity composition. Relative scheduling is a request to
2//! the emitting actor's interpreter and never observes a clock here.
3
4use std::time::Duration;
5
6use super::domain::{TimerAdmission, TimerLease};
7use super::event::{TimedEvent, TimedReaction};
8use crate::protocol::{ScheduleAfter, TimerId};
9use behavior::Step;
10use behavior::{
11    Actions, Address, Behavior, BirthMode, EventLayer, InterpreterRequests, SendEffects, SendLayer,
12    UserEvent,
13};
14
15pub(crate) type ReceiveTimeoutActions<B> = Actions<
16    behavior::BehaviorAddr<B>,
17    <B as Behavior>::Ph,
18    SendLayer<InterpreterRequests<ScheduleAfter>, <B as Behavior>::Sends>,
19    <B as Behavior>::Birth,
20>;
21
22/// Infallible fold invoked for the accepted inactivity generation.
23///
24/// ```compile_fail,E0308
25/// # struct App;
26/// # impl behavior::Protocol for App { type Addr = behavior::MailAddr; type Msg = (); }
27/// # impl behavior::Behavior for App {
28/// #   type Protocol = Self; type Event = behavior::User<behavior::MailAddr, ()>; type Sends = Vec<behavior::Never>;
29/// #   type Ph = behavior::Never; type Error = behavior::Never; type Birth = behavior::NoBirths;
30/// #   fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> behavior::BehaviorActed<Self> { Ok(behavior::Actions::cont()) }
31/// # }
32/// fn fallible(_: &mut App) -> behavior::BehaviorActed<App> { Ok(behavior::Actions::cont()) }
33/// let _ = behavior_actors::ReceiveTimeout::new(App, behavior_actors::TimerId(1), std::time::Duration::from_secs(1), fallible);
34/// ```
35/// A pure one-notification-per-idle-period receive timeout.
36///
37/// Only successful user communications are activity. Timer, peer, child,
38/// worker, and shutdown interpreter events compose through this wrapper but never
39/// rearm it. A matching timeout invokes the infallible reaction and consumes
40/// the live generation. If that reaction continues, the timeout remains
41/// unarmed until another successful continuing user communication. Ordinary
42/// delegated transitions retain the wrapped error type.
43pub struct ReceiveTimeout<B: Behavior> {
44    inner: B,
45    id: TimerId,
46    after: Duration,
47    timer: TimerLease,
48    on_elapsed: TimedReaction<B>,
49}
50
51impl<B: Behavior> ReceiveTimeout<B> {
52    /// Wrap `inner` with a relative inactivity timeout.
53    ///
54    /// Initialization and each successful continuing user fold stage a fresh
55    /// timer generation. Interpreter events do not reset inactivity.
56    #[must_use]
57    pub fn new(inner: B, id: TimerId, after: Duration, on_elapsed: TimedReaction<B>) -> Self {
58        Self {
59            inner,
60            id,
61            after,
62            timer: TimerLease::new(),
63            on_elapsed,
64        }
65    }
66
67    fn schedule(&mut self) -> InterpreterRequests<ScheduleAfter> {
68        self.timer
69            .arm()
70            .map_or_else(InterpreterRequests::empty, |generation| {
71                InterpreterRequests::one(ScheduleAfter::new(self.id, generation, self.after))
72            })
73    }
74
75    fn wrap(
76        actions: Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
77        own: InterpreterRequests<ScheduleAfter>,
78    ) -> ReceiveTimeoutActions<B> {
79        actions.map_sends(|inner| SendLayer::new(own, inner))
80    }
81
82    fn terminal(actions: &Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>) -> bool {
83        matches!(actions.become_, Step::Stop(_))
84    }
85}
86
87impl<B: Behavior + behavior::BehaviorBase> behavior::BehaviorBase for ReceiveTimeout<B> {
88    type Base = B::Base;
89
90    fn base(&self) -> &Self::Base {
91        self.inner.base()
92    }
93}
94
95impl<B> crate::StashStatus for ReceiveTimeout<B>
96where
97    B: Behavior + crate::StashStatus,
98{
99    fn stashed_messages(&self) -> usize {
100        self.inner.stashed_messages()
101    }
102}
103
104impl<B, A, Ph, Sends, Br> Behavior for ReceiveTimeout<B>
105where
106    A: Address,
107    Sends: SendEffects + behavior::SendsFor<B::Event>,
108    Br: BirthMode,
109    B: Behavior<Ph = Ph, Sends = Sends, Birth = Br>,
110    B::Protocol: behavior::Protocol<Addr = A>,
111{
112    type Protocol = B::Protocol;
113    type Event = TimedEvent<B::Event>;
114    type Sends = SendLayer<InterpreterRequests<ScheduleAfter>, Sends>;
115    type Ph = Ph;
116    type Error = B::Error;
117    type Birth = Br;
118
119    fn init(
120        &mut self,
121        _: behavior::InitializationTurn,
122    ) -> Result<ReceiveTimeoutActions<B>, Self::Error> {
123        let actions = behavior::initialize(&mut self.inner)?;
124        let own = if Self::terminal(&actions) {
125            self.timer.disarm();
126            InterpreterRequests::empty()
127        } else {
128            self.schedule()
129        };
130        Ok(Self::wrap(actions, own))
131    }
132
133    fn transition(
134        &mut self,
135        _: behavior::ActiveTurn,
136        event: Self::Event,
137    ) -> Result<ReceiveTimeoutActions<B>, Self::Error> {
138        match event {
139            EventLayer::Owned(elapsed) if elapsed.id == self.id => {
140                match self.timer.accept(elapsed.generation) {
141                    TimerAdmission::Accepted => {
142                        let actions = (self.on_elapsed)(&mut self.inner);
143                        Ok(Self::wrap(actions, InterpreterRequests::empty()))
144                    }
145                    TimerAdmission::Ignored => Ok(Actions::cont()),
146                }
147            }
148            EventLayer::Owned(_) => Ok(Actions::cont()),
149            EventLayer::Inner(event) => match event.into_user() {
150                Ok(user) => {
151                    let event = B::Event::user(user.from, user.message);
152                    let actions = behavior::delegate_transition(&mut self.inner, event)?;
153                    let own = if Self::terminal(&actions) {
154                        self.timer.disarm();
155                        InterpreterRequests::empty()
156                    } else {
157                        self.schedule()
158                    };
159                    Ok(Self::wrap(actions, own))
160                }
161                Err(service) => {
162                    let actions = behavior::delegate_transition(&mut self.inner, service)?;
163                    if Self::terminal(&actions) {
164                        self.timer.disarm();
165                    }
166                    Ok(Self::wrap(actions, InterpreterRequests::empty()))
167                }
168            },
169        }
170    }
171}
172
173#[cfg(test)]
174mod tests {
175    use super::*;
176    use crate::TimerGeneration;
177    use behavior::{MailAddr, Never, NoBirths, User};
178
179    struct Count(u8);
180
181    impl behavior::Protocol for Count {
182        type Addr = MailAddr;
183        type Msg = ();
184    }
185
186    impl Behavior for Count {
187        type Protocol = Self;
188        type Event = User<MailAddr, ()>;
189        type Sends = Vec<Never>;
190        type Ph = Never;
191        type Error = Never;
192        type Birth = NoBirths;
193
194        fn transition(
195            &mut self,
196            _: behavior::ActiveTurn,
197            _: Self::Event,
198        ) -> behavior::BehaviorActed<Self> {
199            self.0 += 1;
200            Ok(Actions::cont())
201        }
202    }
203
204    type CountBehavior = Count;
205
206    impl behavior::BehaviorBase for Count {
207        type Base = Self;
208
209        fn base(&self) -> &Self {
210            self
211        }
212    }
213
214    fn elapsed(_inner: &mut CountBehavior) -> Actions<MailAddr, Never, Vec<Never>, NoBirths> {
215        Actions::cont()
216    }
217
218    #[tokio::test]
219    async fn exhaustion_retires_only_the_timer_and_preserves_the_inner_transition() {
220        let mut timeout =
221            ReceiveTimeout::new(Count(0), TimerId(0), Duration::from_secs(1), elapsed);
222        let initialized = behavior::initialize(&mut timeout).unwrap();
223        assert_eq!(initialized.sends.owned.len(), 1);
224        assert!(initialized.sends.inner.is_empty());
225        assert!(initialized.creates.is_empty());
226        assert!(matches!(initialized.become_, Step::Continue));
227        timeout.timer = TimerLease::idle(TimerGeneration(u64::MAX));
228
229        let actions = behavior::delegate_transition(
230            &mut timeout,
231            EventLayer::Inner(User::user(MailAddr(1), ())),
232        )
233        .unwrap();
234
235        assert!(actions.sends.owned.is_empty());
236        assert!(actions.creates.is_empty());
237        assert!(matches!(actions.become_, Step::Continue));
238        assert_eq!(behavior::BehaviorBase::base(&timeout).0, 1);
239        assert_eq!(timeout.timer.live(), None);
240    }
241}