Skip to main content

behavior/
deadlined.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 crate::Step;
7use crate::behavior::{
8    Actions, Address, Become, Behavior, BirthMode, SendAlgebra, SendProduct, ServiceSends, User,
9    UserEvent,
10};
11use crate::protocol::{
12    ChildEvent, ChildStopped, PeerEvent, PeerStopped, ScheduleAt, ShutdownEvent, ShutdownRequested,
13    TimeEvent, TimerElapsed, TimerGeneration, TimerId, WorkerEvent, WorkerStopped,
14};
15
16#[derive(Debug, Clone, PartialEq, Eq)]
17pub enum AtEvent<E> {
18    Inner(E),
19    Reached(TimerElapsed),
20}
21
22impl<E> TimeEvent for AtEvent<E> {
23    fn time_reached(event: TimerElapsed) -> Option<Self> {
24        Some(Self::Reached(event))
25    }
26}
27
28impl<E: UserEvent> UserEvent for AtEvent<E> {
29    type Addr = E::Addr;
30    type Message = E::Message;
31
32    fn user(from: Self::Addr, message: Self::Message) -> Self {
33        Self::Inner(E::user(from, message))
34    }
35
36    fn into_user(self) -> Result<User<Self::Addr, Self::Message>, Self> {
37        match self {
38            Self::Inner(event) => event.into_user().map_err(Self::Inner),
39            reached @ Self::Reached(_) => Err(reached),
40        }
41    }
42}
43
44impl<E: PeerEvent<A>, A: Address> PeerEvent<A> for AtEvent<E> {
45    fn peer_stopped(event: PeerStopped<A>) -> Option<Self> {
46        E::peer_stopped(event).map(Self::Inner)
47    }
48}
49
50impl<E: ChildEvent<A>, A: Address> ChildEvent<A> for AtEvent<E> {
51    fn child_stopped(event: ChildStopped<A>) -> Option<Self> {
52        E::child_stopped(event).map(Self::Inner)
53    }
54}
55
56impl<E: WorkerEvent<A>, A: Address> WorkerEvent<A> for AtEvent<E> {
57    fn worker_stopped(event: WorkerStopped<A>) -> Option<Self> {
58        E::worker_stopped(event).map(Self::Inner)
59    }
60}
61
62impl<E: ShutdownEvent> ShutdownEvent for AtEvent<E> {
63    fn shutdown_requested(event: ShutdownRequested) -> Option<Self> {
64        E::shutdown_requested(event).map(Self::Inner)
65    }
66}
67
68pub type AtReaction<B> =
69    fn(&mut B) -> Result<Become<<B as Behavior>::Addr>, <B as Behavior>::Error>;
70
71pub type AtSends<B> = SendProduct<<B as Behavior>::Sends, ServiceSends<ScheduleAt>>;
72
73pub type AtActions<B> =
74    Actions<<B as Behavior>::Addr, <B as Behavior>::Ph, AtSends<B>, <B as Behavior>::Birth>;
75
76pub struct At<B: Behavior> {
77    inner: B,
78    id: TimerId,
79    scheduled: Option<(TimerGeneration, Instant)>,
80    on_reached: AtReaction<B>,
81}
82
83impl<B: Behavior> At<B> {
84    #[must_use]
85    pub fn new(inner: B, id: TimerId, at: Option<Instant>, on_reached: AtReaction<B>) -> Self {
86        Self {
87            inner,
88            id,
89            scheduled: at.map(|at| (TimerGeneration(0), at)),
90            on_reached,
91        }
92    }
93
94    #[must_use]
95    pub fn inner(&self) -> &B {
96        &self.inner
97    }
98}
99
100impl<B, A, Ph, Sends, Br> Behavior for At<B>
101where
102    A: Address + Send,
103    Sends: SendAlgebra,
104    Br: BirthMode,
105    B: Behavior<Addr = A, Ph = Ph, Sends = Sends, Birth = Br> + Send,
106    B::Event: TimeEvent + Send,
107    B::Msg: Send,
108{
109    type Addr = A;
110    type Msg = B::Msg;
111    type Event = AtEvent<B::Event>;
112    type Sends = SendProduct<Sends, ServiceSends<ScheduleAt>>;
113    type Ph = Ph;
114    type Error = B::Error;
115    type Birth = Br;
116
117    async fn init(&mut self) -> Result<AtActions<B>, B::Error> {
118        let actions = self.inner.init().await?;
119        let own = if matches!(actions.become_, Step::Stop(_)) {
120            self.scheduled = None;
121            ServiceSends::empty()
122        } else {
123            self.scheduled
124                .map_or_else(ServiceSends::empty, |(generation, at)| {
125                    ServiceSends::one(ScheduleAt {
126                        id: self.id,
127                        generation,
128                        at,
129                    })
130                })
131        };
132        Ok(Self::wrap(actions, own))
133    }
134
135    async fn step(&mut self, event: Self::Event) -> Result<AtActions<B>, B::Error> {
136        match event {
137            AtEvent::Reached(event)
138                if event.id == self.id
139                    && self.scheduled.map(|(generation, _)| generation)
140                        == Some(event.generation) =>
141            {
142                self.scheduled = None;
143                let become_ = match (self.on_reached)(&mut self.inner)? {
144                    Step::Continue => Step::Continue,
145                    Step::Goto(never) => match never {},
146                    Step::Stop(exit) => Step::Stop(exit),
147                };
148                Ok(Actions {
149                    sends: SendProduct {
150                        inner: B::Sends::empty(),
151                        own: ServiceSends::empty(),
152                    },
153                    creates: Vec::new(),
154                    become_,
155                })
156            }
157            AtEvent::Reached(event) => match B::Event::time_reached(event) {
158                Some(inner) => {
159                    let actions = self.inner.step(inner).await?;
160                    if matches!(actions.become_, Step::Stop(_)) {
161                        self.scheduled = None;
162                    }
163                    Ok(Self::wrap(actions, ServiceSends::empty()))
164                }
165                None => Ok(Actions::cont()),
166            },
167            AtEvent::Inner(event) => {
168                let actions = self.inner.step(event).await?;
169                if matches!(actions.become_, Step::Stop(_)) {
170                    self.scheduled = None;
171                }
172                Ok(Self::wrap(actions, ServiceSends::empty()))
173            }
174        }
175    }
176}
177
178impl<B: Behavior> At<B> {
179    fn wrap(
180        actions: Actions<B::Addr, B::Ph, B::Sends, B::Birth>,
181        own: ServiceSends<ScheduleAt>,
182    ) -> AtActions<B> {
183        Actions {
184            sends: SendProduct {
185                inner: actions.sends,
186                own,
187            },
188            creates: actions.creates,
189            become_: actions.become_,
190        }
191    }
192}