Skip to main content

aion/engine/
api.rs

1//! `Engine` start, cancel, result, list, and shutdown support.
2
3use std::collections::HashMap;
4use std::sync::Arc;
5
6use aion_core::{
7    Event, Payload, RunId, SearchAttributeSchema, SearchAttributeValue, WorkflowError,
8    WorkflowFilter, WorkflowId, WorkflowSummary,
9};
10use tokio::sync::Mutex as AsyncMutex;
11use tokio::task::JoinHandle;
12
13use crate::durability::Recorder;
14use crate::schedule::ScheduleEvaluator;
15use aion_store::EventStore;
16use aion_store::visibility::VisibilityStore;
17
18use crate::lifecycle::continue_as_new::{self, ContinueAsNewContext, ContinueAsNewRequest};
19use crate::lifecycle::start::{self, StartWorkflowContext};
20use crate::lifecycle::terminate::{self, TerminateWorkflowContext};
21use crate::lifecycle::transition;
22use crate::registry::{TerminalOutcome, WorkflowHandle};
23use crate::{
24    EngineError, Registry, RuntimeHandle, SupervisionTree, WorkflowCatalog,
25    signal::SignalResumeHandoff,
26};
27
28use super::api_schedule::{
29    ScheduleRuntimeDeps, default_schedule_evaluator, schedule_coordinator_workflow_id,
30};
31use super::delegated::DelegatedSeams;
32use super::shutdown_gate::ShutdownGate;
33
34/// Live embedded workflow engine assembled by [`crate::EngineBuilder`].
35pub struct Engine {
36    store: Arc<dyn EventStore>,
37    visibility_store: Arc<dyn VisibilityStore>,
38    pub(super) schedule_recorder: Arc<AsyncMutex<Recorder>>,
39    pub(super) schedule_evaluator: Arc<AsyncMutex<ScheduleEvaluator>>,
40    pub(super) schedule_coordinator_workflow_id: WorkflowId,
41    runtime: Arc<RuntimeHandle>,
42    catalog: Arc<WorkflowCatalog>,
43    registry: Arc<Registry>,
44    supervision: Arc<SupervisionTree>,
45    delegated: DelegatedSeams,
46    signal_handoff: Arc<SignalResumeHandoff>,
47    search_attribute_schema: Arc<SearchAttributeSchema>,
48    pub(super) shutdown_gate: ShutdownGate,
49    /// Serializes the deploy mutations (load / route / unload) end-to-end
50    /// across BOTH the catalog commit and its store persistence write, so
51    /// the persisted package set and route pointers can never disagree with
52    /// the catalog through interleaving (for example a concurrent re-deploy
53    /// re-persisting a version an unload just deleted). Workflow dispatch
54    /// never takes this lock.
55    pub(super) deploy_mutations: AsyncMutex<()>,
56    visibility_reconciliation_task: Option<JoinHandle<()>>,
57}
58
59/// Components required to construct an [`Engine`].
60pub(crate) struct EngineComponents {
61    pub(crate) store: Arc<dyn EventStore>,
62    pub(crate) visibility_store: Arc<dyn VisibilityStore>,
63    pub(crate) runtime: Arc<RuntimeHandle>,
64    pub(crate) catalog: Arc<WorkflowCatalog>,
65    pub(crate) registry: Arc<Registry>,
66    pub(crate) supervision: Arc<SupervisionTree>,
67    pub(crate) delegated: DelegatedSeams,
68    pub(crate) signal_handoff: Arc<SignalResumeHandoff>,
69    pub(crate) search_attribute_schema: Arc<SearchAttributeSchema>,
70    pub(crate) visibility_reconciliation_task: Option<JoinHandle<()>>,
71}
72
73impl Engine {
74    /// Construct an engine from already-assembled components.
75    #[must_use]
76    pub(crate) fn new(components: EngineComponents) -> Self {
77        let EngineComponents {
78            store,
79            visibility_store,
80            runtime,
81            catalog,
82            registry,
83            supervision,
84            delegated,
85            signal_handoff,
86            search_attribute_schema,
87            visibility_reconciliation_task,
88        } = components;
89        let schedule_coordinator_workflow_id = schedule_coordinator_workflow_id();
90        let schedule_recorder = Arc::new(AsyncMutex::new(Recorder::new(
91            schedule_coordinator_workflow_id.clone(),
92            Arc::clone(&store),
93        )));
94        let runtime_arc = runtime;
95        let registry_arc = registry;
96        let supervision_arc = supervision;
97        let schedule_evaluator = Arc::new(AsyncMutex::new(default_schedule_evaluator(
98            schedule_coordinator_workflow_id.clone(),
99            Arc::clone(&schedule_recorder),
100            ScheduleRuntimeDeps {
101                store: Arc::clone(&store),
102                visibility_store: Arc::clone(&visibility_store),
103                runtime: Arc::clone(&runtime_arc),
104                catalog: Arc::clone(&catalog),
105                registry: Arc::clone(&registry_arc),
106                supervision: Arc::clone(&supervision_arc),
107                search_attribute_schema: Arc::clone(&search_attribute_schema),
108            },
109        )));
110        Self {
111            store,
112            visibility_store,
113            schedule_recorder,
114            schedule_evaluator,
115            schedule_coordinator_workflow_id,
116            runtime: runtime_arc,
117            catalog,
118            registry: registry_arc,
119            supervision: supervision_arc,
120            delegated,
121            signal_handoff,
122            search_attribute_schema,
123            shutdown_gate: ShutdownGate::default(),
124            deploy_mutations: AsyncMutex::new(()),
125            visibility_reconciliation_task,
126        }
127    }
128
129    /// Advance the schedule coordinator's recorder head to match persisted
130    /// events so that a rebuilt engine resumes appending at the correct
131    /// sequence rather than conflicting at head 0.
132    ///
133    /// # Errors
134    ///
135    /// Returns store read errors.
136    pub(crate) async fn catchup_schedule_coordinator(&self) -> Result<(), EngineError> {
137        let history = self
138            .store
139            .read_history(&self.schedule_coordinator_workflow_id)
140            .await?;
141        let head = u64::try_from(history.len()).unwrap_or(u64::MAX);
142        if head > 0 {
143            let mut recorder = self.schedule_recorder.lock().await;
144            *recorder = Recorder::resume_at(
145                self.schedule_coordinator_workflow_id.clone(),
146                Arc::clone(&self.store),
147                head,
148            );
149        }
150        Ok(())
151    }
152
153    /// Event store used by lifecycle and delegated AD/AT operations.
154    #[must_use]
155    pub fn store(&self) -> Arc<dyn EventStore> {
156        Arc::clone(&self.store)
157    }
158
159    /// Visibility store used for workflow summary projections.
160    #[must_use]
161    pub fn visibility_store(&self) -> Arc<dyn VisibilityStore> {
162        Arc::clone(&self.visibility_store)
163    }
164
165    /// Runtime boundary assembled for this engine.
166    #[must_use]
167    pub fn runtime(&self) -> &RuntimeHandle {
168        &self.runtime
169    }
170
171    /// Shared workflow package catalog: loaded versions and routing.
172    #[must_use]
173    pub fn workflow_catalog(&self) -> &Arc<WorkflowCatalog> {
174        &self.catalog
175    }
176
177    /// Active execution registry.
178    #[must_use]
179    pub fn registry(&self) -> &Registry {
180        &self.registry
181    }
182
183    /// Supervision tree snapshot/model.
184    #[must_use]
185    pub fn supervision(&self) -> &SupervisionTree {
186        &self.supervision
187    }
188
189    /// Delegated signal/query/subscribe seams installed for AT/AD integration.
190    #[must_use]
191    pub const fn delegated(&self) -> &DelegatedSeams {
192        &self.delegated
193    }
194
195    /// Shared in-memory handoff for already-recorded non-resident signals.
196    #[must_use]
197    pub fn signal_handoff(&self) -> Arc<SignalResumeHandoff> {
198        Arc::clone(&self.signal_handoff)
199    }
200
201    /// Start a loaded workflow type as a new BEAM process.
202    ///
203    /// `search_attributes` are validated against the engine's configured
204    /// [`SearchAttributeSchema`] and recorded atomically with the
205    /// `WorkflowStarted` event, so visibility metadata can never be lost to a
206    /// crash between start and a later attribute update.
207    ///
208    /// # Errors
209    ///
210    /// Returns [`EngineError::ShuttingDown`] after shutdown begins, and
211    /// [`EngineError::Durability`] when a search attribute is unregistered or
212    /// mistyped (nothing is appended and no process is spawned). Otherwise
213    /// delegates to the start lifecycle transition and returns its typed errors.
214    pub async fn start_workflow(
215        &self,
216        workflow_type: &str,
217        input: Payload,
218        search_attributes: HashMap<String, SearchAttributeValue>,
219    ) -> Result<WorkflowHandle, EngineError> {
220        let operation = self.shutdown_gate.begin_start()?;
221        let result = start::start_workflow_with_options(
222            StartWorkflowContext {
223                store: self.store(),
224                visibility_store: self.visibility_store(),
225                catalog: Arc::clone(&self.catalog),
226                runtime: Arc::clone(&self.runtime),
227                supervision: Arc::clone(&self.supervision),
228                registry: Arc::clone(&self.registry),
229                signal_handoff: Some(self.signal_handoff()),
230                search_attribute_schema: Arc::clone(&self.search_attribute_schema),
231                monitor_tokio_handle: tokio::runtime::Handle::current(),
232            },
233            workflow_type,
234            input,
235            start::StartWorkflowOptions {
236                search_attributes,
237                ..start::StartWorkflowOptions::default()
238            },
239        )
240        .await;
241        drop(operation);
242        result
243    }
244
245    /// Resume a suspended workflow run and flush deferred signals through its mailbox.
246    ///
247    /// # Errors
248    ///
249    /// Returns [`EngineError::WorkflowNotFound`] when the `(workflow, run)` pair
250    /// is absent, or registry errors from the residency transition. Deferred
251    /// delivery failures are logged and dropped because signals are already durable.
252    pub fn resume_workflow(
253        &self,
254        id: &WorkflowId,
255        run: &RunId,
256    ) -> Result<WorkflowHandle, EngineError> {
257        let handle = transition::resume(self.registry(), id, run)?;
258        if let Err(error) = self.signal_handoff.deliver_deferred(self, id) {
259            tracing::warn!(
260                workflow_id = %id,
261                run_id = %run,
262                error = %error,
263                "failed to flush deferred signals after workflow resume"
264            );
265        }
266        Ok(handle)
267    }
268
269    /// Cancel a live workflow run by killing its runtime process.
270    ///
271    /// # Errors
272    ///
273    /// Returns [`EngineError::ShuttingDown`] after shutdown begins, and
274    /// [`EngineError::WorkflowNotFound`] when the `(workflow, run)` pair
275    /// is not live. Other typed errors come from the cancel transition.
276    pub async fn cancel(
277        &self,
278        id: &WorkflowId,
279        run: &RunId,
280        reason: impl Into<String>,
281    ) -> Result<(), EngineError> {
282        let operation = self.shutdown_gate.begin_operation()?;
283        let result = terminate::cancel(
284            TerminateWorkflowContext {
285                runtime: &self.runtime,
286                store: self.store(),
287                visibility_store: self.visibility_store(),
288                registry: &self.registry,
289            },
290            id,
291            run,
292            reason,
293        )
294        .await;
295        drop(operation);
296        result
297    }
298
299    /// Continue a live workflow run as a new run under the same workflow id.
300    ///
301    /// # Errors
302    ///
303    /// Returns [`EngineError::ShuttingDown`] after shutdown begins, and
304    /// [`EngineError::WorkflowNotFound`] when the `(workflow, run)` pair
305    /// is not live. Other typed errors come from the continue-as-new transition.
306    pub async fn continue_as_new(
307        &self,
308        id: &WorkflowId,
309        run: &RunId,
310        input: Payload,
311        workflow_type: Option<String>,
312    ) -> Result<WorkflowHandle, EngineError> {
313        let operation = self.shutdown_gate.begin_operation()?;
314        let result = continue_as_new::continue_as_new(
315            ContinueAsNewContext {
316                store: self.store(),
317                visibility_store: Arc::clone(&self.visibility_store),
318                catalog: Arc::clone(&self.catalog),
319                runtime: &self.runtime,
320                supervision: Arc::clone(&self.supervision),
321                registry: &self.registry,
322                search_attribute_schema: Arc::clone(&self.search_attribute_schema),
323            },
324            id,
325            run,
326            ContinueAsNewRequest {
327                input,
328                workflow_type,
329            },
330        )
331        .await;
332        drop(operation);
333        result
334    }
335
336    /// Await a workflow run's terminal result.
337    ///
338    /// Already-terminal histories return immediately. Live workflows await their
339    /// completion notifier. Unknown workflow/run pairs return not found.
340    ///
341    /// # Errors
342    ///
343    /// Returns store, registry, or runtime channel errors as typed [`EngineError`]
344    /// variants, or [`EngineError::WorkflowNotFound`] when no live handle or
345    /// terminal history exists for the requested pair.
346    pub async fn result(
347        &self,
348        id: &WorkflowId,
349        run: &RunId,
350    ) -> Result<Result<Payload, WorkflowError>, EngineError> {
351        let history = self.store.read_history(id).await?;
352        if let Some(outcome) = terminal_outcome_from_history(&history) {
353            return Ok(outcome_to_result(outcome));
354        }
355
356        let handle = match self.registry.get(id, run)? {
357            Some(handle) => handle,
358            // Registration birth window: the run is durably started but its
359            // handle insert has not landed yet (see
360            // `Engine::handle_after_birth_window`).
361            None => self
362                .handle_after_birth_window(id, run, &history)
363                .await?
364                .ok_or_else(|| workflow_not_found(id, run))?,
365        };
366        let mut receiver = handle.completion().subscribe();
367        loop {
368            if let Some(outcome) = receiver.borrow().clone() {
369                return Ok(outcome_to_result(outcome));
370            }
371            if receiver.changed().await.is_err() {
372                if let Some(outcome) =
373                    terminal_outcome_from_history(&self.store.read_history(id).await?)
374                {
375                    return Ok(outcome_to_result(outcome));
376                }
377                return Err(EngineError::Runtime {
378                    reason: format!(
379                        "completion channel closed before workflow `{id}/{run}` finished"
380                    ),
381                });
382            }
383        }
384    }
385
386    /// List live and terminal workflow summaries matching `filter`.
387    ///
388    /// Store projections are authoritative; live registry entries are projected
389    /// from durable history before being merged and deduplicated.
390    ///
391    /// # Errors
392    ///
393    /// Returns typed store or registry errors when visibility data cannot be read.
394    pub async fn list_workflows(
395        &self,
396        filter: WorkflowFilter,
397    ) -> Result<Vec<WorkflowSummary>, EngineError> {
398        let mut summaries = self
399            .store
400            .query(&filter)
401            .await?
402            .into_iter()
403            .map(|summary| (summary.workflow_id.clone(), summary))
404            .collect::<HashMap<_, _>>();
405
406        for handle in self.registry.list()? {
407            let history = self.store.read_history(handle.workflow_id()).await?;
408            self.registry
409                .reconcile(handle.workflow_id(), handle.run_id(), &history)?;
410            if let Some(summary) = WorkflowSummary::from_history(&history) {
411                if filter.matches(&summary) {
412                    summaries.insert(summary.workflow_id.clone(), summary);
413                }
414            }
415        }
416
417        let mut summaries = summaries.into_values().collect::<Vec<_>>();
418        summaries.sort_by(|left, right| {
419            left.started_at.cmp(&right.started_at).then_with(|| {
420                left.workflow_id
421                    .to_string()
422                    .cmp(&right.workflow_id.to_string())
423            })
424        });
425        Ok(summaries)
426    }
427
428    /// Gracefully stop accepting new starts and shut down the embedded runtime.
429    ///
430    /// # Errors
431    ///
432    /// Returns registry poison or runtime shutdown failures as typed errors.
433    pub fn shutdown(&self) -> Result<(), EngineError> {
434        if let Some(task) = &self.visibility_reconciliation_task {
435            task.abort();
436        }
437        self.shutdown_gate.close_and_wait()?;
438        // Epoch close for engine-side child tasks (F4): the scheduler stops
439        // first (so no NIF can arm a new watcher mid-shutdown), then every
440        // watcher and spawn-recovery task is aborted AND awaited to
441        // quiescence — a task still mid-record after shutdown could
442        // double-write a parent history a successor engine over the same
443        // store also records into. Arming is additionally gated inside the
444        // task registry the moment shutdown begins.
445        self.runtime.shutdown()?;
446        self.runtime.nif_state().shutdown_child_tasks();
447        Ok(())
448    }
449}
450
451pub(crate) fn terminal_outcome_from_history(events: &[Event]) -> Option<TerminalOutcome> {
452    for event in events.iter().rev() {
453        match event {
454            Event::WorkflowStarted { .. } => return None,
455            Event::WorkflowCompleted { result, .. } => {
456                return Some(TerminalOutcome::Completed(result.clone()));
457            }
458            Event::WorkflowFailed { error, .. } => {
459                return Some(TerminalOutcome::Failed(error.clone()));
460            }
461            Event::WorkflowCancelled { reason, .. } => {
462                return Some(TerminalOutcome::Cancelled(reason.clone()));
463            }
464            Event::WorkflowTimedOut { timeout, .. } => {
465                return Some(TerminalOutcome::TimedOut(timeout.clone()));
466            }
467            Event::WorkflowContinuedAsNew {
468                input,
469                workflow_type,
470                parent_run_id,
471                ..
472            } => {
473                return Some(TerminalOutcome::ContinuedAsNew {
474                    input: input.clone(),
475                    workflow_type: workflow_type.clone(),
476                    parent_run_id: parent_run_id.clone(),
477                });
478            }
479            Event::SearchAttributesUpdated { .. }
480            | Event::ActivityScheduled { .. }
481            | Event::ActivityStarted { .. }
482            | Event::ActivityCompleted { .. }
483            | Event::ActivityFailed { .. }
484            | Event::ActivityCancelled { .. }
485            | Event::TimerStarted { .. }
486            | Event::TimerFired { .. }
487            | Event::TimerCancelled { .. }
488            | Event::WithTimeoutCompleted { .. }
489            | Event::SignalReceived { .. }
490            | Event::SignalSent { .. }
491            | Event::ChildWorkflowStarted { .. }
492            | Event::ChildWorkflowCompleted { .. }
493            | Event::ChildWorkflowFailed { .. }
494            | Event::ChildWorkflowCancelled { .. }
495            | Event::ScheduleCreated { .. }
496            | Event::ScheduleUpdated { .. }
497            | Event::SchedulePaused { .. }
498            | Event::ScheduleResumed { .. }
499            | Event::ScheduleDeleted { .. }
500            | Event::ScheduleTriggered { .. } => {}
501        }
502    }
503    None
504}
505
506fn outcome_to_result(outcome: TerminalOutcome) -> Result<Payload, WorkflowError> {
507    match outcome {
508        TerminalOutcome::Completed(payload) => Ok(payload),
509        TerminalOutcome::Failed(error) => Err(error),
510        TerminalOutcome::Cancelled(reason) => Err(WorkflowError {
511            message: format!("workflow cancelled: {reason}"),
512            details: None,
513        }),
514        TerminalOutcome::TimedOut(timeout) => Err(WorkflowError {
515            message: format!("workflow timed out: {timeout}"),
516            details: None,
517        }),
518        TerminalOutcome::ContinuedAsNew { parent_run_id, .. } => Err(WorkflowError {
519            message: format!("workflow continued as new from run {parent_run_id}"),
520            details: None,
521        }),
522    }
523}
524
525pub(crate) fn workflow_not_found(id: &WorkflowId, run: &RunId) -> EngineError {
526    EngineError::WorkflowNotFound {
527        workflow_type: format!("{id}/{run}"),
528    }
529}
530
531#[cfg(test)]
532mod tests {
533    use std::collections::HashMap;
534    use std::sync::Arc;
535
536    use aion_core::{Event, Payload, SearchAttributeSchema, WorkflowFilter, WorkflowStatus};
537    use aion_package::ContentHash;
538    use aion_store::visibility::VisibilityStore;
539    use aion_store::{EventStore, InMemoryStore};
540    use serde_json::json;
541
542    use super::{DelegatedSeams, Engine, EngineComponents};
543    use crate::durability::Recorder;
544    use crate::lifecycle::terminate::{self, TerminateWorkflowContext};
545    use crate::registry::{CompletionNotifier, HandleResidency, WorkflowHandleParts};
546    use crate::{
547        EngineError, Registry, RuntimeConfig, RuntimeHandle, SupervisionTree, WorkflowCatalog,
548        WorkflowHandle,
549    };
550
551    fn payload(label: &str) -> Result<Payload, aion_core::PayloadError> {
552        Payload::from_json(&json!({ "label": label }))
553    }
554
555    fn workflow_error(message: &str) -> aion_core::WorkflowError {
556        aion_core::WorkflowError {
557            message: message.to_owned(),
558            details: None,
559        }
560    }
561
562    fn workflow_catalog(workflow_type: &str, deployed_module: &str) -> Arc<WorkflowCatalog> {
563        let catalog = Arc::new(WorkflowCatalog::new());
564        catalog.note_loaded_workflow_for_test(
565            workflow_type,
566            deployed_module,
567            "run",
568            ContentHash::from_bytes([5; 32]),
569        );
570        catalog
571    }
572
573    fn engine_with_loaded_workflow(
574        store: Arc<dyn EventStore>,
575        workflow_type: &str,
576        deployed_module: &str,
577    ) -> Result<Engine, EngineError> {
578        let runtime = RuntimeHandle::new(RuntimeConfig::new(Some(1)))?;
579        runtime.register_waiting_test_module(deployed_module, "run");
580        let visibility_store: Arc<dyn VisibilityStore> = Arc::new(InMemoryStore::default());
581        Ok(Engine::new(EngineComponents {
582            store,
583            visibility_store,
584            runtime: Arc::new(runtime),
585            catalog: workflow_catalog(workflow_type, deployed_module),
586            registry: Arc::new(Registry::default()),
587            supervision: Arc::new(SupervisionTree::new()),
588            delegated: DelegatedSeams::default(),
589            signal_handoff: Arc::new(crate::signal::SignalResumeHandoff::new()),
590            search_attribute_schema: Arc::new(SearchAttributeSchema::new()),
591            visibility_reconciliation_task: None,
592        }))
593    }
594
595    fn termination_context(engine: &Engine) -> TerminateWorkflowContext<'_> {
596        TerminateWorkflowContext {
597            runtime: engine.runtime(),
598            store: engine.store(),
599            visibility_store: engine.visibility_store(),
600            registry: engine.registry(),
601        }
602    }
603
604    async fn insert_active_handle(
605        engine: &Engine,
606        store: Arc<dyn EventStore>,
607        workflow_type: &str,
608    ) -> Result<WorkflowHandle, Box<dyn std::error::Error>> {
609        let workflow_id = aion_core::WorkflowId::new_v4();
610        let run_id = aion_core::RunId::new_v4();
611        let mut recorder = Recorder::new(workflow_id.clone(), store);
612        recorder
613            .record_workflow_started(
614                chrono::Utc::now(),
615                crate::durability::WorkflowStartRecord {
616                    workflow_type: workflow_type.to_owned(),
617                    input: payload("input")?,
618                    run_id: run_id.clone(),
619                    parent_run_id: None,
620                    package_version: aion_core::PackageVersion::new("a".repeat(64)),
621                },
622            )
623            .await?;
624        let pid = engine.runtime().spawn_test_process_with_trap_exit(true)?;
625        let handle = WorkflowHandle::new(WorkflowHandleParts {
626            workflow_id: workflow_id.clone(),
627            run_id: run_id.clone(),
628            pid,
629            workflow_type: workflow_type.to_owned(),
630            loaded_version: ContentHash::from_bytes([9; 32]),
631            cached_status: WorkflowStatus::Running,
632            residency: HandleResidency::Resident,
633            recorder,
634            completion: CompletionNotifier::new(),
635        });
636        engine
637            .registry()
638            .insert((workflow_id, run_id), handle.clone())?;
639        Ok(handle)
640    }
641
642    #[tokio::test]
643    async fn start_then_cancel_records_started_then_cancelled()
644    -> Result<(), Box<dyn std::error::Error>> {
645        let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
646        let engine =
647            engine_with_loaded_workflow(Arc::clone(&store), "checkout", "checkout_deployed")?;
648        let handle = engine
649            .start_workflow("checkout", payload("input")?, HashMap::new())
650            .await?;
651
652        engine
653            .cancel(
654                handle.workflow_id(),
655                handle.run_id(),
656                "caller requested cancellation",
657            )
658            .await?;
659
660        let history = store.read_history(handle.workflow_id()).await?;
661        match history.as_slice() {
662            [
663                Event::WorkflowStarted { .. },
664                Event::WorkflowCancelled { reason, .. },
665            ] => {
666                assert_eq!(reason, "caller requested cancellation");
667            }
668            other => return Err(format!("expected started then cancelled, found {other:?}").into()),
669        }
670        engine.shutdown()?;
671        Ok(())
672    }
673
674    #[tokio::test]
675    async fn result_returns_completed_payload() -> Result<(), Box<dyn std::error::Error>> {
676        let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
677        let engine =
678            engine_with_loaded_workflow(Arc::clone(&store), "checkout", "checkout_deployed")?;
679        let handle = engine
680            .start_workflow("checkout", payload("input")?, HashMap::new())
681            .await?;
682        let result_payload = payload("result")?;
683
684        terminate::complete(
685            termination_context(&engine),
686            handle.workflow_id(),
687            handle.run_id(),
688            result_payload.clone(),
689        )
690        .await?;
691
692        assert_eq!(
693            engine.result(handle.workflow_id(), handle.run_id()).await?,
694            Ok(result_payload)
695        );
696        engine.shutdown()?;
697        Ok(())
698    }
699
700    #[tokio::test]
701    async fn result_returns_failed_workflow_error() -> Result<(), Box<dyn std::error::Error>> {
702        let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
703        let engine =
704            engine_with_loaded_workflow(Arc::clone(&store), "checkout", "checkout_deployed")?;
705        let handle = engine
706            .start_workflow("checkout", payload("input")?, HashMap::new())
707            .await?;
708        let error = workflow_error("workflow failed");
709
710        terminate::fail(
711            termination_context(&engine),
712            handle.workflow_id(),
713            handle.run_id(),
714            error.clone(),
715        )
716        .await?;
717
718        assert_eq!(
719            engine.result(handle.workflow_id(), handle.run_id()).await?,
720            Err(error)
721        );
722        engine.shutdown()?;
723        Ok(())
724    }
725
726    #[tokio::test]
727    async fn result_unknown_workflow_returns_not_found() -> Result<(), Box<dyn std::error::Error>> {
728        let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
729        let engine = engine_with_loaded_workflow(store, "checkout", "checkout_deployed")?;
730        let workflow_id = aion_core::WorkflowId::new_v4();
731        let run_id = aion_core::RunId::new_v4();
732
733        let result = engine.result(&workflow_id, &run_id).await;
734
735        assert!(matches!(result, Err(EngineError::WorkflowNotFound { .. })));
736        engine.shutdown()?;
737        Ok(())
738    }
739
740    #[tokio::test]
741    async fn continue_as_new_unknown_workflow_returns_not_found()
742    -> Result<(), Box<dyn std::error::Error>> {
743        let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
744        let engine = engine_with_loaded_workflow(store, "checkout", "checkout_deployed")?;
745        let workflow_id = aion_core::WorkflowId::new_v4();
746        let run_id = aion_core::RunId::new_v4();
747
748        let result = engine
749            .continue_as_new(&workflow_id, &run_id, payload("next")?, None)
750            .await;
751
752        assert!(matches!(result, Err(EngineError::WorkflowNotFound { .. })));
753        engine.shutdown()?;
754        Ok(())
755    }
756
757    #[tokio::test]
758    async fn list_workflows_merges_live_and_terminal_without_duplicates()
759    -> Result<(), Box<dyn std::error::Error>> {
760        let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
761        let engine =
762            engine_with_loaded_workflow(Arc::clone(&store), "checkout", "checkout_deployed")?;
763        let running = insert_active_handle(&engine, Arc::clone(&store), "checkout").await?;
764        let completed = engine
765            .start_workflow("checkout", payload("input")?, HashMap::new())
766            .await?;
767        terminate::complete(
768            termination_context(&engine),
769            completed.workflow_id(),
770            completed.run_id(),
771            payload("result")?,
772        )
773        .await?;
774
775        let summaries = engine.list_workflows(WorkflowFilter::default()).await?;
776        assert_eq!(summaries.len(), 2);
777        assert!(summaries.iter().any(|summary| {
778            &summary.workflow_id == running.workflow_id()
779                && summary.status == WorkflowStatus::Running
780        }));
781        assert!(summaries.iter().any(|summary| {
782            &summary.workflow_id == completed.workflow_id()
783                && summary.status == WorkflowStatus::Completed
784        }));
785
786        let completed_only = engine
787            .list_workflows(WorkflowFilter {
788                status: Some(WorkflowStatus::Completed),
789                ..WorkflowFilter::default()
790            })
791            .await?;
792        assert_eq!(completed_only.len(), 1);
793        assert_eq!(&completed_only[0].workflow_id, completed.workflow_id());
794        engine.shutdown()?;
795        Ok(())
796    }
797
798    #[tokio::test]
799    async fn shutdown_rejects_subsequent_starts() -> Result<(), Box<dyn std::error::Error>> {
800        let store: Arc<dyn EventStore> = Arc::new(InMemoryStore::default());
801        let engine =
802            engine_with_loaded_workflow(Arc::clone(&store), "checkout", "checkout_deployed")?;
803        let handle = engine
804            .start_workflow("checkout", payload("input")?, HashMap::new())
805            .await?;
806        terminate::complete(
807            termination_context(&engine),
808            handle.workflow_id(),
809            handle.run_id(),
810            payload("result")?,
811        )
812        .await?;
813
814        engine.shutdown()?;
815        let result = engine
816            .start_workflow("checkout", payload("after-shutdown")?, HashMap::new())
817            .await;
818
819        assert!(matches!(result, Err(EngineError::ShuttingDown)));
820        Ok(())
821    }
822}