Skip to main content

behavior_actors/time/
deadline.rs

1//! Pure one-shot time composition. Scheduling is a request to the emitting
2//! actor's local clock service.
3
4use std::time::Instant;
5
6use super::domain::{OneShotSchedule, TimerAdmission};
7use super::event::TimedEvent;
8use crate::protocol::{ScheduleAt, TimerId};
9use behavior::Step;
10use behavior::{
11    Actions, Address, Become, Behavior, BirthMode, EventLayer, InterpreterRequests, SendEffects,
12    SendLayer,
13};
14
15/// Infallible reaction to one accepted deadline.
16///
17/// ```compile_fail,E0308
18/// # struct App;
19/// # impl behavior::Protocol for App { type Addr = behavior::MailAddr; type Msg = (); }
20/// # impl behavior::Behavior for App {
21/// #   type Protocol = Self; type Event = behavior::User<behavior::MailAddr, ()>; type Sends = Vec<behavior::Never>;
22/// #   type Ph = behavior::Never; type Error = behavior::Never; type Birth = behavior::NoBirths;
23/// #   fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> behavior::BehaviorActed<Self> { Ok(behavior::Actions::cont()) }
24/// # }
25/// fn fallible(_: &mut App) -> Result<behavior::Become, behavior::Never> {
26///     Ok(behavior::Step::Continue)
27/// }
28/// let _ = behavior_actors::Deadline::new(App, behavior_actors::TimerId(1), None, fallible);
29/// ```
30pub type DeadlineReaction<B> = fn(&mut B) -> Become;
31
32pub(crate) type DeadlineActions<B> = Actions<
33    behavior::BehaviorAddr<B>,
34    <B as Behavior>::Ph,
35    SendLayer<InterpreterRequests<ScheduleAt>, <B as Behavior>::Sends>,
36    <B as Behavior>::Birth,
37>;
38
39pub struct Deadline<B: Behavior> {
40    inner: B,
41    schedule: OneShotSchedule,
42    on_reached: DeadlineReaction<B>,
43}
44
45impl<B: Behavior> Deadline<B> {
46    /// Wrap `inner` with one optional absolute deadline and pure reaction.
47    ///
48    /// Initialization stages the schedule when `at` is present. Clock access
49    /// and timer delivery remain interpreter capabilities.
50    /// Nested timers remain independently addressable even when they reuse the
51    /// same `(id, generation)`: each emitted schedule selects this wrapper's
52    /// structural ingress destination.
53    /// The reaction is infallible because it receives mutable access to the
54    /// wrapped behavior; ordinary delegated transitions retain its error type.
55    #[must_use]
56    pub fn new(
57        inner: B,
58        id: TimerId,
59        at: Option<Instant>,
60        on_reached: DeadlineReaction<B>,
61    ) -> Self {
62        Self {
63            inner,
64            schedule: OneShotSchedule::new(id, at),
65            on_reached,
66        }
67    }
68}
69
70impl<B: Behavior + behavior::BehaviorBase> behavior::BehaviorBase for Deadline<B> {
71    type Base = B::Base;
72
73    fn base(&self) -> &Self::Base {
74        self.inner.base()
75    }
76}
77
78impl<B> crate::StashStatus for Deadline<B>
79where
80    B: Behavior + crate::StashStatus,
81{
82    fn stashed_messages(&self) -> usize {
83        self.inner.stashed_messages()
84    }
85}
86
87impl<B, A, Ph, Sends, Br> Behavior for Deadline<B>
88where
89    A: Address,
90    Sends: SendEffects + behavior::SendsFor<B::Event>,
91    Br: BirthMode,
92    B: Behavior<Ph = Ph, Sends = Sends, Birth = Br>,
93    B::Protocol: behavior::Protocol<Addr = A>,
94{
95    type Protocol = B::Protocol;
96    type Event = TimedEvent<B::Event>;
97    type Sends = SendLayer<InterpreterRequests<ScheduleAt>, Sends>;
98    type Ph = Ph;
99    type Error = B::Error;
100    type Birth = Br;
101
102    fn init(&mut self, _: behavior::InitializationTurn) -> Result<DeadlineActions<B>, B::Error> {
103        let actions = behavior::initialize(&mut self.inner)?;
104        let own = if matches!(actions.become_, Step::Stop(_)) {
105            self.schedule.cancel();
106            InterpreterRequests::empty()
107        } else {
108            self.schedule.request().map_or_else(
109                InterpreterRequests::empty,
110                |(id, generation, at)| {
111                    InterpreterRequests::one(ScheduleAt::new(id, generation, at))
112                },
113            )
114        };
115        Ok(Self::wrap(actions, own))
116    }
117
118    fn transition(
119        &mut self,
120        _: behavior::ActiveTurn,
121        event: Self::Event,
122    ) -> Result<DeadlineActions<B>, B::Error> {
123        match event {
124            EventLayer::Owned(event) => match self.schedule.accept(event.id, event.generation) {
125                TimerAdmission::Accepted => {
126                    let become_ = match (self.on_reached)(&mut self.inner) {
127                        Step::Continue => Step::Continue,
128                        Step::Goto(never) => match never {},
129                        Step::Stop(exit) => Step::Stop(exit),
130                    };
131                    Ok(Actions::just(become_))
132                }
133                TimerAdmission::Ignored => Ok(Actions::cont()),
134            },
135            EventLayer::Inner(event) => {
136                let actions = behavior::delegate_transition(&mut self.inner, event)?;
137                if matches!(actions.become_, Step::Stop(_)) {
138                    self.schedule.cancel();
139                }
140                Ok(Self::wrap(actions, InterpreterRequests::empty()))
141            }
142        }
143    }
144}
145
146impl<B: Behavior> Deadline<B> {
147    fn wrap(
148        actions: Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
149        own: InterpreterRequests<ScheduleAt>,
150    ) -> DeadlineActions<B> {
151        actions.map_sends(|inner| SendLayer::new(own, inner))
152    }
153}