Skip to main content

behavior/timing/
deadline.rs

1//! Pure one-shot time composition. Scheduling is a request to the emitting
2//! actor's local clock service.
3
4use 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
18/// Named effect lanes added by [`Deadline`].
19pub 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}