1#[cfg(test)]
4use behavior::Recipient;
5use behavior::{Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, User};
6use thiserror::Error;
7
8use crate::DeliveryRoute;
9
10#[derive(Debug, Clone, PartialEq, Eq)]
12pub struct WorkflowDefinition<K> {
13 pub steps: Vec<K>,
15 pub dependencies: Vec<(K, K)>,
17}
18
19#[derive(Debug, Error, Clone, PartialEq, Eq)]
21pub enum WorkflowConfigError<K> {
22 #[error("workflow requires at least one step")]
24 Empty,
25 #[error("workflow step identity is duplicated")]
27 DuplicateStep { step: K },
28 #[error("workflow dependency names an undefined prerequisite")]
30 UnknownPrerequisite { step: K },
31 #[error("workflow dependency names an undefined dependent")]
33 UnknownDependent { step: K },
34 #[error("workflow step cannot depend upon itself")]
36 SelfDependency { step: K },
37 #[error("workflow dependencies contain a cycle")]
39 Cycle,
40}
41
42#[derive(Debug, Clone, Copy, PartialEq, Eq)]
44pub enum WorkflowStepState {
45 Blocked,
47 Active,
49 Completed,
51}
52
53pub enum WorkflowState<K, Route> {
55 Ready,
57 Running {
59 steps: Vec<(K, WorkflowStepState)>,
61 reply_to: Route,
63 },
64 Succeeded {
66 reply_to: Route,
68 },
69 Failed {
71 step: K,
73 reply_to: Route,
75 },
76 Cancelled {
78 reply_to: Route,
80 },
81}
82
83#[derive(Debug, Clone, PartialEq, Eq)]
85pub enum WorkflowRejection<K> {
86 AlreadyStarted,
88 UnknownStep { step: K },
90 Blocked { step: K },
92 AlreadyCompleted { step: K },
94 Terminal { step: Option<K> },
96}
97
98#[derive(Debug, Clone, PartialEq, Eq)]
100pub enum WorkflowInput<K> {
101 Complete { step: K },
103 Fail { step: K },
105}
106
107#[derive(Debug, Error, Clone, PartialEq, Eq)]
109pub enum WorkflowError<K> {
110 #[error("workflow input arrived before the run was started")]
112 NotStarted(WorkflowInput<K>),
113}
114
115#[derive(Debug, Clone, PartialEq, Eq)]
117pub enum WorkflowOutcome<K> {
118 Started { activated: Vec<K> },
120 Advanced { completed: K, activated: Vec<K> },
122 Succeeded { completed: K },
124 Failed { step: K },
126 Cancelled,
128 Rejected(WorkflowRejection<K>),
130}
131
132pub enum WorkflowMessage<K, Route> {
134 Start { reply_to: Route },
136 Complete { step: K },
138 Fail { step: K },
140 Cancel { reply_to: Route },
142}
143
144pub struct Workflow<
159 A: Address,
160 K,
161 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = WorkflowOutcome<K>>>,
162> {
163 definition: WorkflowDefinition<K>,
164 state: WorkflowState<K, Route>,
165 marker: core::marker::PhantomData<fn() -> A>,
166}
167
168impl<A, K, Route> Workflow<A, K, Route>
169where
170 A: Address,
171 K: Clone + Eq,
172 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = WorkflowOutcome<K>>> + Clone,
173{
174 pub fn new(definition: WorkflowDefinition<K>) -> Result<Self, WorkflowConfigError<K>> {
181 validate(&definition)?;
182 Ok(Self {
183 definition,
184 state: WorkflowState::Ready,
185 marker: core::marker::PhantomData,
186 })
187 }
188
189 #[must_use]
191 pub const fn definition(&self) -> &WorkflowDefinition<K> {
192 &self.definition
193 }
194
195 #[must_use]
197 pub const fn state(&self) -> &WorkflowState<K, Route> {
198 &self.state
199 }
200
201 fn reply(
202 reply_to: Route,
203 outcome: WorkflowOutcome<K>,
204 ) -> Actions<A, Never, Route::Sends, NoBirths> {
205 Actions::send(reply_to.deliver(outcome))
206 }
207
208 fn start(&mut self, reply_to: Route) -> Actions<A, Never, Route::Sends, NoBirths> {
209 if !matches!(self.state, WorkflowState::Ready) {
210 return Self::reply(
211 reply_to,
212 WorkflowOutcome::Rejected(WorkflowRejection::AlreadyStarted),
213 );
214 }
215 let mut steps: Vec<_> = self
216 .definition
217 .steps
218 .iter()
219 .cloned()
220 .map(|step| (step, WorkflowStepState::Blocked))
221 .collect();
222 let mut activated = Vec::new();
223 for (step, state) in &mut steps {
224 if !self
225 .definition
226 .dependencies
227 .iter()
228 .any(|(_, dependent)| dependent == step)
229 {
230 *state = WorkflowStepState::Active;
231 activated.push(step.clone());
232 }
233 }
234 self.state = WorkflowState::Running {
235 steps,
236 reply_to: reply_to.clone(),
237 };
238 Self::reply(reply_to, WorkflowOutcome::Started { activated })
239 }
240
241 fn complete(&mut self, step: K) -> Actions<A, Never, Route::Sends, NoBirths> {
242 let terminal_reply = match &self.state {
243 WorkflowState::Succeeded { reply_to }
244 | WorkflowState::Failed { reply_to, .. }
245 | WorkflowState::Cancelled { reply_to } => Some(reply_to.clone()),
246 WorkflowState::Ready | WorkflowState::Running { .. } => None,
247 };
248 let WorkflowState::Running { steps, reply_to } = &mut self.state else {
249 return terminal_reply.map_or_else(Actions::cont, |reply_to| {
250 Self::reply(
251 reply_to,
252 WorkflowOutcome::Rejected(WorkflowRejection::Terminal { step: Some(step) }),
253 )
254 });
255 };
256 let reply_to = reply_to.clone();
257 let Some(index) = steps.iter().position(|(candidate, _)| candidate == &step) else {
258 return Self::reply(
259 reply_to,
260 WorkflowOutcome::Rejected(WorkflowRejection::UnknownStep { step }),
261 );
262 };
263 match steps[index].1 {
264 WorkflowStepState::Blocked => {
265 return Self::reply(
266 reply_to,
267 WorkflowOutcome::Rejected(WorkflowRejection::Blocked { step }),
268 );
269 }
270 WorkflowStepState::Completed => {
271 return Self::reply(
272 reply_to,
273 WorkflowOutcome::Rejected(WorkflowRejection::AlreadyCompleted { step }),
274 );
275 }
276 WorkflowStepState::Active => steps[index].1 = WorkflowStepState::Completed,
277 }
278 if steps
279 .iter()
280 .all(|(_, state)| *state == WorkflowStepState::Completed)
281 {
282 self.state = WorkflowState::Succeeded {
283 reply_to: reply_to.clone(),
284 };
285 return Self::reply(reply_to, WorkflowOutcome::Succeeded { completed: step });
286 }
287 let completed: Vec<_> = steps
288 .iter()
289 .filter(|(_, state)| *state == WorkflowStepState::Completed)
290 .map(|(key, _)| key.clone())
291 .collect();
292 let mut activated = Vec::new();
293 for (candidate, state) in steps
294 .iter_mut()
295 .filter(|(_, state)| *state == WorkflowStepState::Blocked)
296 {
297 let ready = self
298 .definition
299 .dependencies
300 .iter()
301 .filter(|(_, dependent)| dependent == candidate)
302 .all(|(required, _)| completed.contains(required));
303 if ready {
304 *state = WorkflowStepState::Active;
305 activated.push(candidate.clone());
306 }
307 }
308 Self::reply(
309 reply_to,
310 WorkflowOutcome::Advanced {
311 completed: step,
312 activated,
313 },
314 )
315 }
316
317 fn fail(&mut self, step: K) -> Actions<A, Never, Route::Sends, NoBirths> {
318 let terminal_reply = match &self.state {
319 WorkflowState::Succeeded { reply_to }
320 | WorkflowState::Failed { reply_to, .. }
321 | WorkflowState::Cancelled { reply_to } => Some(reply_to.clone()),
322 WorkflowState::Ready | WorkflowState::Running { .. } => None,
323 };
324 let WorkflowState::Running { steps, reply_to } = &self.state else {
325 return terminal_reply.map_or_else(Actions::cont, |reply_to| {
326 Self::reply(
327 reply_to,
328 WorkflowOutcome::Rejected(WorkflowRejection::Terminal { step: Some(step) }),
329 )
330 });
331 };
332 let reply_to = reply_to.clone();
333 let Some((_, phase)) = steps.iter().find(|(candidate, _)| candidate == &step) else {
334 return Self::reply(
335 reply_to,
336 WorkflowOutcome::Rejected(WorkflowRejection::UnknownStep { step }),
337 );
338 };
339 match phase {
340 WorkflowStepState::Active => {
341 self.state = WorkflowState::Failed {
342 step: step.clone(),
343 reply_to: reply_to.clone(),
344 };
345 Self::reply(reply_to, WorkflowOutcome::Failed { step })
346 }
347 WorkflowStepState::Blocked => Self::reply(
348 reply_to,
349 WorkflowOutcome::Rejected(WorkflowRejection::Blocked { step }),
350 ),
351 WorkflowStepState::Completed => Self::reply(
352 reply_to,
353 WorkflowOutcome::Rejected(WorkflowRejection::AlreadyCompleted { step }),
354 ),
355 }
356 }
357}
358
359fn validate<K: Clone + Eq>(
360 definition: &WorkflowDefinition<K>,
361) -> Result<(), WorkflowConfigError<K>> {
362 if definition.steps.is_empty() {
363 return Err(WorkflowConfigError::Empty);
364 }
365 for (index, step) in definition.steps.iter().enumerate() {
366 if definition.steps[..index].contains(step) {
367 return Err(WorkflowConfigError::DuplicateStep { step: step.clone() });
368 }
369 }
370 for (required, dependent) in &definition.dependencies {
371 if !definition.steps.contains(required) {
372 return Err(WorkflowConfigError::UnknownPrerequisite {
373 step: required.clone(),
374 });
375 }
376 if !definition.steps.contains(dependent) {
377 return Err(WorkflowConfigError::UnknownDependent {
378 step: dependent.clone(),
379 });
380 }
381 if required == dependent {
382 return Err(WorkflowConfigError::SelfDependency {
383 step: required.clone(),
384 });
385 }
386 }
387 let mut reached = Vec::new();
388 while reached.len() < definition.steps.len() {
389 let before = reached.len();
390 for step in &definition.steps {
391 if !reached.contains(step)
392 && definition
393 .dependencies
394 .iter()
395 .filter(|(_, dependent)| dependent == step)
396 .all(|(required, _)| reached.contains(required))
397 {
398 reached.push(step.clone());
399 }
400 }
401 if reached.len() == before {
402 return Err(WorkflowConfigError::Cycle);
403 }
404 }
405 Ok(())
406}
407
408impl<A, K, Route> BehaviorBase for Workflow<A, K, Route>
409where
410 A: Address,
411 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = WorkflowOutcome<K>>>,
412{
413 type Base = Self;
414 fn base(&self) -> &Self {
415 self
416 }
417}
418
419impl<A, K, Route> behavior::Protocol for Workflow<A, K, Route>
420where
421 A: Address,
422 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = WorkflowOutcome<K>>>,
423{
424 type Addr = A;
425 type Msg = WorkflowMessage<K, Route>;
426}
427
428impl<A, K, Route> Behavior for Workflow<A, K, Route>
429where
430 A: Address,
431 K: Clone + Eq,
432 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = WorkflowOutcome<K>>> + Clone,
433 Route::Sends: behavior::SendsFor<User<A, WorkflowMessage<K, Route>>>,
434{
435 type Protocol = Self;
436 type Event = User<A, behavior::BehaviorMessage<Self>>;
437 type Sends = Route::Sends;
438 type Ph = Never;
439 type Error = WorkflowError<K>;
440 type Birth = NoBirths;
441 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
442 Ok(match event.message {
443 WorkflowMessage::Start { reply_to } => self.start(reply_to),
444 WorkflowMessage::Complete { step } => {
445 if matches!(self.state, WorkflowState::Ready) {
446 return Err(WorkflowError::NotStarted(WorkflowInput::Complete { step }));
447 }
448 self.complete(step)
449 }
450 WorkflowMessage::Fail { step } => {
451 if matches!(self.state, WorkflowState::Ready) {
452 return Err(WorkflowError::NotStarted(WorkflowInput::Fail { step }));
453 }
454 self.fail(step)
455 }
456 WorkflowMessage::Cancel { reply_to } => match &self.state {
457 WorkflowState::Ready | WorkflowState::Running { .. } => {
458 self.state = WorkflowState::Cancelled {
459 reply_to: reply_to.clone(),
460 };
461 Self::reply(reply_to, WorkflowOutcome::Cancelled)
462 }
463 _ => Self::reply(
464 reply_to,
465 WorkflowOutcome::Rejected(WorkflowRejection::Terminal { step: None }),
466 ),
467 },
468 })
469 }
470}
471
472#[cfg(test)]
473mod tests {
474 use super::*;
475 use crate::Activate as _;
476 use behavior::MailAddr;
477
478 struct Reply;
479 impl behavior::Protocol for Reply {
480 type Addr = MailAddr;
481 type Msg = WorkflowOutcome<&'static str>;
482 }
483
484 impl Behavior for Reply {
485 type Protocol = Self;
486 type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
487 type Sends = Vec<Never>;
488 type Ph = Never;
489 type Error = Never;
490 type Birth = NoBirths;
491 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
492 Ok(Actions::cont())
493 }
494 }
495 type Subject = Workflow<MailAddr, &'static str, Recipient<Reply>>;
496 fn reply() -> Recipient<Reply> {
497 Recipient::global(MailAddr(9))
498 }
499 fn diamond() -> WorkflowDefinition<&'static str> {
500 WorkflowDefinition {
501 steps: vec!["root", "left", "right", "join"],
502 dependencies: vec![
503 ("root", "left"),
504 ("root", "right"),
505 ("left", "join"),
506 ("right", "join"),
507 ],
508 }
509 }
510
511 #[test]
512 fn validates_empty_duplicate_unknown_and_cyclic_graphs() {
513 assert!(matches!(
514 Subject::new(WorkflowDefinition {
515 steps: vec![],
516 dependencies: vec![]
517 }),
518 Err(WorkflowConfigError::Empty)
519 ));
520 assert!(matches!(
521 Subject::new(WorkflowDefinition {
522 steps: vec!["a", "a"],
523 dependencies: vec![]
524 }),
525 Err(WorkflowConfigError::DuplicateStep { .. })
526 ));
527 assert!(matches!(
528 Subject::new(WorkflowDefinition {
529 steps: vec!["a"],
530 dependencies: vec![("missing", "a")]
531 }),
532 Err(WorkflowConfigError::UnknownPrerequisite { .. })
533 ));
534 assert!(matches!(
535 Subject::new(WorkflowDefinition {
536 steps: vec!["a", "b"],
537 dependencies: vec![("a", "b"), ("b", "a")]
538 }),
539 Err(WorkflowConfigError::Cycle)
540 ));
541 }
542
543 #[test]
544 fn pre_start_step_inputs_are_returned_with_their_operation_and_do_not_start_a_run() {
545 let mut subject = Subject::new(diamond())
546 .unwrap()
547 .initialize()
548 .unwrap()
549 .behavior;
550 let rejection = subject.receive(MailAddr(0), WorkflowMessage::Complete { step: "root" });
551 assert!(matches!(
552 rejection,
553 Err(WorkflowError::NotStarted(WorkflowInput::Complete {
554 step: "root"
555 }))
556 ));
557 let rejection = subject.receive(MailAddr(0), WorkflowMessage::Fail { step: "left" });
558 assert!(matches!(
559 rejection,
560 Err(WorkflowError::NotStarted(WorkflowInput::Fail {
561 step: "left"
562 }))
563 ));
564 assert!(matches!(subject.state(), WorkflowState::Ready));
565 }
566
567 #[test]
568 fn diamond_activates_each_step_once_after_all_prerequisites() {
569 let mut subject = (Subject::new(diamond()).unwrap())
570 .initialize()
571 .unwrap()
572 .behavior;
573 let started = subject
574 .receive(MailAddr(0), WorkflowMessage::Start { reply_to: reply() })
575 .unwrap();
576 assert!(
577 matches!(&started.sends[0].message, WorkflowOutcome::Started { activated } if activated == &["root"])
578 );
579 let root = subject
580 .receive(MailAddr(0), WorkflowMessage::Complete { step: "root" })
581 .unwrap();
582 assert!(
583 matches!(&root.sends[0].message, WorkflowOutcome::Advanced { activated, .. } if activated == &["left", "right"])
584 );
585 let left = subject
586 .receive(MailAddr(0), WorkflowMessage::Complete { step: "left" })
587 .unwrap();
588 assert!(
589 matches!(&left.sends[0].message, WorkflowOutcome::Advanced { activated, .. } if activated.is_empty())
590 );
591 let right = subject
592 .receive(MailAddr(0), WorkflowMessage::Complete { step: "right" })
593 .unwrap();
594 assert!(
595 matches!(&right.sends[0].message, WorkflowOutcome::Advanced { activated, .. } if activated == &["join"])
596 );
597 let finished = subject
598 .receive(MailAddr(0), WorkflowMessage::Complete { step: "join" })
599 .unwrap();
600 assert!(matches!(
601 finished.sends[0].message,
602 WorkflowOutcome::Succeeded { completed: "join" }
603 ));
604 }
605
606 #[test]
607 fn blocked_failure_and_duplicate_completion_are_atomic() {
608 let mut subject = (Subject::new(diamond()).unwrap())
609 .initialize()
610 .unwrap()
611 .behavior;
612 let started = subject
613 .receive(MailAddr(0), WorkflowMessage::Start { reply_to: reply() })
614 .unwrap();
615 assert_eq!(started.sends.len(), 1);
616 assert!(matches!(
617 started.sends[0].message,
618 WorkflowOutcome::Started { .. }
619 ));
620 assert!(started.creates.is_empty());
621 assert_eq!(started.become_, behavior::Step::Continue);
622 let blocked = subject
623 .receive(MailAddr(0), WorkflowMessage::Fail { step: "join" })
624 .unwrap();
625 assert!(matches!(
626 blocked.sends[0].message,
627 WorkflowOutcome::Rejected(WorkflowRejection::Blocked { step: "join" })
628 ));
629 let advanced = subject
630 .receive(MailAddr(0), WorkflowMessage::Complete { step: "root" })
631 .unwrap();
632 assert_eq!(advanced.sends.len(), 1);
633 assert!(matches!(
634 advanced.sends[0].message,
635 WorkflowOutcome::Advanced {
636 completed: "root",
637 ..
638 }
639 ));
640 assert!(advanced.creates.is_empty());
641 assert_eq!(advanced.become_, behavior::Step::Continue);
642 let duplicate = subject
643 .receive(MailAddr(0), WorkflowMessage::Complete { step: "root" })
644 .unwrap();
645 assert!(matches!(
646 duplicate.sends[0].message,
647 WorkflowOutcome::Rejected(WorkflowRejection::AlreadyCompleted { step: "root" })
648 ));
649 }
650}