Skip to main content

behavior/timing/
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::TimerLease;
7use super::event::TimedEvent;
8use crate::Step;
9use crate::behavior::{
10    Actions, Address, Behavior, BirthMode, SendAlgebra, ServiceSends, UserEvent,
11};
12use crate::protocol::{ScheduleAfter, TimerElapsed, TimerId};
13use crate::{Own, RouteInput, SendInput};
14
15pub type ReceiveTimeoutEvent<E> = TimedEvent<E>;
16
17pub type ReceiveTimeoutReaction<B> = fn(
18    &mut B,
19) -> Result<
20    Actions<
21        <B as Behavior>::Addr,
22        <B as Behavior>::Ph,
23        <B as Behavior>::Sends,
24        <B as Behavior>::Birth,
25    >,
26    <B as Behavior>::Error,
27>;
28
29/// Named effect lanes added by [`ReceiveTimeout`].
30#[derive(Debug, Clone, PartialEq, Eq)]
31pub struct ReceiveTimeoutSends<Sends> {
32    pub behavior: Sends,
33    pub schedules: ServiceSends<ScheduleAfter>,
34}
35
36impl<Sends: SendAlgebra> SendAlgebra for ReceiveTimeoutSends<Sends> {
37    fn empty() -> Self {
38        Self {
39            behavior: Sends::empty(),
40            schedules: ServiceSends::empty(),
41        }
42    }
43
44    fn append(&mut self, other: Self) {
45        self.behavior.append(other.behavior);
46        self.schedules.append(other.schedules);
47    }
48}
49
50impl<Sends> SendInput<ScheduleAfter, Own> for ReceiveTimeoutSends<Sends> {
51    fn emit(&mut self, input: ScheduleAfter) {
52        self.schedules.send(input);
53    }
54}
55
56pub(crate) type ReceiveTimeoutActions<B> = Actions<
57    <B as Behavior>::Addr,
58    <B as Behavior>::Ph,
59    ReceiveTimeoutSends<<B as Behavior>::Sends>,
60    <B as Behavior>::Birth,
61>;
62
63/// A pure one-notification-per-idle-period receive timeout.
64///
65/// Only successful user communications are activity. Timer, peer, child,
66/// worker, and shutdown service events compose through this wrapper but never
67/// rearm it. A matching timeout consumes the live generation before invoking
68/// the reaction; if that reaction continues, the timeout remains unarmed until
69/// another successful continuing user communication.
70pub struct ReceiveTimeout<B: Behavior> {
71    inner: B,
72    id: TimerId,
73    after: Duration,
74    timer: TimerLease,
75    on_elapsed: ReceiveTimeoutReaction<B>,
76}
77
78impl<B: Behavior> ReceiveTimeout<B> {
79    #[must_use]
80    pub(crate) fn new(
81        inner: B,
82        id: TimerId,
83        after: Duration,
84        on_elapsed: ReceiveTimeoutReaction<B>,
85    ) -> Self {
86        Self {
87            inner,
88            id,
89            after,
90            timer: TimerLease::new(),
91            on_elapsed,
92        }
93    }
94
95    fn schedule(&mut self) -> ServiceSends<ScheduleAfter> {
96        self.timer
97            .arm()
98            .map_or_else(ServiceSends::empty, |generation| {
99                ServiceSends::one(ScheduleAfter::new(self.id, generation, self.after))
100            })
101    }
102
103    fn wrap(
104        actions: Actions<B::Addr, B::Ph, B::Sends, B::Birth>,
105        own: ServiceSends<ScheduleAfter>,
106    ) -> ReceiveTimeoutActions<B> {
107        actions.map_sends(|behavior| ReceiveTimeoutSends {
108            behavior,
109            schedules: own,
110        })
111    }
112
113    fn terminal(actions: &Actions<B::Addr, B::Ph, B::Sends, B::Birth>) -> bool {
114        matches!(actions.become_, Step::Stop(_))
115    }
116}
117
118impl<B: Behavior + crate::BehaviorBase> crate::BehaviorBase for ReceiveTimeout<B> {
119    type Base = B::Base;
120
121    fn base(&self) -> &Self::Base {
122        self.inner.base()
123    }
124}
125
126impl<B> crate::StashStatus for ReceiveTimeout<B>
127where
128    B: Behavior + crate::StashStatus,
129{
130    fn stashed_messages(&self) -> usize {
131        self.inner.stashed_messages()
132    }
133}
134
135impl<B, A, Ph, Sends, Br> Behavior for ReceiveTimeout<B>
136where
137    A: Address,
138    Sends: SendAlgebra,
139    Br: BirthMode,
140    B: Behavior<Addr = A, Ph = Ph, Sends = Sends, Birth = Br>,
141    B::Event: crate::RouteInput<TimerElapsed>,
142{
143    type Addr = A;
144    type Msg = B::Msg;
145    type Event = ReceiveTimeoutEvent<B::Event>;
146    type Sends = ReceiveTimeoutSends<Sends>;
147    type Ph = Ph;
148    type Error = B::Error;
149    type Birth = Br;
150
151    fn init(
152        &mut self,
153        _: crate::InitializationTurn,
154    ) -> Result<ReceiveTimeoutActions<B>, Self::Error> {
155        let actions = crate::calculus::initialize(&mut self.inner)?;
156        let own = if Self::terminal(&actions) {
157            self.timer.disarm();
158            ServiceSends::empty()
159        } else {
160            self.schedule()
161        };
162        Ok(Self::wrap(actions, own))
163    }
164
165    fn transition(
166        &mut self,
167        _: crate::ActiveTurn,
168        event: Self::Event,
169    ) -> Result<ReceiveTimeoutActions<B>, Self::Error> {
170        match event {
171            ReceiveTimeoutEvent::Elapsed(elapsed)
172                if elapsed.id == self.id && self.timer.accept(elapsed.generation) =>
173            {
174                let actions = (self.on_elapsed)(&mut self.inner)?;
175                Ok(Self::wrap(actions, ServiceSends::empty()))
176            }
177            ReceiveTimeoutEvent::Elapsed(elapsed) if elapsed.id == self.id => Ok(Actions::cont()),
178            ReceiveTimeoutEvent::Elapsed(elapsed) => {
179                let Ok(inner) = B::Event::route(elapsed) else {
180                    return Ok(Actions::cont());
181                };
182                let actions = crate::calculus::delegate_transition(&mut self.inner, inner)?;
183                if Self::terminal(&actions) {
184                    self.timer.disarm();
185                }
186                Ok(Self::wrap(actions, ServiceSends::empty()))
187            }
188            ReceiveTimeoutEvent::Behavior(event) => match event.into_user() {
189                Ok(user) => {
190                    let event = B::Event::user(user.from, user.message);
191                    let actions = crate::calculus::delegate_transition(&mut self.inner, event)?;
192                    let own = if Self::terminal(&actions) {
193                        self.timer.disarm();
194                        ServiceSends::empty()
195                    } else {
196                        self.schedule()
197                    };
198                    Ok(Self::wrap(actions, own))
199                }
200                Err(service) => {
201                    let actions = crate::calculus::delegate_transition(&mut self.inner, service)?;
202                    if Self::terminal(&actions) {
203                        self.timer.disarm();
204                    }
205                    Ok(Self::wrap(actions, ServiceSends::empty()))
206                }
207            },
208        }
209    }
210}
211
212#[cfg(test)]
213mod tests {
214    use super::*;
215    use crate::{Acted, MailAddr, Never, NoBirths, TimerGeneration, User};
216
217    struct Count(u8);
218
219    impl Behavior for Count {
220        type Addr = MailAddr;
221        type Msg = ();
222        type Event = User<MailAddr, ()>;
223        type Sends = Vec<Never>;
224        type Ph = Never;
225        type Error = Never;
226        type Birth = NoBirths;
227
228        fn transition(
229            &mut self,
230            _: crate::ActiveTurn,
231            _: Self::Event,
232        ) -> crate::BehaviorActed<Self> {
233            self.0 += 1;
234            Ok(Actions::cont())
235        }
236    }
237
238    type CountBehavior = Count;
239
240    impl crate::BehaviorBase for Count {
241        type Base = Self;
242
243        fn base(&self) -> &Self {
244            self
245        }
246    }
247
248    #[allow(
249        clippy::unnecessary_wraps,
250        reason = "the reaction fixture must implement the fallible reaction signature"
251    )]
252    fn elapsed(_inner: &mut CountBehavior) -> Acted<MailAddr, Never, Vec<Never>, NoBirths, Never> {
253        Ok(Actions::cont())
254    }
255
256    #[tokio::test]
257    async fn exhaustion_retires_only_the_timer_and_preserves_the_inner_fold() {
258        let mut timeout =
259            ReceiveTimeout::new(Count(0), TimerId(0), Duration::from_secs(1), elapsed);
260        crate::calculus::initialize(&mut timeout).unwrap();
261        timeout.timer = TimerLease::idle(TimerGeneration(u64::MAX));
262
263        let actions = crate::calculus::delegate_transition(
264            &mut timeout,
265            ReceiveTimeoutEvent::Behavior(User::user(MailAddr(1), ())),
266        )
267        .unwrap();
268
269        assert!(actions.sends.schedules.is_empty());
270        assert!(actions.creates.is_empty());
271        assert!(matches!(actions.become_, Step::Continue));
272        assert_eq!(crate::BehaviorBase::base(&timeout).0, 1);
273        assert_eq!(timeout.timer.live(), None);
274    }
275}