behavior_actors/lifecycle/
task.rs1use behavior::{
4 Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Step, User,
5};
6#[cfg(test)]
7use behavior::{Delivery, Recipient};
8use thiserror::Error;
9
10use crate::DeliveryRoute;
11
12#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
14pub enum TaskState {
15 Pending,
17 Completed,
19 Cancelled,
21}
22
23#[derive(Debug, Clone, PartialEq, Eq)]
25pub enum TaskResult<R> {
26 Completed(R),
28 Cancelled,
30}
31
32pub enum TaskMessage<R, Route> {
34 Complete {
36 result: R,
38 reply_to: Route,
40 },
41 Cancel {
43 reply_to: Route,
45 },
46}
47
48#[derive(Debug, Error, Clone, PartialEq, Eq)]
51pub enum TaskError<R, Route> {
52 #[error("task result arrived after completion")]
54 ResultAfterCompletion { result: R, reply_to: Route },
55 #[error("task result arrived after cancellation")]
57 ResultAfterCancellation { result: R, reply_to: Route },
58 #[error("task cancellation arrived after completion")]
60 CancellationAfterCompletion { reply_to: Route },
61 #[error("task cancellation arrived after cancellation")]
63 CancellationAfterCancellation { reply_to: Route },
64}
65
66pub struct Task<A, R, Route>
80where
81 A: Address,
82 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = TaskResult<R>>>,
83{
84 state: TaskState,
85 protocol: core::marker::PhantomData<fn() -> (A, R, Route)>,
86}
87
88impl<A, R, Route> Task<A, R, Route>
89where
90 A: Address,
91 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = TaskResult<R>>>,
92{
93 #[must_use]
95 pub const fn new() -> Self {
96 Self {
97 state: TaskState::Pending,
98 protocol: core::marker::PhantomData,
99 }
100 }
101
102 #[must_use]
104 pub const fn state(&self) -> TaskState {
105 self.state
106 }
107}
108
109impl<A, R, Route> Default for Task<A, R, Route>
110where
111 A: Address,
112 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = TaskResult<R>>>,
113{
114 fn default() -> Self {
115 Self::new()
116 }
117}
118
119impl<A, R, Route> BehaviorBase for Task<A, R, Route>
120where
121 A: Address,
122 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = TaskResult<R>>>,
123{
124 type Base = Self;
125
126 fn base(&self) -> &Self {
127 self
128 }
129}
130
131impl<A, R, Route> behavior::Protocol for Task<A, R, Route>
132where
133 A: Address,
134 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = TaskResult<R>>>,
135{
136 type Addr = A;
137 type Msg = TaskMessage<R, Route>;
138}
139
140impl<A, R, Route> Behavior for Task<A, R, Route>
141where
142 A: Address,
143 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = TaskResult<R>>>,
144 Route::Sends: behavior::SendsFor<User<A, TaskMessage<R, Route>>>,
145{
146 type Protocol = Self;
147 type Event = User<A, behavior::BehaviorMessage<Self>>;
148 type Sends = Route::Sends;
149 type Ph = Never;
150 type Error = TaskError<R, Route>;
151 type Birth = NoBirths;
152
153 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
154 match (self.state, event.message) {
155 (TaskState::Pending, TaskMessage::Complete { result, reply_to }) => {
156 self.state = TaskState::Completed;
157 Ok(Actions::new(
158 reply_to.deliver(TaskResult::Completed(result)),
159 behavior::Creations::empty(),
160 Step::Stop(behavior::Stopped),
161 ))
162 }
163 (TaskState::Pending, TaskMessage::Cancel { reply_to }) => {
164 self.state = TaskState::Cancelled;
165 Ok(Actions::new(
166 reply_to.deliver(TaskResult::Cancelled),
167 behavior::Creations::empty(),
168 Step::Stop(behavior::Stopped),
169 ))
170 }
171 (TaskState::Completed, TaskMessage::Complete { result, reply_to }) => {
172 Err(TaskError::ResultAfterCompletion { result, reply_to })
173 }
174 (TaskState::Cancelled, TaskMessage::Complete { result, reply_to }) => {
175 Err(TaskError::ResultAfterCancellation { result, reply_to })
176 }
177 (TaskState::Completed, TaskMessage::Cancel { reply_to }) => {
178 Err(TaskError::CancellationAfterCompletion { reply_to })
179 }
180 (TaskState::Cancelled, TaskMessage::Cancel { reply_to }) => {
181 Err(TaskError::CancellationAfterCancellation { reply_to })
182 }
183 }
184 }
185}
186
187#[cfg(test)]
188mod tests {
189 use super::*;
190 use crate::Activate as _;
191 use behavior::MailAddr;
192
193 struct Reply;
194
195 impl behavior::Protocol for Reply {
196 type Addr = MailAddr;
197 type Msg = TaskResult<u8>;
198 }
199
200 impl Behavior for Reply {
201 type Protocol = Self;
202 type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
203 type Sends = Vec<Never>;
204 type Ph = Never;
205 type Error = Never;
206 type Birth = NoBirths;
207
208 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
209 Ok(Actions::cont())
210 }
211 }
212
213 type TestTask = Task<MailAddr, u8, Recipient<Reply>>;
214
215 #[test]
216 fn completion_reports_owned_result_and_stops_atomically() {
217 let reply = Recipient::<Reply>::global(MailAddr(8));
218 let mut task = (TestTask::new()).initialize().unwrap().behavior;
219 let completed = task
220 .receive(
221 MailAddr(9),
222 TaskMessage::Complete {
223 result: 7,
224 reply_to: reply,
225 },
226 )
227 .unwrap();
228 assert!(completed.sends == vec![Delivery::new(reply, TaskResult::Completed(7))]);
229 assert!(completed.creates.is_empty());
230 assert!(matches!(completed.become_, Step::Stop(_)));
231 assert_eq!(task.state(), TaskState::Completed);
232 let rejection = task.receive(
233 MailAddr(9),
234 TaskMessage::Complete {
235 result: 8,
236 reply_to: reply,
237 },
238 );
239 assert!(matches!(
240 rejection,
241 Err(TaskError::ResultAfterCompletion { result: 8, reply_to }) if reply_to == reply
242 ));
243 }
244
245 #[test]
246 fn cancellation_is_a_distinct_terminal_outcome() {
247 let reply = Recipient::<Reply>::global(MailAddr(8));
248 let mut task = (TestTask::new()).initialize().unwrap().behavior;
249 let cancelled = task
250 .receive(MailAddr(9), TaskMessage::Cancel { reply_to: reply })
251 .unwrap();
252 assert!(cancelled.sends == vec![Delivery::new(reply, TaskResult::Cancelled)]);
253 assert!(matches!(cancelled.become_, Step::Stop(_)));
254 assert_eq!(task.state(), TaskState::Cancelled);
255 let rejection = task.receive(
256 MailAddr(9),
257 TaskMessage::Complete {
258 result: 9,
259 reply_to: reply,
260 },
261 );
262 assert!(matches!(
263 rejection,
264 Err(TaskError::ResultAfterCancellation { result: 9, reply_to }) if reply_to == reply
265 ));
266 }
267}