1use 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
34pub 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 pub(super) deploy_mutations: AsyncMutex<()>,
56 visibility_reconciliation_task: Option<JoinHandle<()>>,
57}
58
59pub(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 #[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(®istry_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 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 #[must_use]
155 pub fn store(&self) -> Arc<dyn EventStore> {
156 Arc::clone(&self.store)
157 }
158
159 #[must_use]
161 pub fn visibility_store(&self) -> Arc<dyn VisibilityStore> {
162 Arc::clone(&self.visibility_store)
163 }
164
165 #[must_use]
167 pub fn runtime(&self) -> &RuntimeHandle {
168 &self.runtime
169 }
170
171 #[must_use]
173 pub fn workflow_catalog(&self) -> &Arc<WorkflowCatalog> {
174 &self.catalog
175 }
176
177 #[must_use]
179 pub fn registry(&self) -> &Registry {
180 &self.registry
181 }
182
183 #[must_use]
185 pub fn supervision(&self) -> &SupervisionTree {
186 &self.supervision
187 }
188
189 #[must_use]
191 pub const fn delegated(&self) -> &DelegatedSeams {
192 &self.delegated
193 }
194
195 #[must_use]
197 pub fn signal_handoff(&self) -> Arc<SignalResumeHandoff> {
198 Arc::clone(&self.signal_handoff)
199 }
200
201 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 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 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 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 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 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 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 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 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}