lash_core/runtime/effect/
inline_host.rs1use 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#[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 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}