Skip to main content

lash_core/runtime/effect/
inline_host.rs

1use std::sync::Arc;
2use std::time::Instant;
3
4use tokio_util::sync::CancellationToken;
5
6use super::{
7    AwaitEventKey, AwaitEventResolver, AwaitEventWaitIdentity, BoundaryReason, EffectHost,
8    ExecutionScope, InlineRuntimeEffectController, Resolution, ResolveOutcome,
9    RuntimeEffectController, RuntimeEffectControllerError, RuntimeEffectEnvelope,
10    RuntimeEffectLocalExecutor, RuntimeEffectOutcome, ScopedEffectController, SegmentProgress,
11};
12use crate::RuntimeError;
13
14/// In-process deployment effect host.
15#[derive(Clone)]
16pub struct InlineEffectHost {
17    controller: Arc<dyn RuntimeEffectController>,
18    allow_process_lifetime_completion_keys: Arc<std::sync::atomic::AtomicBool>,
19}
20
21impl InlineEffectHost {
22    pub fn new(controller: Arc<dyn RuntimeEffectController>) -> Self {
23        Self {
24            controller,
25            allow_process_lifetime_completion_keys: Arc::new(std::sync::atomic::AtomicBool::new(
26                false,
27            )),
28        }
29    }
30
31    /// Explicitly accept that externally routed completion keys die with this
32    /// process. Intended only for deliberately single-process embeddings.
33    pub fn allow_process_lifetime_completion_keys(self) -> Self {
34        self.allow_process_lifetime_completion_keys
35            .store(true, std::sync::atomic::Ordering::Relaxed);
36        self
37    }
38}
39
40impl Default for InlineEffectHost {
41    fn default() -> Self {
42        Self::new(Arc::new(InlineRuntimeEffectController::default()))
43    }
44}
45
46#[async_trait::async_trait]
47impl AwaitEventResolver for InlineEffectHost {
48    fn durability_tier(&self) -> crate::DurabilityTier {
49        self.controller.durability_tier()
50    }
51
52    fn allows_process_lifetime_completion_keys(&self) -> bool {
53        self.controller.allows_process_lifetime_completion_keys()
54            || self
55                .allow_process_lifetime_completion_keys
56                .load(std::sync::atomic::Ordering::Relaxed)
57    }
58
59    async fn await_event_key(
60        &self,
61        scope: &ExecutionScope,
62        wait: AwaitEventWaitIdentity,
63    ) -> Result<AwaitEventKey, RuntimeError> {
64        self.controller.await_event_key(scope, wait).await
65    }
66
67    async fn resolve_await_event(
68        &self,
69        key: &AwaitEventKey,
70        resolution: Resolution,
71    ) -> Result<ResolveOutcome, RuntimeError> {
72        self.controller.resolve_await_event(key, resolution).await
73    }
74
75    async fn peek_await_event(
76        &self,
77        key: &AwaitEventKey,
78    ) -> Result<Option<Resolution>, RuntimeError> {
79        self.controller.peek_await_event(key).await
80    }
81
82    async fn await_await_event(
83        &self,
84        key: &AwaitEventKey,
85        cancel: CancellationToken,
86        deadline: Option<Instant>,
87    ) -> Result<Resolution, RuntimeError> {
88        self.controller
89            .await_await_event(key, cancel, deadline)
90            .await
91    }
92
93    async fn revoke_await_events_for_session(&self, session_id: &str) -> Result<(), RuntimeError> {
94        self.controller
95            .revoke_await_events_for_session(session_id)
96            .await
97    }
98
99    async fn cancel_await_events_for_session(&self, session_id: &str) -> Result<(), RuntimeError> {
100        self.controller
101            .cancel_await_events_for_session(session_id)
102            .await
103    }
104}
105
106#[async_trait::async_trait]
107impl EffectHost for InlineEffectHost {
108    fn scoped<'run>(
109        &'run self,
110        scope: ExecutionScope,
111    ) -> Result<ScopedEffectController<'run>, RuntimeError> {
112        ScopedEffectController::shared(
113            Arc::new(InlineHostScopedController {
114                controller: Arc::clone(&self.controller),
115                allow_process_lifetime_completion_keys: Arc::clone(
116                    &self.allow_process_lifetime_completion_keys,
117                ),
118            }),
119            scope,
120        )
121    }
122
123    fn scoped_static(
124        &self,
125        scope: ExecutionScope,
126    ) -> Result<Option<ScopedEffectController<'static>>, RuntimeError> {
127        Ok(Some(ScopedEffectController::shared(
128            Arc::new(InlineHostScopedController {
129                controller: Arc::clone(&self.controller),
130                allow_process_lifetime_completion_keys: Arc::clone(
131                    &self.allow_process_lifetime_completion_keys,
132                ),
133            }),
134            scope,
135        )?))
136    }
137}
138
139#[derive(Clone)]
140struct InlineHostScopedController {
141    controller: Arc<dyn RuntimeEffectController>,
142    allow_process_lifetime_completion_keys: Arc<std::sync::atomic::AtomicBool>,
143}
144
145#[async_trait::async_trait]
146impl AwaitEventResolver for InlineHostScopedController {
147    fn durability_tier(&self) -> crate::DurabilityTier {
148        self.controller.durability_tier()
149    }
150
151    fn allows_process_lifetime_completion_keys(&self) -> bool {
152        self.controller.allows_process_lifetime_completion_keys()
153            || self
154                .allow_process_lifetime_completion_keys
155                .load(std::sync::atomic::Ordering::Relaxed)
156    }
157
158    async fn await_event_key(
159        &self,
160        scope: &ExecutionScope,
161        wait: AwaitEventWaitIdentity,
162    ) -> Result<AwaitEventKey, RuntimeError> {
163        self.controller.await_event_key(scope, wait).await
164    }
165
166    async fn resolve_await_event(
167        &self,
168        key: &AwaitEventKey,
169        resolution: Resolution,
170    ) -> Result<ResolveOutcome, RuntimeError> {
171        self.controller.resolve_await_event(key, resolution).await
172    }
173
174    async fn peek_await_event(
175        &self,
176        key: &AwaitEventKey,
177    ) -> Result<Option<Resolution>, RuntimeError> {
178        self.controller.peek_await_event(key).await
179    }
180
181    async fn await_await_event(
182        &self,
183        key: &AwaitEventKey,
184        cancel: CancellationToken,
185        deadline: Option<Instant>,
186    ) -> Result<Resolution, RuntimeError> {
187        self.controller
188            .await_await_event(key, cancel, deadline)
189            .await
190    }
191
192    async fn revoke_await_events_for_session(&self, session_id: &str) -> Result<(), RuntimeError> {
193        self.controller
194            .revoke_await_events_for_session(session_id)
195            .await
196    }
197
198    async fn cancel_await_events_for_session(&self, session_id: &str) -> Result<(), RuntimeError> {
199        self.controller
200            .cancel_await_events_for_session(session_id)
201            .await
202    }
203}
204
205#[async_trait::async_trait]
206impl RuntimeEffectController for InlineHostScopedController {
207    fn wants_segment_boundary(&self, progress: &SegmentProgress) -> Option<BoundaryReason> {
208        self.controller.wants_segment_boundary(progress)
209    }
210
211    fn supports_concurrent_effects(&self) -> bool {
212        self.controller.supports_concurrent_effects()
213    }
214
215    async fn execute_effect(
216        &self,
217        envelope: RuntimeEffectEnvelope,
218        local_executor: RuntimeEffectLocalExecutor<'_>,
219    ) -> Result<RuntimeEffectOutcome, RuntimeEffectControllerError> {
220        self.controller
221            .execute_effect(envelope, local_executor)
222            .await
223    }
224}