1use std::path::Path;
14use std::sync::Arc;
15use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
16use std::time::{Duration, Instant};
17
18use lash_core::{
19 AwaitEventKey, AwaitEventResolver, AwaitEventWaitIdentity, CanonicalRuntimeEffectEnvelope,
20 DurabilityTier, EffectHost, ExecutionScope, LeaseTimings, Resolution, ResolveOutcome,
21 RuntimeAwaitEventOptions, RuntimeEffectCommand, RuntimeEffectController,
22 RuntimeEffectControllerError, RuntimeEffectEnvelope, RuntimeEffectLocalExecutor,
23 RuntimeEffectOutcome, RuntimeError, ScopedEffectController, validate_replayed_effect_envelope,
24};
25use tokio_util::sync::CancellationToken;
26
27use super::*;
28use crate::await_event::SqliteAwaitEvents;
29
30const STATUS_IN_PROGRESS: &str = "in_progress";
31const STATUS_COMPLETED: &str = "completed";
32const STATUS_FAILED: &str = "failed";
33const BUSY_POLL: Duration = Duration::from_millis(25);
34
35static EFFECT_OWNER_COUNTER: AtomicU64 = AtomicU64::new(1);
36
37#[derive(Clone, Debug, Default)]
39pub struct SqliteEffectReplayOptions {
40 pub lease_timings: LeaseTimings,
44}
45
46struct SqliteEffectReplayInner {
47 conn: SqliteConnection,
48 clock: Arc<dyn lash_core::Clock>,
49 owner_id: String,
50 lease_counter: AtomicU64,
51 replay_mode: AtomicBool,
52 lease_timings: LeaseTimings,
53 await_events: SqliteAwaitEvents,
54}
55
56#[derive(Clone)]
62pub struct SqliteEffectHost {
63 inner: Arc<SqliteEffectReplayInner>,
64}
65
66#[derive(Clone)]
68pub struct SqliteRuntimeEffectController {
69 inner: Arc<SqliteEffectReplayInner>,
70 scope: ExecutionScope,
71}
72
73struct ClaimedEffect {
74 scope_id: String,
75 replay_key: String,
76 envelope_hash: String,
77 lease_token: String,
78 due_at_ms: Option<u64>,
79}
80
81enum PreparedEffect {
82 ReplayMismatch {
83 recorded_envelope: Box<CanonicalRuntimeEffectEnvelope>,
84 stored_envelope_hash: String,
85 },
86 ReplayOutcome {
87 outcome: Box<RuntimeEffectOutcome>,
88 due_at_ms: Option<u64>,
89 },
90 ReplayError(RuntimeEffectControllerError),
91 Claimed(ClaimedEffect),
92 Busy {
93 retry_at_ms: u64,
94 },
95}
96
97impl SqliteEffectHost {
98 pub async fn open(path: &Path) -> tokio_rusqlite::Result<Self> {
99 Self::open_with_options(path, SqliteEffectReplayOptions::default()).await
100 }
101
102 pub async fn open_with_clock(
103 path: &Path,
104 clock: Arc<dyn lash_core::Clock>,
105 ) -> tokio_rusqlite::Result<Self> {
106 Self::open_with_options_and_clock(path, SqliteEffectReplayOptions::default(), clock).await
107 }
108
109 pub async fn open_with_options(
110 path: &Path,
111 options: SqliteEffectReplayOptions,
112 ) -> tokio_rusqlite::Result<Self> {
113 Self::open_with_options_and_clock(path, options, Arc::new(lash_core::SystemClock)).await
114 }
115
116 pub async fn open_with_options_and_clock(
117 path: &Path,
118 options: SqliteEffectReplayOptions,
119 clock: Arc<dyn lash_core::Clock>,
120 ) -> tokio_rusqlite::Result<Self> {
121 Ok(Self {
122 inner: open_effect_replay_inner(path, StoreBacking::File, options, clock).await?,
123 })
124 }
125
126 pub async fn memory() -> tokio_rusqlite::Result<Self> {
127 Self::memory_with_options(SqliteEffectReplayOptions::default()).await
128 }
129
130 pub async fn memory_with_clock(
131 clock: Arc<dyn lash_core::Clock>,
132 ) -> tokio_rusqlite::Result<Self> {
133 Self::memory_with_options_and_clock(SqliteEffectReplayOptions::default(), clock).await
134 }
135
136 pub async fn memory_with_options(
137 options: SqliteEffectReplayOptions,
138 ) -> tokio_rusqlite::Result<Self> {
139 Self::memory_with_options_and_clock(options, Arc::new(lash_core::SystemClock)).await
140 }
141
142 pub async fn memory_with_options_and_clock(
143 options: SqliteEffectReplayOptions,
144 clock: Arc<dyn lash_core::Clock>,
145 ) -> tokio_rusqlite::Result<Self> {
146 Ok(Self {
147 inner: open_effect_replay_memory_inner(options, clock).await?,
148 })
149 }
150
151 pub fn start_replay(&self) {
154 self.inner.replay_mode.store(true, Ordering::SeqCst);
155 }
156}
157
158#[async_trait::async_trait]
159impl AwaitEventResolver for SqliteEffectHost {
160 fn durability_tier(&self) -> DurabilityTier {
161 DurabilityTier::Durable
162 }
163
164 async fn await_event_key(
165 &self,
166 scope: &ExecutionScope,
167 wait: AwaitEventWaitIdentity,
168 ) -> Result<AwaitEventKey, RuntimeError> {
169 self.inner.await_events.key_for(scope, wait).await
170 }
171
172 async fn resolve_await_event(
173 &self,
174 key: &AwaitEventKey,
175 resolution: Resolution,
176 ) -> Result<ResolveOutcome, RuntimeError> {
177 self.inner.await_events.resolve(key, resolution).await
178 }
179
180 async fn peek_await_event(
181 &self,
182 key: &AwaitEventKey,
183 ) -> Result<Option<Resolution>, RuntimeError> {
184 self.inner.await_events.peek(key).await
185 }
186
187 async fn await_await_event(
188 &self,
189 key: &AwaitEventKey,
190 cancel: CancellationToken,
191 deadline: Option<Instant>,
192 ) -> Result<Resolution, RuntimeError> {
193 self.inner
194 .await_events
195 .await_resolution(key, cancel, deadline)
196 .await
197 }
198
199 async fn revoke_await_events_for_session(&self, session_id: &str) -> Result<(), RuntimeError> {
200 self.inner.await_events.revoke_session(session_id).await
201 }
202
203 async fn cancel_await_events_for_session(&self, session_id: &str) -> Result<(), RuntimeError> {
204 self.inner.await_events.cancel_session(session_id).await
205 }
206}
207
208impl EffectHost for SqliteEffectHost {
209 fn scoped<'run>(
210 &'run self,
211 scope: ExecutionScope,
212 ) -> Result<ScopedEffectController<'run>, RuntimeError> {
213 let controller = SqliteRuntimeEffectController {
214 inner: Arc::clone(&self.inner),
215 scope: scope.clone(),
216 };
217 ScopedEffectController::shared(Arc::new(controller), scope)
218 }
219
220 fn scoped_static(
221 &self,
222 scope: ExecutionScope,
223 ) -> Result<Option<ScopedEffectController<'static>>, RuntimeError> {
224 let controller = SqliteRuntimeEffectController {
225 inner: Arc::clone(&self.inner),
226 scope: scope.clone(),
227 };
228 Ok(Some(ScopedEffectController::shared(
229 Arc::new(controller),
230 scope,
231 )?))
232 }
233}
234
235impl SqliteRuntimeEffectController {
236 pub async fn open(path: &Path, scope: ExecutionScope) -> tokio_rusqlite::Result<Self> {
237 Self::open_with_options(path, scope, SqliteEffectReplayOptions::default()).await
238 }
239
240 pub async fn open_with_clock(
241 path: &Path,
242 scope: ExecutionScope,
243 clock: Arc<dyn lash_core::Clock>,
244 ) -> tokio_rusqlite::Result<Self> {
245 Self::open_with_options_and_clock(path, scope, SqliteEffectReplayOptions::default(), clock)
246 .await
247 }
248
249 pub async fn open_with_options(
250 path: &Path,
251 scope: ExecutionScope,
252 options: SqliteEffectReplayOptions,
253 ) -> tokio_rusqlite::Result<Self> {
254 Self::open_with_options_and_clock(path, scope, options, Arc::new(lash_core::SystemClock))
255 .await
256 }
257
258 pub async fn open_with_options_and_clock(
259 path: &Path,
260 scope: ExecutionScope,
261 options: SqliteEffectReplayOptions,
262 clock: Arc<dyn lash_core::Clock>,
263 ) -> tokio_rusqlite::Result<Self> {
264 Ok(Self {
265 inner: open_effect_replay_inner(path, StoreBacking::File, options, clock).await?,
266 scope,
267 })
268 }
269
270 pub async fn memory(scope: ExecutionScope) -> tokio_rusqlite::Result<Self> {
271 Self::memory_with_options(scope, SqliteEffectReplayOptions::default()).await
272 }
273
274 pub async fn memory_with_clock(
275 scope: ExecutionScope,
276 clock: Arc<dyn lash_core::Clock>,
277 ) -> tokio_rusqlite::Result<Self> {
278 Self::memory_with_options_and_clock(scope, SqliteEffectReplayOptions::default(), clock)
279 .await
280 }
281
282 pub async fn memory_with_options(
283 scope: ExecutionScope,
284 options: SqliteEffectReplayOptions,
285 ) -> tokio_rusqlite::Result<Self> {
286 Self::memory_with_options_and_clock(scope, options, Arc::new(lash_core::SystemClock)).await
287 }
288
289 pub async fn memory_with_options_and_clock(
290 scope: ExecutionScope,
291 options: SqliteEffectReplayOptions,
292 clock: Arc<dyn lash_core::Clock>,
293 ) -> tokio_rusqlite::Result<Self> {
294 Ok(Self {
295 inner: open_effect_replay_memory_inner(options, clock).await?,
296 scope,
297 })
298 }
299
300 pub fn start_replay(&self) {
303 self.inner.replay_mode.store(true, Ordering::SeqCst);
304 }
305
306 async fn prepare_effect(
307 &self,
308 envelope: &RuntimeEffectEnvelope,
309 reconstructed_envelope: &CanonicalRuntimeEffectEnvelope,
310 ) -> Result<PreparedEffect, RuntimeEffectControllerError> {
311 let replay_key = envelope
312 .invocation
313 .replay_key()
314 .ok_or_else(|| {
315 RuntimeEffectControllerError::new(
316 "sqlite_effect_replay_key_missing",
317 "runtime effect envelope requires replay.key",
318 )
319 })?
320 .to_string();
321 let envelope_hash = reconstructed_envelope.hash().to_string();
322 let envelope_json =
323 serde_json::to_string(reconstructed_envelope).map_err(effect_encode_error)?;
324 let scope_id = self.scope.id().to_string();
325 let now = self.inner.clock.timestamp_ms();
326 let lease_token = self.inner.next_lease_token();
327 let due_at_ms = sleep_due_at_ms(envelope, now);
328 let lease_ttl_ms = self.inner.lease_timings.ttl_ms();
329 let lease_expires_at_ms = now.saturating_add(lease_ttl_ms);
330 let replay_mode = self.inner.replay_mode.load(Ordering::SeqCst);
331 let owner_id = self.inner.owner_id.clone();
332
333 let outcome: Result<PreparedEffect, RuntimeEffectControllerError> = self
341 .inner
342 .conn
343 .write(move |tx| {
344 let row = tx
345 .query_row(
346 "SELECT envelope_hash, envelope_json, status, outcome_json, error_json,
347 lease_owner_id, lease_token, lease_expires_at_ms, due_at_ms
348 FROM runtime_effect_replay
349 WHERE scope_id = ?1 AND replay_key = ?2",
350 params![scope_id.as_str(), replay_key.as_str()],
351 |row| {
352 Ok((
353 row.get::<_, String>(0)?,
354 row.get::<_, String>(1)?,
355 row.get::<_, String>(2)?,
356 row.get::<_, Option<String>>(3)?,
357 row.get::<_, Option<String>>(4)?,
358 row.get::<_, i64>(7)?,
359 row.get::<_, Option<i64>>(8)?,
360 ))
361 },
362 )
363 .optional()?;
364
365 let Some((
366 existing_hash,
367 existing_envelope_json,
368 status,
369 outcome_json,
370 error_json,
371 lease_expires_row,
372 existing_due_row,
373 )) = row
374 else {
375 if replay_mode {
376 return Ok(Err(RuntimeEffectControllerError::new(
377 "sqlite_effect_replay_missing",
378 format!(
379 "no recorded runtime effect for scope `{scope_id}` and replay key `{replay_key}`"
380 ),
381 )));
382 }
383 let due_at_param = due_at_ms.map(|value| value as i64);
384 tx.execute(
385 "INSERT INTO runtime_effect_replay (
386 scope_id, replay_key, envelope_hash, envelope_json, status, outcome_json,
387 error_json, lease_owner_id, lease_token, lease_expires_at_ms,
388 due_at_ms, created_at_ms, updated_at_ms
389 )
390 VALUES (?1, ?2, ?3, ?4, ?5, NULL, NULL, ?6, ?7, ?8, ?9, ?10, ?11)",
391 params![
392 scope_id.as_str(),
393 replay_key.as_str(),
394 envelope_hash.as_str(),
395 envelope_json.as_str(),
396 STATUS_IN_PROGRESS,
397 owner_id.as_str(),
398 lease_token.as_str(),
399 lease_expires_at_ms as i64,
400 due_at_param,
401 now as i64,
402 now as i64,
403 ],
404 )?;
405 return Ok(Ok(PreparedEffect::Claimed(ClaimedEffect {
406 scope_id,
407 replay_key,
408 envelope_hash,
409 lease_token,
410 due_at_ms,
411 })));
412 };
413
414 if existing_hash != envelope_hash {
415 let recorded_envelope: CanonicalRuntimeEffectEnvelope =
416 match serde_json::from_str(&existing_envelope_json) {
417 Ok(envelope) => envelope,
418 Err(err) => return Ok(Err(effect_decode_error(err))),
419 };
420 return Ok(Ok(PreparedEffect::ReplayMismatch {
421 recorded_envelope: Box::new(recorded_envelope),
422 stored_envelope_hash: existing_hash,
423 }));
424 }
425
426 let lease_expires_at_ms = lease_expires_row as u64;
427 let existing_due_at_ms = existing_due_row.map(|value| value as u64);
428
429 match status.as_str() {
430 STATUS_COMPLETED => {
431 let Some(json) = outcome_json else {
432 return Ok(Err(RuntimeEffectControllerError::new(
433 "sqlite_effect_replay_corrupt_row",
434 "completed runtime effect row is missing outcome_json",
435 )));
436 };
437 let outcome = match serde_json::from_str(&json) {
438 Ok(outcome) => outcome,
439 Err(err) => return Ok(Err(effect_decode_error(err))),
440 };
441 Ok(Ok(PreparedEffect::ReplayOutcome {
442 outcome: Box::new(outcome),
443 due_at_ms: existing_due_at_ms,
444 }))
445 }
446 STATUS_FAILED => {
447 let Some(json) = error_json else {
448 return Ok(Err(RuntimeEffectControllerError::new(
449 "sqlite_effect_replay_corrupt_row",
450 "failed runtime effect row is missing error_json",
451 )));
452 };
453 let err = match serde_json::from_str(&json) {
454 Ok(err) => err,
455 Err(err) => return Ok(Err(effect_decode_error(err))),
456 };
457 Ok(Ok(PreparedEffect::ReplayError(err)))
458 }
459 STATUS_IN_PROGRESS if lease_expires_at_ms > now => Ok(Ok(PreparedEffect::Busy {
460 retry_at_ms: lease_expires_at_ms,
461 })),
462 STATUS_IN_PROGRESS => {
463 let due_at_ms = existing_due_at_ms.or(due_at_ms);
464 let due_at_param = due_at_ms.map(|value| value as i64);
465 tx.execute(
466 "UPDATE runtime_effect_replay
467 SET lease_owner_id = ?3,
468 lease_token = ?4,
469 lease_expires_at_ms = ?5,
470 due_at_ms = ?6,
471 updated_at_ms = ?7
472 WHERE scope_id = ?1 AND replay_key = ?2",
473 params![
474 scope_id.as_str(),
475 replay_key.as_str(),
476 owner_id.as_str(),
477 lease_token.as_str(),
478 now.saturating_add(lease_ttl_ms) as i64,
479 due_at_param,
480 now as i64,
481 ],
482 )?;
483 Ok(Ok(PreparedEffect::Claimed(ClaimedEffect {
484 scope_id,
485 replay_key,
486 envelope_hash,
487 lease_token,
488 due_at_ms,
489 })))
490 }
491 other => Ok(Err(RuntimeEffectControllerError::new(
492 "sqlite_effect_replay_corrupt_row",
493 format!("unknown runtime effect replay status `{other}`"),
494 ))),
495 }
496 })
497 .await
498 .map_err(effect_sqlite_error)?;
499 outcome
500 }
501
502 async fn finalize_effect(
503 &self,
504 claim: &ClaimedEffect,
505 outcome: &Result<RuntimeEffectOutcome, RuntimeEffectControllerError>,
506 ) -> Result<(), RuntimeEffectControllerError> {
507 let (status, outcome_json, error_json) = match outcome {
508 Ok(outcome) => (
509 STATUS_COMPLETED,
510 Some(serde_json::to_string(outcome).map_err(effect_encode_error)?),
511 None,
512 ),
513 Err(err) => (
514 STATUS_FAILED,
515 None,
516 Some(serde_json::to_string(err).map_err(effect_encode_error)?),
517 ),
518 };
519 let now = self.inner.clock.timestamp_ms();
520 let scope_id = claim.scope_id.clone();
521 let replay_key = claim.replay_key.clone();
522 let envelope_hash = claim.envelope_hash.clone();
523 let owner_id = self.inner.owner_id.clone();
524 let lease_token = claim.lease_token.clone();
525 let status = status.to_string();
526
527 let result: Result<(), RuntimeEffectControllerError> = self
528 .inner
529 .conn
530 .write(move |tx| {
531 let changed = tx.execute(
532 "UPDATE runtime_effect_replay
533 SET status = ?6,
534 outcome_json = ?7,
535 error_json = ?8,
536 lease_owner_id = NULL,
537 lease_token = NULL,
538 lease_expires_at_ms = 0,
539 updated_at_ms = ?9
540 WHERE scope_id = ?1
541 AND replay_key = ?2
542 AND envelope_hash = ?3
543 AND lease_owner_id = ?4
544 AND lease_token = ?5
545 AND status = 'in_progress'
546 AND lease_expires_at_ms > ?10",
547 params![
548 scope_id.as_str(),
549 replay_key.as_str(),
550 envelope_hash.as_str(),
551 owner_id.as_str(),
552 lease_token.as_str(),
553 status.as_str(),
554 outcome_json,
555 error_json,
556 now as i64,
557 now as i64,
558 ],
559 )?;
560 if changed != 1 {
561 return Ok(Err(RuntimeEffectControllerError::new(
562 "sqlite_effect_replay_lease_lost",
563 format!(
564 "runtime effect replay lease was lost before finalizing scope `{scope_id}` replay key `{replay_key}`"
565 ),
566 )));
567 }
568 Ok(Ok(()))
569 })
570 .await
571 .map_err(effect_sqlite_error)?;
572 result
573 }
574
575 async fn renew_effect_lease(
576 &self,
577 claim: &ClaimedEffect,
578 ) -> Result<(), RuntimeEffectControllerError> {
579 let now = self.inner.clock.timestamp_ms();
580 let renewed_expires_at = now.saturating_add(self.inner.lease_timings.ttl_ms());
581 let scope_id = claim.scope_id.clone();
582 let replay_key = claim.replay_key.clone();
583 let envelope_hash = claim.envelope_hash.clone();
584 let owner_id = self.inner.owner_id.clone();
585 let lease_token = claim.lease_token.clone();
586
587 let result: Result<(), RuntimeEffectControllerError> = self
588 .inner
589 .conn
590 .write(move |tx| {
591 let changed = tx.execute(
592 "UPDATE runtime_effect_replay
593 SET lease_expires_at_ms = ?6,
594 updated_at_ms = ?7
595 WHERE scope_id = ?1
596 AND replay_key = ?2
597 AND envelope_hash = ?3
598 AND lease_owner_id = ?4
599 AND lease_token = ?5
600 AND status = 'in_progress'
601 AND lease_expires_at_ms > ?8",
602 params![
603 scope_id.as_str(),
604 replay_key.as_str(),
605 envelope_hash.as_str(),
606 owner_id.as_str(),
607 lease_token.as_str(),
608 renewed_expires_at as i64,
609 now as i64,
610 now as i64,
611 ],
612 )?;
613 if changed != 1 {
614 return Ok(Err(RuntimeEffectControllerError::new(
615 "sqlite_effect_replay_lease_lost",
616 format!(
617 "runtime effect replay lease was lost while executing scope `{scope_id}` replay key `{replay_key}`"
618 ),
619 )));
620 }
621 Ok(Ok(()))
622 })
623 .await
624 .map_err(effect_sqlite_error)?;
625 result
626 }
627
628 async fn execute_claimed_effect_with_renewal(
629 &self,
630 claim: &ClaimedEffect,
631 envelope: RuntimeEffectEnvelope,
632 local_executor: RuntimeEffectLocalExecutor<'_>,
633 ) -> Result<RuntimeEffectOutcome, RuntimeEffectControllerError> {
634 let renew_every = self.inner.lease_timings.renew_interval();
635 let effect = self.execute_claimed_effect(claim, envelope, local_executor);
636 tokio::pin!(effect);
637
638 loop {
639 tokio::select! {
640 result = &mut effect => return result,
641 _ = self.inner.clock.sleep(renew_every) => {
642 self.renew_effect_lease(claim).await?;
643 }
644 }
645 }
646 }
647
648 async fn execute_claimed_effect(
649 &self,
650 claim: &ClaimedEffect,
651 envelope: RuntimeEffectEnvelope,
652 local_executor: RuntimeEffectLocalExecutor<'_>,
653 ) -> Result<RuntimeEffectOutcome, RuntimeEffectControllerError> {
654 if matches!(envelope.command, RuntimeEffectCommand::Sleep { .. }) {
655 sleep_until_due(self.inner.clock.as_ref(), claim.due_at_ms).await;
656 return Ok(RuntimeEffectOutcome::Sleep);
657 }
658 match envelope.command {
659 RuntimeEffectCommand::PeekAwaitEvent { key } => {
660 let resolution = self
661 .peek_await_event(&key)
662 .await
663 .map_err(RuntimeEffectControllerError::from)?;
664 Ok(RuntimeEffectOutcome::PeekAwaitEvent { resolution })
665 }
666 RuntimeEffectCommand::AwaitEvent { key } => {
667 let RuntimeAwaitEventOptions {
668 cancellation,
669 deadline,
670 clock,
671 ..
672 } = local_executor.into_await_event_options()?;
673 let resolution = self
674 .inner
675 .await_events
676 .await_resolution_with_clock(&key, cancellation, deadline, clock.as_ref())
677 .await
678 .map_err(RuntimeEffectControllerError::from)?;
679 Ok(RuntimeEffectOutcome::AwaitEvent { resolution })
680 }
681 RuntimeEffectCommand::Process { command } => {
682 let result = local_executor.into_process()?.execute(*command).await?;
683 Ok(RuntimeEffectOutcome::Process { result })
684 }
685 _ => local_executor.execute(envelope).await,
686 }
687 }
688}
689
690#[async_trait::async_trait]
691impl AwaitEventResolver for SqliteRuntimeEffectController {
692 fn durability_tier(&self) -> DurabilityTier {
693 DurabilityTier::Durable
694 }
695
696 async fn await_event_key(
697 &self,
698 scope: &ExecutionScope,
699 wait: AwaitEventWaitIdentity,
700 ) -> Result<AwaitEventKey, RuntimeError> {
701 self.inner.await_events.key_for(scope, wait).await
702 }
703
704 async fn resolve_await_event(
705 &self,
706 key: &AwaitEventKey,
707 resolution: Resolution,
708 ) -> Result<ResolveOutcome, RuntimeError> {
709 self.inner.await_events.resolve(key, resolution).await
710 }
711
712 async fn peek_await_event(
713 &self,
714 key: &AwaitEventKey,
715 ) -> Result<Option<Resolution>, RuntimeError> {
716 self.inner.await_events.peek(key).await
717 }
718
719 async fn await_await_event(
720 &self,
721 key: &AwaitEventKey,
722 cancel: CancellationToken,
723 deadline: Option<Instant>,
724 ) -> Result<Resolution, RuntimeError> {
725 self.inner
726 .await_events
727 .await_resolution(key, cancel, deadline)
728 .await
729 }
730
731 async fn revoke_await_events_for_session(&self, session_id: &str) -> Result<(), RuntimeError> {
732 self.inner.await_events.revoke_session(session_id).await
733 }
734
735 async fn cancel_await_events_for_session(&self, session_id: &str) -> Result<(), RuntimeError> {
736 self.inner.await_events.cancel_session(session_id).await
737 }
738}
739
740#[async_trait::async_trait]
741impl RuntimeEffectController for SqliteRuntimeEffectController {
742 async fn execute_effect(
743 &self,
744 envelope: RuntimeEffectEnvelope,
745 local_executor: RuntimeEffectLocalExecutor<'_>,
746 ) -> Result<RuntimeEffectOutcome, RuntimeEffectControllerError> {
747 let reconstructed_envelope = envelope.canonical_form()?;
748 let replay_trace = local_executor.replay_validation_trace().cloned();
749 loop {
750 match self
751 .prepare_effect(&envelope, &reconstructed_envelope)
752 .await?
753 {
754 PreparedEffect::ReplayMismatch {
755 recorded_envelope,
756 stored_envelope_hash,
757 } => {
758 validate_replayed_effect_envelope(
759 recorded_envelope.as_ref(),
760 &reconstructed_envelope,
761 "sqlite_effect_replay_hash_conflict",
762 replay_trace.as_ref(),
763 )?;
764 return Err(RuntimeEffectControllerError::new(
765 "runtime_effect_envelope_canonical_hash_invariant",
766 format!(
767 "stored envelope_hash {stored_envelope_hash} did not match the persisted canonical envelope hash {}",
768 recorded_envelope.hash()
769 ),
770 ));
771 }
772 PreparedEffect::ReplayOutcome { outcome, due_at_ms } => {
773 sleep_until_due(self.inner.clock.as_ref(), due_at_ms).await;
774 return Ok(*outcome);
775 }
776 PreparedEffect::ReplayError(err) => return Err(err),
777 PreparedEffect::Claimed(claim) => {
778 let result = self
779 .execute_claimed_effect_with_renewal(&claim, envelope, local_executor)
780 .await;
781 let finalize = self.finalize_effect(&claim, &result).await;
782 return match (result, finalize) {
783 (Ok(outcome), Ok(())) => Ok(outcome),
784 (Err(err), Ok(())) => Err(err),
785 (_, Err(err)) => Err(err),
786 };
787 }
788 PreparedEffect::Busy { retry_at_ms } => {
789 sleep_until_retry(self.inner.clock.as_ref(), retry_at_ms).await;
790 }
791 }
792 }
793 }
794}
795
796async fn open_effect_replay_inner(
797 path: &Path,
798 backing: StoreBacking,
799 options: SqliteEffectReplayOptions,
800 clock: Arc<dyn lash_core::Clock>,
801) -> tokio_rusqlite::Result<Arc<SqliteEffectReplayInner>> {
802 let conn = SqliteConnection::open(path).await?;
803 let signing_secret = ensure_effect_schema(&conn).await?;
804 apply_pragmas(&conn, backing).await?;
805 Ok(Arc::new(SqliteEffectReplayInner::new(
806 conn,
807 options,
808 clock,
809 signing_secret,
810 )))
811}
812
813async fn open_effect_replay_memory_inner(
814 options: SqliteEffectReplayOptions,
815 clock: Arc<dyn lash_core::Clock>,
816) -> tokio_rusqlite::Result<Arc<SqliteEffectReplayInner>> {
817 let conn = SqliteConnection::open_in_memory().await?;
818 let signing_secret = ensure_effect_schema(&conn).await?;
819 apply_pragmas(&conn, StoreBacking::Memory).await?;
820 Ok(Arc::new(SqliteEffectReplayInner::new(
821 conn,
822 options,
823 clock,
824 signing_secret,
825 )))
826}
827
828impl SqliteEffectReplayInner {
829 fn new(
830 conn: SqliteConnection,
831 options: SqliteEffectReplayOptions,
832 clock: Arc<dyn lash_core::Clock>,
833 signing_secret: Vec<u8>,
834 ) -> Self {
835 let sequence = EFFECT_OWNER_COUNTER.fetch_add(1, Ordering::SeqCst);
836 let timestamp_ms = clock.timestamp_ms();
837 let await_events = SqliteAwaitEvents::new(conn.clone(), signing_secret, Arc::clone(&clock));
838 Self {
839 conn,
840 clock,
841 owner_id: format!("pid{}-{sequence}-{}", std::process::id(), timestamp_ms),
842 lease_counter: AtomicU64::new(1),
843 replay_mode: AtomicBool::new(false),
844 lease_timings: options.lease_timings,
845 await_events,
846 }
847 }
848
849 fn next_lease_token(&self) -> String {
850 let sequence = self.lease_counter.fetch_add(1, Ordering::SeqCst);
851 format!("{}:{sequence}", self.owner_id)
852 }
853}
854
855fn sleep_due_at_ms(envelope: &RuntimeEffectEnvelope, now: u64) -> Option<u64> {
856 match envelope.command {
857 RuntimeEffectCommand::Sleep { duration_ms } => Some(now.saturating_add(duration_ms)),
858 _ => None,
859 }
860}
861
862async fn sleep_until_due(clock: &dyn lash_core::Clock, due_at_ms: Option<u64>) {
863 let Some(due_at_ms) = due_at_ms else {
864 return;
865 };
866 let now = clock.timestamp_ms();
867 if due_at_ms > now {
868 clock.sleep(Duration::from_millis(due_at_ms - now)).await;
869 }
870}
871
872async fn sleep_until_retry(clock: &dyn lash_core::Clock, retry_at_ms: u64) {
873 let now = clock.timestamp_ms();
874 let delay = if retry_at_ms > now {
875 Duration::from_millis(retry_at_ms - now).min(BUSY_POLL)
876 } else {
877 BUSY_POLL
878 };
879 clock.sleep(delay).await;
880}
881
882fn effect_sqlite_error(err: rusqlite::Error) -> RuntimeEffectControllerError {
883 RuntimeEffectControllerError::new("sqlite_effect_replay_store", err.to_string())
884}
885
886fn effect_encode_error(err: serde_json::Error) -> RuntimeEffectControllerError {
887 RuntimeEffectControllerError::new(
888 "sqlite_effect_replay_encode",
889 format!("failed to encode runtime effect replay row: {err}"),
890 )
891}
892
893fn effect_decode_error(err: serde_json::Error) -> RuntimeEffectControllerError {
894 RuntimeEffectControllerError::new(
895 "sqlite_effect_replay_decode",
896 format!("failed to decode runtime effect replay row: {err}"),
897 )
898}