1use 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
14pub 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 #[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}