Skip to main content

everruns_engine/
phase_effects.rs

1//! Host-applied effects produced while an engine phase is running.
2//!
3//! Unlike transition lifecycle effects, phase effects cannot be buffered until
4//! Input/Reason/Act returns: output deltas, tool progress, and tool completion
5//! are live protocol signals. This contract still keeps recording outside the
6//! phase algorithm. The phase describes an effect and the injected host sink
7//! applies it immediately, preserving streaming and durable event order.
8
9use 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/// One effect produced during Input/Reason/Act execution.
16#[derive(Debug)]
17pub enum PhaseEffect {
18    /// Append and publish one canonical event.
19    EmitEvent(EventRequest),
20}
21
22/// Host-applied sink for live phase effects.
23///
24/// Engine phase constructors adapt the host's existing [`EventEmitter`] to
25/// this contract, so custom hosts retain the open event-storage boundary while
26/// phase algorithms no longer invoke the persistence-shaped operation directly.
27#[async_trait]
28pub trait PhaseEffectSink: Send + Sync {
29    /// Apply one phase effect and return the accepted canonical event.
30    async fn apply_phase_effect(&self, effect: PhaseEffect) -> Result<Event>;
31
32    /// Convenience for the common canonical-event effect.
33    async fn emit_phase_event(&self, request: EventRequest) -> Result<Event> {
34        self.apply_phase_effect(PhaseEffect::EmitEvent(request))
35            .await
36    }
37}
38
39/// Adapter installed by engine phases around a host event emitter.
40pub(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}