everruns_engine/
phase_effects.rs1use async_trait::async_trait;
10use everruns_core::event_emitter::EventEmitter;
11use everruns_core::events::{Event, EventRequest};
12use everruns_provider::error::Result;
13use std::sync::Arc;
14
15#[derive(Debug)]
17pub enum PhaseEffect {
18 EmitEvent(EventRequest),
20}
21
22#[async_trait]
28pub trait PhaseEffectSink: Send + Sync {
29 async fn apply_phase_effect(&self, effect: PhaseEffect) -> Result<Event>;
31
32 async fn emit_phase_event(&self, request: EventRequest) -> Result<Event> {
34 self.apply_phase_effect(PhaseEffect::EmitEvent(request))
35 .await
36 }
37}
38
39pub(crate) struct PhaseEffectEmitter<E: ?Sized> {
41 inner: Arc<E>,
42}
43
44impl<E: ?Sized> Clone for PhaseEffectEmitter<E> {
45 fn clone(&self) -> Self {
46 Self {
47 inner: self.inner.clone(),
48 }
49 }
50}
51
52impl<E: ?Sized> PhaseEffectEmitter<E> {
53 pub(crate) fn new(inner: Arc<E>) -> Self {
54 Self { inner }
55 }
56
57 pub(crate) async fn emit(&self, request: EventRequest) -> Result<Event>
58 where
59 E: PhaseEffectSink,
60 {
61 self.inner
62 .apply_phase_effect(PhaseEffect::EmitEvent(request))
63 .await
64 }
65}
66
67impl AsRef<dyn EventEmitter> for PhaseEffectEmitter<dyn PhaseEffectSink> {
68 fn as_ref(&self) -> &(dyn EventEmitter + 'static) {
69 self
70 }
71}
72
73#[async_trait]
74impl<T: EventEmitter + ?Sized> PhaseEffectSink for T {
75 async fn apply_phase_effect(&self, effect: PhaseEffect) -> Result<Event> {
76 match effect {
77 PhaseEffect::EmitEvent(request) => self.emit(request).await,
78 }
79 }
80}
81
82#[async_trait]
83impl<T: PhaseEffectSink + ?Sized> EventEmitter for PhaseEffectEmitter<T> {
84 async fn emit(&self, request: EventRequest) -> Result<Event> {
85 PhaseEffectEmitter::emit(self, request).await
86 }
87}
88
89#[cfg(test)]
90mod tests {
91 use super::*;
92 use everruns_core::Message;
93 use everruns_core::events::{EventContext, INPUT_MESSAGE, InputMessageData};
94 use everruns_provider::typed_id::SessionId;
95
96 #[tokio::test]
97 async fn phase_effect_adapter_preserves_host_sequence_and_event_type() {
98 let host = Arc::new(crate::test_fixtures::TestEventEmitter::new());
99 let sink = PhaseEffectEmitter::new(host.clone());
100 let session_id = SessionId::new();
101
102 let first = sink
103 .emit_phase_event(EventRequest::new(
104 session_id,
105 EventContext::empty(),
106 InputMessageData::new(Message::user("first")),
107 ))
108 .await
109 .unwrap();
110 let second = sink
111 .emit(EventRequest::new(
112 session_id,
113 EventContext::empty(),
114 InputMessageData::new(Message::user("second")),
115 ))
116 .await
117 .unwrap();
118
119 assert_eq!((first.sequence, second.sequence), (Some(1), Some(2)));
120 assert_eq!(
121 (first.event_type.as_str(), second.event_type.as_str()),
122 (INPUT_MESSAGE, INPUT_MESSAGE)
123 );
124 let recorded = host.events().await;
125 assert_eq!(recorded.len(), 2);
126 assert_eq!(recorded[0].id, first.id);
127 assert_eq!(recorded[1].id, second.id);
128 }
129}