Skip to main content

lash_sqlite_store/
effect_replay.rs

1//! SQLite-backed runtime effect replay host (tokio-rusqlite port of the the prior store
2//! store's `effect_replay`).
3//!
4//! The public surface is identical to the prior implementation — same type names
5//! (`SqliteEffectHost`, `SqliteRuntimeEffectController`, `SqliteEffectReplayOptions`),
6//! same async signatures — so consumers swap backends with a path rename only.
7//! The only mechanical change is the database layer: the prior store ran every op directly
8//! on `&Connection` with `.await`; here every database body is a *synchronous*
9//! rusqlite closure handed to the shared [`SqliteConnection`] wrapper. The
10//! claim/finalize paths run inside `conn.write` (`BEGIN IMMEDIATE`) so the
11//! cross-process write lock is taken up front — the lease fence WAL gives us.
12
13use 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/// Options for SQLite-backed runtime effect replay.
38#[derive(Clone, Debug, Default)]
39pub struct SqliteEffectReplayOptions {
40    /// Effect-replay lease timing capability. Hosts share the same
41    /// [`LeaseTimings`] they configure on the runtime so effect leases expire
42    /// on the same failover window as session and process leases.
43    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/// Deployment-level SQLite effect host.
57///
58/// This host persists runtime effect history in a local SQLite database and
59/// returns scoped controllers that replay completed outcomes by
60/// `(scope_id, replay_key)`.
61#[derive(Clone)]
62pub struct SqliteEffectHost {
63    inner: Arc<SqliteEffectReplayInner>,
64}
65
66/// Scoped SQLite-backed runtime effect controller.
67#[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    /// Force strict replay mode: missing effect history fails instead of
152    /// executing locally. Normal operation still replays any completed row.
153    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    /// Force strict replay mode: missing effect history fails instead of
301    /// executing locally. Normal operation still replays any completed row.
302    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        // The `BEGIN IMMEDIATE` transaction is run on the connection thread via
334        // `conn.write`. The closure returns a `rusqlite::Result` carrying our
335        // own outcome `Result<PreparedEffect, RuntimeEffectControllerError>`:
336        // the SQL committing is independent of whether the recorded effect was
337        // a success, a failure, or a fresh claim — exactly the prior behaviour
338        // (commit on `Ok(_)` for any `PreparedEffect` variant). Only a real
339        // SQLite error rolls back.
340        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}