Skip to main content

behavior_actors/time/
periodic.rs

1//! Pure generation-safe periodic timer composition.
2
3use std::time::Duration;
4
5use super::domain::{TimerAdmission, TimerLease};
6use super::event::{TimedEvent, TimedReaction};
7use crate::protocol::{ScheduleAfter, TimerId};
8use behavior::Step;
9use behavior::{
10    Actions, Address, Behavior, BehaviorActed, BirthMode, EventLayer, InterpreterRequests,
11    SendEffects, SendLayer,
12};
13
14/// Infallible fold invoked for each accepted periodic generation.
15///
16/// ```compile_fail,E0308
17/// # struct App;
18/// # impl behavior::Protocol for App { type Addr = behavior::MailAddr; type Msg = (); }
19/// # impl behavior::Behavior for App {
20/// #   type Protocol = Self; type Event = behavior::User<behavior::MailAddr, ()>; type Sends = Vec<behavior::Never>;
21/// #   type Ph = behavior::Never; type Error = behavior::Never; type Birth = behavior::NoBirths;
22/// #   fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> behavior::BehaviorActed<Self> { Ok(behavior::Actions::cont()) }
23/// # }
24/// fn fallible(_: &mut App) -> behavior::BehaviorActed<App> { Ok(behavior::Actions::cont()) }
25/// let _ = behavior_actors::Periodic::new(App, behavior_actors::TimerId(1), std::time::Duration::from_secs(1), fallible);
26/// ```
27/// Repeatedly notify a wrapped behavior at a relative interval.
28///
29/// Initialization preserves inner effects and appends generation zero. Each
30/// matching timer event is consumed once, folds `on_elapsed`, and—only when
31/// that fold continues—appends the next generation. Stale, duplicate, and
32/// wrong-ID observations are inert unless the inner protocol accepts them.
33/// Termination emits no replacement. Generation exhaustion leaves the
34/// timer retired while preserving the inner continuation. These are Bombay
35/// policies; clock access and sleeping remain `bombay-timers` capabilities.
36/// Reactions are infallible because they receive mutable access to the wrapped
37/// behavior; ordinary delegated transitions retain the wrapped error type.
38pub struct Periodic<B: Behavior> {
39    inner: B,
40    id: TimerId,
41    every: Duration,
42    lease: TimerLease,
43    on_elapsed: TimedReaction<B>,
44}
45
46impl<B: Behavior> Periodic<B> {
47    /// Wrap `inner` with a relative periodic timer and pure reaction.
48    ///
49    /// Accepted timer generations are rearmed only after a continuing
50    /// reaction. Clock access and scheduling remain interpreter capabilities.
51    #[must_use]
52    pub fn new(inner: B, id: TimerId, every: Duration, on_elapsed: TimedReaction<B>) -> Self {
53        Self {
54            inner,
55            id,
56            every,
57            lease: TimerLease::new(),
58            on_elapsed,
59        }
60    }
61
62    fn schedule(&mut self) -> InterpreterRequests<ScheduleAfter> {
63        self.lease
64            .arm()
65            .map_or_else(InterpreterRequests::empty, |generation| {
66                InterpreterRequests::one(ScheduleAfter::new(self.id, generation, self.every))
67            })
68    }
69
70    fn wrap(
71        actions: Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
72        schedules: InterpreterRequests<ScheduleAfter>,
73    ) -> Actions<
74        behavior::BehaviorAddr<B>,
75        B::Ph,
76        SendLayer<InterpreterRequests<ScheduleAfter>, B::Sends>,
77        B::Birth,
78    > {
79        actions.map_sends(|inner| SendLayer::new(schedules, inner))
80    }
81
82    fn wrap_and_rearm(
83        &mut self,
84        actions: Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
85    ) -> Actions<
86        behavior::BehaviorAddr<B>,
87        B::Ph,
88        SendLayer<InterpreterRequests<ScheduleAfter>, B::Sends>,
89        B::Birth,
90    > {
91        let schedules = if matches!(actions.become_, Step::Stop(_)) {
92            self.lease.disarm();
93            InterpreterRequests::empty()
94        } else {
95            self.schedule()
96        };
97        Self::wrap(actions, schedules)
98    }
99}
100
101impl<B: Behavior + behavior::BehaviorBase> behavior::BehaviorBase for Periodic<B> {
102    type Base = B::Base;
103
104    fn base(&self) -> &Self::Base {
105        self.inner.base()
106    }
107}
108
109impl<B, A, Ph, Sends, Br> Behavior for Periodic<B>
110where
111    A: Address,
112    Sends: SendEffects + behavior::SendsFor<B::Event>,
113    Br: BirthMode,
114    B: Behavior<Ph = Ph, Sends = Sends, Birth = Br>,
115    B::Protocol: behavior::Protocol<Addr = A>,
116{
117    type Protocol = B::Protocol;
118    type Event = TimedEvent<B::Event>;
119    type Sends = SendLayer<InterpreterRequests<ScheduleAfter>, Sends>;
120    type Ph = Ph;
121    type Error = B::Error;
122    type Birth = Br;
123
124    fn init(&mut self, _: behavior::InitializationTurn) -> BehaviorActed<Self> {
125        let actions = behavior::initialize(&mut self.inner)?;
126        Ok(self.wrap_and_rearm(actions))
127    }
128
129    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
130        match event {
131            EventLayer::Owned(elapsed) if elapsed.id == self.id => {
132                match self.lease.accept(elapsed.generation) {
133                    TimerAdmission::Accepted => {
134                        let actions = (self.on_elapsed)(&mut self.inner);
135                        Ok(self.wrap_and_rearm(actions))
136                    }
137                    TimerAdmission::Ignored => Ok(Actions::cont()),
138                }
139            }
140            EventLayer::Owned(_) => Ok(Actions::cont()),
141            EventLayer::Inner(event) => {
142                let actions = behavior::delegate_transition(&mut self.inner, event)?;
143                if matches!(actions.become_, Step::Stop(_)) {
144                    self.lease.disarm();
145                }
146                Ok(Self::wrap(actions, InterpreterRequests::empty()))
147            }
148        }
149    }
150}
151
152#[cfg(test)]
153mod tests {
154    use super::*;
155    use crate::{Activate as _, TimerElapsed};
156    use behavior::{MailAddr, Never, NoBirths, Step, User};
157
158    struct Probe(usize);
159
160    impl behavior::BehaviorBase for Probe {
161        type Base = Self;
162
163        fn base(&self) -> &Self {
164            self
165        }
166    }
167
168    impl behavior::Protocol for Probe {
169        type Addr = MailAddr;
170        type Msg = ();
171    }
172
173    impl Behavior for Probe {
174        type Protocol = Self;
175        type Event = User<MailAddr, ()>;
176        type Sends = Vec<Never>;
177        type Ph = Never;
178        type Error = Never;
179        type Birth = NoBirths;
180
181        fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
182            Ok(Actions::cont())
183        }
184    }
185
186    fn tick(probe: &mut Probe) -> Actions<MailAddr, Never, Vec<Never>, NoBirths> {
187        probe.0 += 1;
188        if probe.0 == 2 {
189            Actions::stop()
190        } else {
191            Actions::cont()
192        }
193    }
194
195    #[test]
196    fn accepted_ticks_rearm_until_the_reaction_stops() {
197        let every = Duration::from_secs(3);
198        let initialized = crate::Periodic::new(Probe(0), TimerId(2), every, tick)
199            .initialize()
200            .unwrap();
201        assert_eq!(
202            initialized.actions.sends.owned.as_slice(),
203            [ScheduleAfter::new(
204                TimerId(2),
205                crate::TimerGeneration(0),
206                every
207            )]
208        );
209        let mut active = initialized.behavior;
210
211        let first = active
212            .on_path(TimerElapsed::new(TimerId(2), crate::TimerGeneration(0)))
213            .unwrap();
214        assert_eq!(
215            first.sends.owned.as_slice(),
216            [ScheduleAfter::new(
217                TimerId(2),
218                crate::TimerGeneration(1),
219                every
220            )]
221        );
222        assert!(matches!(first.become_, Step::Continue));
223        assert_eq!(active.base().0, 1);
224
225        let duplicate = active
226            .on_path(TimerElapsed::new(TimerId(2), crate::TimerGeneration(0)))
227            .unwrap();
228        assert!(duplicate.sends.owned.is_empty());
229        assert_eq!(active.base().0, 1);
230
231        let stopped = active
232            .on_path(TimerElapsed::new(TimerId(2), crate::TimerGeneration(1)))
233            .unwrap();
234        assert!(stopped.sends.owned.is_empty());
235        assert!(matches!(stopped.become_, Step::Stop(_)));
236        assert_eq!(active.base().0, 2);
237    }
238}