behavior/timing/
deadline.rs1use tokio::time::Instant;
5
6use super::domain::OneShotSchedule;
7use super::event::TimedEvent;
8use crate::Step;
9use crate::behavior::{Actions, Address, Become, Behavior, BirthMode, SendAlgebra, ServiceSends};
10use crate::protocol::{ScheduleAt, TimeEvent, TimerId};
11use crate::{Inner, Own, SendInput};
12
13pub type DeadlineEvent<E> = TimedEvent<E>;
14
15pub type DeadlineReaction<B> =
16 fn(&mut B) -> Result<Become<<B as Behavior>::Addr>, <B as Behavior>::Error>;
17
18pub struct DeadlineSends<Sends> {
20 pub behavior: Sends,
21 pub schedules: ServiceSends<ScheduleAt>,
22}
23
24impl<Sends: SendAlgebra> SendAlgebra for DeadlineSends<Sends> {
25 fn empty() -> Self {
26 Self {
27 behavior: Sends::empty(),
28 schedules: ServiceSends::empty(),
29 }
30 }
31
32 fn append(&mut self, other: Self) {
33 self.behavior.append(other.behavior);
34 self.schedules.append(other.schedules);
35 }
36}
37
38impl<Sends> SendInput<ScheduleAt, Own> for DeadlineSends<Sends> {
39 fn emit(&mut self, input: ScheduleAt) {
40 self.schedules.send(input);
41 }
42}
43
44impl<Sends, Input, Path> SendInput<Input, Inner<Path>> for DeadlineSends<Sends>
45where
46 Sends: SendInput<Input, Path>,
47{
48 fn emit(&mut self, input: Input) {
49 <Sends as SendInput<Input, Path>>::emit(&mut self.behavior, input);
50 }
51}
52
53pub type DeadlineActions<B> = Actions<
54 <B as Behavior>::Addr,
55 <B as Behavior>::Ph,
56 DeadlineSends<<B as Behavior>::Sends>,
57 <B as Behavior>::Birth,
58>;
59
60pub struct Deadline<B: Behavior> {
61 inner: B,
62 schedule: OneShotSchedule,
63 on_reached: DeadlineReaction<B>,
64}
65
66impl<B: Behavior> Deadline<B> {
67 #[must_use]
68 pub fn new(
69 inner: B,
70 id: TimerId,
71 at: Option<Instant>,
72 on_reached: DeadlineReaction<B>,
73 ) -> Self {
74 Self {
75 inner,
76 schedule: OneShotSchedule::new(id, at),
77 on_reached,
78 }
79 }
80
81 #[must_use]
82 pub fn inner(&self) -> &B {
83 &self.inner
84 }
85}
86
87impl<B, A, Ph, Sends, Br> Behavior for Deadline<B>
88where
89 A: Address,
90 Sends: SendAlgebra,
91 Br: BirthMode,
92 B: Behavior<Addr = A, Ph = Ph, Sends = Sends, Birth = Br>,
93 B::Event: TimeEvent,
94{
95 type Addr = A;
96 type Msg = B::Msg;
97 type Event = DeadlineEvent<B::Event>;
98 type Sends = DeadlineSends<Sends>;
99 type Ph = Ph;
100 type Error = B::Error;
101 type Birth = Br;
102
103 fn init(&mut self) -> Result<DeadlineActions<B>, B::Error> {
104 let actions = self.inner.init()?;
105 let own = if matches!(actions.become_, Step::Stop(_)) {
106 self.schedule.cancel();
107 ServiceSends::empty()
108 } else {
109 self.schedule
110 .request()
111 .map_or_else(ServiceSends::empty, |(id, generation, at)| {
112 ServiceSends::one(ScheduleAt::new(id, generation, at))
113 })
114 };
115 Ok(Self::wrap(actions, own))
116 }
117
118 fn transition(&mut self, event: Self::Event) -> Result<DeadlineActions<B>, B::Error> {
119 match event {
120 DeadlineEvent::Elapsed(event) if self.schedule.accept(event.id, event.generation) => {
121 let become_ = match (self.on_reached)(&mut self.inner)? {
122 Step::Continue => Step::Continue,
123 Step::Goto(never) => match never {},
124 Step::Stop(exit) => Step::Stop(exit),
125 };
126 Ok(Actions::just(become_))
127 }
128 DeadlineEvent::Elapsed(event) => match B::Event::time_reached(event) {
129 Some(inner) => {
130 let actions = self.inner.transition(inner)?;
131 if matches!(actions.become_, Step::Stop(_)) {
132 self.schedule.cancel();
133 }
134 Ok(Self::wrap(actions, ServiceSends::empty()))
135 }
136 None => Ok(Actions::cont()),
137 },
138 DeadlineEvent::Inner(event) => {
139 let actions = self.inner.transition(event)?;
140 if matches!(actions.become_, Step::Stop(_)) {
141 self.schedule.cancel();
142 }
143 Ok(Self::wrap(actions, ServiceSends::empty()))
144 }
145 }
146 }
147}
148
149impl<B: Behavior> Deadline<B> {
150 fn wrap(
151 actions: Actions<B::Addr, B::Ph, B::Sends, B::Birth>,
152 own: ServiceSends<ScheduleAt>,
153 ) -> DeadlineActions<B> {
154 actions.map_sends(|behavior| DeadlineSends {
155 behavior,
156 schedules: own,
157 })
158 }
159}