1use aion_core::{Event, Payload, RunId, WorkflowId, current_lease_terminal, run_segment};
4use async_trait::async_trait;
5use futures::stream::{self, BoxStream};
6use serde::{Deserialize, Serialize};
7use std::sync::Arc;
8
9use crate::{Engine, EngineError, SignalRouterError, WorkflowHandle};
10
11use super::api::workflow_not_found;
12
13#[derive(Serialize, Deserialize, Clone, Debug, Default, PartialEq, Eq)]
19pub struct EventFilter {
20 pub workflow_id: Option<WorkflowId>,
22 pub run: Option<RunId>,
24 pub family: Option<EventFamily>,
26}
27
28impl EventFilter {
29 #[must_use]
35 pub fn matches(&self, event: &Event) -> bool {
36 self.workflow_id
37 .as_ref()
38 .is_none_or(|workflow_id| event.workflow_id() == workflow_id)
39 && self
40 .family
41 .is_none_or(|family| family == event_family(event))
42 }
43}
44
45#[derive(Serialize, Deserialize, Copy, Clone, Debug, PartialEq, Eq)]
47pub enum EventFamily {
48 Workflow,
50 Activity,
52 Timer,
54 Signal,
56 ChildWorkflow,
58 Schedule,
60 Workloop,
63}
64
65#[async_trait]
71pub trait SignalRouter: Send + Sync {
72 async fn route(
74 &self,
75 target: &WorkflowHandle,
76 name: String,
77 payload: Payload,
78 ) -> Result<(), EngineError>;
79}
80
81#[async_trait]
87pub trait QueryService: Send + Sync {
88 async fn query(
96 &self,
97 target: &WorkflowHandle,
98 name: String,
99 arguments: Payload,
100 ) -> Result<Payload, EngineError>;
101}
102
103#[derive(thiserror::Error, Clone, Copy, Debug, PartialEq, Eq)]
109#[error("event subscription lagged behind the live stream and skipped {skipped} events")]
110pub struct EventStreamLagged {
111 pub skipped: u64,
113}
114
115pub trait EventPublisher: Send + Sync {
120 fn subscribe(
127 &self,
128 filter: EventFilter,
129 ) -> BoxStream<'static, Result<Event, EventStreamLagged>>;
130}
131
132#[derive(Clone)]
134pub struct DelegatedSeams {
135 signal_router: Arc<dyn SignalRouter>,
136 query_service: Arc<dyn QueryService>,
137 event_publisher: Arc<dyn EventPublisher>,
138}
139
140impl DelegatedSeams {
141 #[must_use]
143 pub const fn new(
144 signal_router: Arc<dyn SignalRouter>,
145 query_service: Arc<dyn QueryService>,
146 event_publisher: Arc<dyn EventPublisher>,
147 ) -> Self {
148 Self {
149 signal_router,
150 query_service,
151 event_publisher,
152 }
153 }
154
155 #[must_use]
157 pub fn signal_router(&self) -> &dyn SignalRouter {
158 self.signal_router.as_ref()
159 }
160
161 #[must_use]
163 pub fn query_service(&self) -> &dyn QueryService {
164 self.query_service.as_ref()
165 }
166
167 #[must_use]
169 pub fn event_publisher(&self) -> &dyn EventPublisher {
170 self.event_publisher.as_ref()
171 }
172
173 pub(crate) fn signal_router_arc(&self) -> Arc<dyn SignalRouter> {
174 Arc::clone(&self.signal_router)
175 }
176
177 pub(crate) fn query_service_arc(&self) -> Arc<dyn QueryService> {
178 Arc::clone(&self.query_service)
179 }
180
181 pub(crate) fn event_publisher_arc(&self) -> Arc<dyn EventPublisher> {
182 Arc::clone(&self.event_publisher)
183 }
184}
185
186impl Default for DelegatedSeams {
187 fn default() -> Self {
188 Self::new(
189 Arc::new(DeferredSignalRouter),
190 Arc::new(DeferredQueryService),
191 Arc::new(DeferredEventPublisher),
192 )
193 }
194}
195
196#[derive(Debug, Default)]
198pub struct DeferredSignalRouter;
199
200#[async_trait]
201impl SignalRouter for DeferredSignalRouter {
202 async fn route(
203 &self,
204 target: &WorkflowHandle,
205 name: String,
206 payload: Payload,
207 ) -> Result<(), EngineError> {
208 let _ = (target, name, payload);
209 Err(EngineError::Runtime {
210 reason: "signal routing seam is not configured".to_owned(),
211 })
212 }
213}
214
215#[derive(Debug, Default)]
217pub struct DeferredQueryService;
218
219#[async_trait]
220impl QueryService for DeferredQueryService {
221 async fn query(
222 &self,
223 target: &WorkflowHandle,
224 name: String,
225 arguments: Payload,
226 ) -> Result<Payload, EngineError> {
227 let _ = (target, name, arguments);
228 Err(EngineError::Runtime {
229 reason: "query service seam is not configured".to_owned(),
230 })
231 }
232}
233
234#[derive(Debug, Default)]
236pub struct DeferredEventPublisher;
237
238impl EventPublisher for DeferredEventPublisher {
239 fn subscribe(
240 &self,
241 filter: EventFilter,
242 ) -> BoxStream<'static, Result<Event, EventStreamLagged>> {
243 let _ = filter;
244 Box::pin(stream::empty())
245 }
246}
247
248impl Engine {
249 pub async fn signal(
267 &self,
268 id: &WorkflowId,
269 run: &RunId,
270 name: impl Into<String>,
271 payload: Payload,
272 ) -> Result<(), EngineError> {
273 let name = name.into();
274 let handle =
275 if let Some(handle) = self.registry().get(id, run)? {
276 crate::signal::admission::admit_against_handle(
279 self.workflow_catalog(),
280 &handle,
281 &name,
282 &payload,
283 )?;
284 handle
285 } else {
286 let history = self.store().read_history(id).await?;
287 if run_has_terminal_history(&history, run) {
288 return Err(SignalRouterError::Terminal {
289 workflow_id: id.clone(),
290 run_id: run.clone(),
291 }
292 .into());
293 }
294 crate::signal::admission::admit_against_history(
300 self.workflow_catalog(),
301 id,
302 run,
303 &history,
304 &name,
305 &payload,
306 )?;
307 let segment = aion_core::run_segment(&history, run);
314 if matches!(
315 aion_core::status_from_events(segment),
316 aion_core::WorkflowStatus::Paused
317 ) {
318 let head = history.last().map(Event::seq).unwrap_or_default();
319 let mut recorder =
320 crate::durability::Recorder::resume_at(id.clone(), self.store(), head)
321 .with_visibility(run.clone(), self.visibility_store());
322 recorder
323 .record_signal_received(chrono::Utc::now(), name, payload)
324 .await?;
325 return Ok(());
326 }
327 if let Some(workloop) = &self.workloop
335 && let Some(record) = workloop.store.get_workloop(id).await?
336 {
337 let head = history.last().map(Event::seq).unwrap_or_default();
338 let mut recorder =
339 crate::durability::Recorder::resume_at(id.clone(), self.store(), head)
340 .with_visibility(run.clone(), self.visibility_store());
341 recorder
342 .record_signal_received(chrono::Utc::now(), name.clone(), payload)
343 .await?;
344 if record.spec.arming().signals().contains(&name) {
345 workloop.service.wake_now(id).await.map_err(|error| {
346 EngineError::Runtime {
347 reason: format!("signal-armed workloop wake failed: {error}"),
348 }
349 })?;
350 }
351 return Ok(());
352 }
353 self.handle_after_birth_window(id, run, &history)
354 .await?
355 .ok_or_else(|| workflow_not_found(id, run))?
356 };
357 self.delegated()
358 .signal_router()
359 .route(&handle, name, payload)
360 .await
361 }
362
363 pub async fn query(
383 &self,
384 id: &WorkflowId,
385 run: &RunId,
386 name: impl Into<String>,
387 arguments: Payload,
388 ) -> Result<Payload, EngineError> {
389 let handle = if let Some(handle) = self.registry().get(id, run)? {
390 handle
391 } else {
392 let history = self.store().read_history(id).await?;
395 if run_has_terminal_history(&history, run) {
396 return Err(EngineError::Query(crate::query::QueryError::NotRunning(
397 id.clone(),
398 )));
399 }
400 match self.handle_after_birth_window(id, run, &history).await? {
401 Some(handle) => handle,
402 None if run_segment(&history, run).is_empty() => {
411 return Err(workflow_not_found(id, run));
412 }
413 None => {
414 return Err(EngineError::Query(crate::query::QueryError::NotRunning(
415 id.clone(),
416 )));
417 }
418 }
419 };
420 self.delegated()
421 .query_service()
422 .query(&handle, name.into(), arguments)
423 .await
424 }
425
426 pub(crate) async fn handle_after_birth_window(
437 &self,
438 id: &WorkflowId,
439 run: &RunId,
440 history: &[Event],
441 ) -> Result<Option<WorkflowHandle>, EngineError> {
442 let started = history
443 .iter()
444 .any(|event| matches!(event, Event::WorkflowStarted { run_id, .. } if run_id == run));
445 if !started {
446 return Ok(None);
447 }
448 wait_for_registered_handle(self.registry(), id, run, self.runtime().signal_delivery()).await
449 }
450
451 #[must_use]
456 pub fn subscribe(
457 &self,
458 filter: EventFilter,
459 ) -> BoxStream<'static, Result<Event, EventStreamLagged>> {
460 self.delegated().event_publisher().subscribe(filter)
461 }
462}
463
464pub(crate) fn run_has_terminal_history(history: &[Event], run: &RunId) -> bool {
465 current_lease_terminal(run_segment(history, run)).is_some()
470}
471
472pub(crate) async fn wait_for_registered_handle(
482 registry: &crate::registry::Registry,
483 id: &WorkflowId,
484 run: &RunId,
485 policy: crate::runtime::SignalDeliveryConfig,
486) -> Result<Option<WorkflowHandle>, EngineError> {
487 let budget = policy
488 .ready_timeout
489 .saturating_mul(policy.max_enqueue_attempts.max(1));
490 let deadline = std::time::Instant::now() + budget;
491 let mut backoff = policy.initial_backoff;
492 loop {
493 if let Some(handle) = registry.get(id, run)? {
494 return Ok(Some(handle));
495 }
496 if std::time::Instant::now() >= deadline {
497 return Ok(None);
498 }
499 tokio::time::sleep(backoff).await;
500 let doubled = backoff.saturating_mul(2);
501 backoff = if doubled > policy.max_backoff {
502 policy.max_backoff
503 } else {
504 doubled
505 };
506 }
507}
508
509const fn event_family(event: &Event) -> EventFamily {
510 match event {
511 Event::WorkflowStarted { .. }
512 | Event::WorkflowCompleted { .. }
513 | Event::WorkflowFailed { .. }
514 | Event::WorkflowCancelled { .. }
515 | Event::WorkflowTimedOut { .. }
516 | Event::WorkflowContinuedAsNew { .. }
517 | Event::WorkflowReopened { .. }
518 | Event::WorkflowPaused { .. }
519 | Event::WorkflowResumed { .. }
520 | Event::SearchAttributesUpdated { .. }
521 | Event::WorkflowHatched { .. } => EventFamily::Workflow,
524 Event::ActivityScheduled { .. }
525 | Event::ActivityStarted { .. }
526 | Event::ActivityLeased { .. }
527 | Event::ActivityAdoptionOffered { .. }
528 | Event::ActivityCompleted { .. }
529 | Event::ActivityFailed { .. }
530 | Event::ActivityAdvisoryExhausted { .. }
531 | Event::ActivityFallbackRouted { .. }
532 | Event::ActivityCancelled { .. } => EventFamily::Activity,
533 Event::TimerStarted { .. }
534 | Event::TimerFired { .. }
535 | Event::TimerCancelled { .. }
536 | Event::WithTimeoutCompleted { .. } => EventFamily::Timer,
537 Event::SignalReceived { .. } | Event::SignalSent { .. } => EventFamily::Signal,
538 Event::ChildWorkflowStarted { .. }
539 | Event::ChildWorkflowCompleted { .. }
540 | Event::ChildWorkflowFailed { .. }
541 | Event::ChildWorkflowCancelled { .. } => EventFamily::ChildWorkflow,
542 Event::ScheduleCreated { .. }
543 | Event::ScheduleUpdated { .. }
544 | Event::SchedulePaused { .. }
545 | Event::ScheduleResumed { .. }
546 | Event::ScheduleDeleted { .. }
547 | Event::ScheduleTriggered { .. } => EventFamily::Schedule,
548 Event::CadenceFired { .. }
549 | Event::IterationClosed { .. }
550 | Event::LoopRetired { .. }
551 | Event::InvariantUnconfirmed { .. } => EventFamily::Workloop,
552 }
553}
554
555#[cfg(test)]
556mod tests {
557 use std::sync::{Arc, Mutex};
558 use std::task::{Context, Waker};
559
560 use aion_core::{EventEnvelope, WorkflowStatus};
561 use aion_package::ContentHash;
562 use aion_store::visibility::VisibilityStore;
563 use aion_store::{EventStore, InMemoryStore};
564 use futures::{StreamExt, stream};
565 use serde_json::json;
566
567 use crate::durability::Recorder;
568 use crate::engine::api::EngineComponents;
569 use crate::registry::{CompletionNotifier, HandleResidency, WorkflowHandleParts};
570 use crate::{
571 Registry, RuntimeConfig, RuntimeHandle, SupervisionTree, WorkflowCatalog, WorkflowHandle,
572 };
573
574 use super::*;
575
576 #[derive(Debug, Default)]
577 struct SignalCapture {
578 calls: Mutex<Vec<(u64, String, Payload)>>,
579 }
580
581 #[async_trait]
582 impl SignalRouter for SignalCapture {
583 async fn route(
584 &self,
585 target: &WorkflowHandle,
586 name: String,
587 payload: Payload,
588 ) -> Result<(), EngineError> {
589 self.calls
590 .lock()
591 .map_err(|_| EngineError::RegistryPoisoned)?
592 .push((target.pid(), name, payload));
593 Ok(())
594 }
595 }
596
597 #[derive(Debug)]
598 struct QueryCapture {
599 calls: Mutex<Vec<(u64, String, Payload)>>,
600 reply: Payload,
601 }
602
603 #[async_trait]
604 impl QueryService for QueryCapture {
605 async fn query(
606 &self,
607 target: &WorkflowHandle,
608 name: String,
609 arguments: Payload,
610 ) -> Result<Payload, EngineError> {
611 self.calls
612 .lock()
613 .map_err(|_| EngineError::RegistryPoisoned)?
614 .push((target.pid(), name, arguments));
615 Ok(self.reply.clone())
616 }
617 }
618
619 #[derive(Debug)]
620 struct FakePublisher {
621 events: Vec<Event>,
622 }
623
624 impl EventPublisher for FakePublisher {
625 fn subscribe(
626 &self,
627 filter: EventFilter,
628 ) -> BoxStream<'static, Result<Event, EventStreamLagged>> {
629 let events = self
630 .events
631 .iter()
632 .filter(|event| filter.matches(event))
633 .cloned()
634 .map(Ok)
635 .collect::<Vec<_>>();
636 stream::iter(events).boxed()
637 }
638 }
639
640 fn payload(label: &str) -> Result<Payload, aion_core::PayloadError> {
641 Payload::from_json(&json!({ "label": label }))
642 }
643
644 fn engine_with_seams(
645 signal_router: Arc<dyn SignalRouter>,
646 query_service: Arc<dyn QueryService>,
647 event_publisher: Arc<dyn EventPublisher>,
648 ) -> Result<Engine, EngineError> {
649 let backing = Arc::new(InMemoryStore::default());
650 let store: Arc<dyn EventStore> = Arc::clone(&backing) as _;
651 let visibility_store: Arc<dyn VisibilityStore> = backing;
652 Ok(Engine::new(EngineComponents {
653 store,
654 visibility_store,
655 runtime: Arc::new(RuntimeHandle::new(RuntimeConfig::new(
656 Some(1),
657 crate::runtime::config::TEST_STOP_DRAIN_TIMEOUT,
658 ))?),
659 catalog: Arc::new(WorkflowCatalog::new()),
660 registry: Arc::new(Registry::default()),
661 supervision: Arc::new(SupervisionTree::new()),
662 delegated: DelegatedSeams::new(signal_router, query_service, event_publisher),
663 signal_handoff: Arc::new(crate::signal::SignalResumeHandoff::new()),
664 search_attribute_schema: Arc::new(aion_core::SearchAttributeSchema::new()),
665 visibility_reconciliation_task: None,
666 deferred_startup_recovery: None,
667 workloop: None,
668 }))
669 }
670
671 async fn recorded_active_handle(
675 engine: &Engine,
676 ) -> Result<WorkflowHandle, Box<dyn std::error::Error>> {
677 let workflow_id = WorkflowId::new_v4();
678 let run_id = RunId::new_v4();
679 let store = engine.store();
680 let mut recorder = Recorder::new(workflow_id.clone(), Arc::clone(&store));
681 recorder
682 .record_workflow_started(
683 chrono::Utc::now(),
684 crate::durability::WorkflowStartRecord {
685 workflow_type: "checkout".to_owned(),
686 input: payload("input")?,
687 run_id: run_id.clone(),
688 parent_run_id: None,
689 parent_workflow_id: None,
690 package_version: aion_core::PackageVersion::new("a".repeat(64)),
691 },
692 )
693 .await?;
694 Ok(WorkflowHandle::new(WorkflowHandleParts {
695 workflow_id,
696 run_id,
697 pid: engine.runtime().spawn_test_process_with_trap_exit(true)?,
698 workflow_type: "checkout".to_owned(),
699 namespace: String::from("default"),
700 loaded_version: ContentHash::from_bytes([1; 32]),
701 cached_status: WorkflowStatus::Running,
702 residency: HandleResidency::Resident,
703 recorder,
704 completion: CompletionNotifier::new(),
705 }))
706 }
707
708 async fn insert_active_handle(
709 engine: &Engine,
710 ) -> Result<WorkflowHandle, Box<dyn std::error::Error>> {
711 let handle = recorded_active_handle(engine).await?;
712 engine.registry().insert(
713 (handle.workflow_id().clone(), handle.run_id().clone()),
714 handle.clone(),
715 )?;
716 Ok(handle)
717 }
718
719 fn envelope(seq: u64, workflow_id: &WorkflowId) -> EventEnvelope {
720 EventEnvelope {
721 seq,
722 recorded_at: chrono::Utc::now(),
723 workflow_id: workflow_id.clone(),
724 }
725 }
726
727 #[tokio::test(flavor = "multi_thread")]
734 async fn signal_inside_the_registration_birth_window_waits_for_the_handle()
735 -> Result<(), Box<dyn std::error::Error>> {
736 let signal = Arc::new(SignalCapture::default());
737 let engine = Arc::new(engine_with_seams(
738 signal.clone(),
739 Arc::new(DeferredQueryService),
740 Arc::new(DeferredEventPublisher),
741 )?);
742 let handle = recorded_active_handle(&engine).await?;
743
744 let signal_call = engine.signal(
752 handle.workflow_id(),
753 handle.run_id(),
754 "approve",
755 payload("birth")?,
756 );
757 let mut signal_call = std::pin::pin!(signal_call);
758 let mut probe = Context::from_waker(Waker::noop());
759 assert!(
760 signal_call.as_mut().poll(&mut probe).is_pending(),
761 "the signal must wait the handle out, not answer before the insert"
762 );
763 engine.registry().insert(
764 (handle.workflow_id().clone(), handle.run_id().clone()),
765 handle.clone(),
766 )?;
767 signal_call.await?;
768
769 let calls = signal
770 .calls
771 .lock()
772 .map_err(|_| EngineError::RegistryPoisoned)?;
773 assert_eq!(calls.len(), 1, "the signal must reach the routed handle");
774 drop(calls);
775 engine.shutdown()?;
776 Ok(())
777 }
778
779 #[tokio::test(flavor = "multi_thread")]
783 async fn signal_for_a_started_run_with_no_handle_fails_typed_after_the_budget()
784 -> Result<(), Box<dyn std::error::Error>> {
785 let engine = engine_with_seams(
786 Arc::new(SignalCapture::default()),
787 Arc::new(DeferredQueryService),
788 Arc::new(DeferredEventPublisher),
789 )?;
790 let handle = recorded_active_handle(&engine).await?;
791
792 let outcome = engine
793 .signal(
794 handle.workflow_id(),
795 handle.run_id(),
796 "approve",
797 payload("never")?,
798 )
799 .await;
800
801 assert!(matches!(outcome, Err(EngineError::WorkflowNotFound { .. })));
802 engine.shutdown()?;
803 Ok(())
804 }
805
806 #[tokio::test]
807 async fn signal_delegates_to_router_and_unknown_returns_not_found()
808 -> Result<(), Box<dyn std::error::Error>> {
809 let signal = Arc::new(SignalCapture::default());
810 let engine = engine_with_seams(
811 signal.clone(),
812 Arc::new(DeferredQueryService),
813 Arc::new(DeferredEventPublisher),
814 )?;
815 let handle = insert_active_handle(&engine).await?;
816 let sent_payload = payload("signal")?;
817
818 engine
819 .signal(
820 handle.workflow_id(),
821 handle.run_id(),
822 "approve",
823 sent_payload.clone(),
824 )
825 .await?;
826
827 {
828 let calls = signal
829 .calls
830 .lock()
831 .map_err(|_| EngineError::RegistryPoisoned)?;
832 assert_eq!(
833 calls.as_slice(),
834 &[(handle.pid(), "approve".to_owned(), sent_payload)]
835 );
836 }
837 let unknown = engine
838 .signal(
839 &WorkflowId::new_v4(),
840 &RunId::new_v4(),
841 "approve",
842 payload("unknown")?,
843 )
844 .await;
845 assert!(matches!(unknown, Err(EngineError::WorkflowNotFound { .. })));
846 engine.shutdown()?;
847 Ok(())
848 }
849
850 #[tokio::test]
851 async fn query_delegates_to_service_and_returns_payload()
852 -> Result<(), Box<dyn std::error::Error>> {
853 let reply = payload("reply")?;
854 let query = Arc::new(QueryCapture {
855 calls: Mutex::new(Vec::new()),
856 reply: reply.clone(),
857 });
858 let engine = engine_with_seams(
859 Arc::new(DeferredSignalRouter),
860 query.clone(),
861 Arc::new(DeferredEventPublisher),
862 )?;
863 let handle = insert_active_handle(&engine).await?;
864
865 let arguments = payload("arguments")?;
866 let returned = engine
867 .query(
868 handle.workflow_id(),
869 handle.run_id(),
870 "state",
871 arguments.clone(),
872 )
873 .await?;
874
875 assert_eq!(returned, reply);
876 let calls = query
877 .calls
878 .lock()
879 .map_err(|_| EngineError::RegistryPoisoned)?;
880 assert_eq!(
883 calls.as_slice(),
884 &[(handle.pid(), "state".to_owned(), arguments)]
885 );
886 drop(calls);
887 engine.shutdown()?;
888 Ok(())
889 }
890
891 #[tokio::test]
892 async fn query_terminal_run_is_not_running_and_unknown_is_not_found()
893 -> Result<(), Box<dyn std::error::Error>> {
894 let engine = engine_with_seams(
895 Arc::new(DeferredSignalRouter),
896 Arc::new(DeferredQueryService),
897 Arc::new(DeferredEventPublisher),
898 )?;
899 let workflow_id = WorkflowId::new_v4();
901 let run_id = aion_core::RunId::new_v4();
902 let mut recorder = crate::durability::Recorder::new(workflow_id.clone(), engine.store());
903 recorder
904 .record_workflow_started(
905 chrono::Utc::now(),
906 crate::durability::WorkflowStartRecord {
907 workflow_type: "checkout".to_owned(),
908 input: payload("input")?,
909 run_id: run_id.clone(),
910 parent_run_id: None,
911 parent_workflow_id: None,
912 package_version: aion_core::PackageVersion::new("a".repeat(64)),
913 },
914 )
915 .await?;
916 recorder
917 .record_workflow_completed(chrono::Utc::now(), payload("result")?)
918 .await?;
919
920 let terminal = engine
921 .query(&workflow_id, &run_id, "state", Payload::json_null())
922 .await;
923 assert!(matches!(
924 terminal,
925 Err(EngineError::Query(crate::query::QueryError::NotRunning(id))) if id == workflow_id
926 ));
927
928 let unknown = engine
929 .query(
930 &WorkflowId::new_v4(),
931 &RunId::new_v4(),
932 "state",
933 Payload::json_null(),
934 )
935 .await;
936 assert!(matches!(unknown, Err(EngineError::WorkflowNotFound { .. })));
937 engine.shutdown()?;
938 Ok(())
939 }
940
941 #[tokio::test]
942 async fn subscribe_delegates_to_publisher_stream_with_filter()
943 -> Result<(), Box<dyn std::error::Error>> {
944 let workflow_id = WorkflowId::new_v4();
945 let other_id = WorkflowId::new_v4();
946 let matching = Event::SignalReceived {
947 envelope: envelope(1, &workflow_id),
948 name: "approved".to_owned(),
949 payload: payload("signal")?,
950 };
951 let filtered = Event::WorkflowStarted {
952 envelope: envelope(1, &other_id),
953 workflow_type: "checkout".to_owned(),
954 input: payload("input")?,
955 run_id: aion_core::RunId::new(uuid::Uuid::from_u128(1)),
956 parent_run_id: None,
957 parent_workflow_id: None,
958 package_version: aion_core::PackageVersion::new("a".repeat(64)),
959 };
960 let engine = engine_with_seams(
961 Arc::new(DeferredSignalRouter),
962 Arc::new(DeferredQueryService),
963 Arc::new(FakePublisher {
964 events: vec![matching.clone(), filtered],
965 }),
966 )?;
967
968 let events = engine
969 .subscribe(EventFilter {
970 workflow_id: Some(workflow_id),
971 run: None,
972 family: Some(EventFamily::Signal),
973 })
974 .collect::<Vec<_>>()
975 .await;
976
977 assert_eq!(events, vec![Ok(matching)]);
978 engine.shutdown()?;
979 Ok(())
980 }
981}