Skip to main content

behavior_actors/workflow/
coordinator.rs

1//! Dependency-ordered workflow activation as a pure fold.
2
3#[cfg(test)]
4use behavior::Recipient;
5use behavior::{Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, User};
6use thiserror::Error;
7
8use crate::DeliveryRoute;
9
10/// Bombay-owned immutable dependency-graph product.
11#[derive(Debug, Clone, PartialEq, Eq)]
12pub struct WorkflowDefinition<K> {
13    /// Steps in deterministic definition order.
14    pub steps: Vec<K>,
15    /// Directed prerequisite edges `(prerequisite, dependent)`.
16    pub dependencies: Vec<(K, K)>,
17}
18
19/// Construction rejection for a malformed workflow graph.
20#[derive(Debug, Error, Clone, PartialEq, Eq)]
21pub enum WorkflowConfigError<K> {
22    /// At least one step is required.
23    #[error("workflow requires at least one step")]
24    Empty,
25    /// Step identity occurs more than once.
26    #[error("workflow step identity is duplicated")]
27    DuplicateStep { step: K },
28    /// One edge names an undefined prerequisite.
29    #[error("workflow dependency names an undefined prerequisite")]
30    UnknownPrerequisite { step: K },
31    /// One edge names an undefined dependent.
32    #[error("workflow dependency names an undefined dependent")]
33    UnknownDependent { step: K },
34    /// One edge depends upon itself.
35    #[error("workflow step cannot depend upon itself")]
36    SelfDependency { step: K },
37    /// The dependency relation contains a cycle.
38    #[error("workflow dependencies contain a cycle")]
39    Cycle,
40}
41
42/// Complete lifecycle of one workflow step.
43#[derive(Debug, Clone, Copy, PartialEq, Eq)]
44pub enum WorkflowStepState {
45    /// At least one prerequisite has not completed.
46    Blocked,
47    /// Activation was emitted and completion is pending.
48    Active,
49    /// Successful completion was committed.
50    Completed,
51}
52
53/// Complete workflow phase sum.
54pub enum WorkflowState<K, Route> {
55    /// Validated definition has not started.
56    Ready,
57    /// At least one step remains incomplete.
58    Running {
59        /// Per-step states in definition order.
60        steps: Vec<(K, WorkflowStepState)>,
61        /// Recipient for every activation and terminal fact.
62        reply_to: Route,
63    },
64    /// Every step completed successfully.
65    Succeeded {
66        /// Retained result route for typed rejection of later stale input.
67        reply_to: Route,
68    },
69    /// One active step failed; no further activations are possible.
70    Failed {
71        /// Step whose failure terminated the run.
72        step: K,
73        /// Retained result route for typed rejection of later stale input.
74        reply_to: Route,
75    },
76    /// Explicit cancellation terminated the run.
77    Cancelled {
78        /// Retained result route for typed rejection of later stale input.
79        reply_to: Route,
80    },
81}
82
83/// Rejected workflow input.
84#[derive(Debug, Clone, PartialEq, Eq)]
85pub enum WorkflowRejection<K> {
86    /// Start was requested after the only run had started or terminated.
87    AlreadyStarted,
88    /// Completion names no defined step.
89    UnknownStep { step: K },
90    /// A blocked step cannot complete or fail.
91    Blocked { step: K },
92    /// The step already completed.
93    AlreadyCompleted { step: K },
94    /// The workflow is terminal.
95    Terminal { step: Option<K> },
96}
97
98/// A command that requires an established workflow run.
99#[derive(Debug, Clone, PartialEq, Eq)]
100pub enum WorkflowInput<K> {
101    /// Successful completion evidence for one step.
102    Complete { step: K },
103    /// Failure evidence for one step.
104    Fail { step: K },
105}
106
107/// Typed refusal of a command for which no result route exists yet.
108#[derive(Debug, Error, Clone, PartialEq, Eq)]
109pub enum WorkflowError<K> {
110    /// `Start` has not established a run and its result route.
111    #[error("workflow input arrived before the run was started")]
112    NotStarted(WorkflowInput<K>),
113}
114
115/// Activation and terminal facts emitted by [`Workflow`].
116#[derive(Debug, Clone, PartialEq, Eq)]
117pub enum WorkflowOutcome<K> {
118    /// A run began and these root steps became active in definition order.
119    Started { activated: Vec<K> },
120    /// One completion committed and may have activated dependents.
121    Advanced { completed: K, activated: Vec<K> },
122    /// The final completion committed.
123    Succeeded { completed: K },
124    /// One active step failed and terminated the run.
125    Failed { step: K },
126    /// Explicit cancellation terminated the run.
127    Cancelled,
128    /// Input was rejected without mutation.
129    Rejected(WorkflowRejection<K>),
130}
131
132/// Closed workflow protocol.
133pub enum WorkflowMessage<K, Route> {
134    /// Begin the single workflow run.
135    Start { reply_to: Route },
136    /// Report successful completion of an active step.
137    Complete { step: K },
138    /// Report failure of an active step.
139    Fail { step: K },
140    /// Cancel a ready or running workflow.
141    Cancel { reply_to: Route },
142}
143
144/// Deterministic dependency-ordered workflow coordinator.
145///
146/// Construction validates a finite acyclic graph. `Start` activates every
147/// root in definition order. Completion activates a blocked step exactly once
148/// when all of its prerequisites are complete. Unknown, blocked, duplicate,
149/// and post-start terminal input is explicitly rejected without mutation.
150/// Completion or failure input before `Start` is returned as a typed
151/// [`WorkflowError::NotStarted`] because no result route has been established.
152/// Failure and cancellation are terminal and prevent later activation. Initialization is
153/// empty and the actor itself remains available to report terminal rejection.
154/// Graph validation, one-run retention, activation ordering, and terminal
155/// policy are Bombay choices. The actor model requires only that each event be
156/// processed as one transition; participant execution and durable saga state
157/// belong to the Driver/Mnesis boundaries. No transition panics.
158pub 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    /// Validate and retain one dependency graph.
175    ///
176    /// # Errors
177    ///
178    /// Returns a concrete [`WorkflowConfigError`] for empty, duplicate,
179    /// unknown, self-dependent, or cyclic definitions.
180    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    /// Borrow the validated definition.
190    #[must_use]
191    pub const fn definition(&self) -> &WorkflowDefinition<K> {
192        &self.definition
193    }
194
195    /// Borrow the complete workflow phase.
196    #[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}