Skip to main content

meerkat_runtime/store/
memory.rs

1//! InMemoryRuntimeStore — in-memory implementation for testing/ephemeral.
2//!
3//! Uses `tokio::sync::Mutex` per the in-memory concurrency rule.
4//! All mutations complete inside one lock acquisition (no lock held across .await).
5
6use std::collections::{HashMap, HashSet};
7use std::sync::Arc;
8use std::sync::Mutex as StdMutex;
9
10use indexmap::IndexMap;
11use meerkat_core::lifecycle::{InputId, RunBoundaryReceipt, RunId};
12#[cfg(not(target_arch = "wasm32"))]
13use tokio::sync::Mutex;
14#[cfg(target_arch = "wasm32")]
15use tokio_with_wasm::alias::sync::Mutex;
16
17use super::{
18    AuthOAuthFlowSnapshotUpdate, MachineLifecycleCommit, MachineLifecycleSnapshot,
19    MachineLifecycleStoreRecord, RuntimeStore, RuntimeStoreError, SessionDelta,
20};
21use crate::identifiers::LogicalRuntimeId;
22use crate::input_state::{InputStatePersistenceRecord, StoredInputState};
23use crate::ops_lifecycle::PersistedOpsSnapshot;
24
25/// Receipt key: (runtime_id, run_id, sequence).
26#[derive(Debug, Clone, PartialEq, Eq, Hash)]
27struct ReceiptKey {
28    runtime_id: String,
29    run_id: RunId,
30    sequence: u64,
31}
32
33#[derive(Debug, Clone)]
34struct CompactionOutboxEntry {
35    intent: meerkat_core::CompactionProjectionIntent,
36    finalized: bool,
37}
38
39/// Inner state protected by the mutex.
40#[derive(Debug, Default)]
41struct Inner {
42    /// runtime_id → (input_id → StoredInputState). IndexMap for deterministic iteration order.
43    input_states: HashMap<String, IndexMap<InputId, StoredInputState>>,
44    /// Receipt storage.
45    receipts: HashMap<ReceiptKey, RunBoundaryReceipt>,
46    /// Runtime session snapshots keyed by canonical runtime id.
47    sessions: HashMap<String, Vec<u8>>,
48    /// Canonical runtime ids whose projection fallback is quarantined.
49    ///
50    /// Mirrors the durable SQLite `runtime_projection_quarantine` table: set
51    /// when a rejected runtime snapshot is cleared via
52    /// `clear_session_snapshot_if_current`, cleared whenever a live snapshot is
53    /// written for the runtime.
54    projection_quarantine: HashSet<String>,
55    /// Persisted machine lifecycle snapshots.
56    runtime_lifecycle: HashMap<String, MachineLifecycleSnapshot>,
57    /// Persisted ops lifecycle snapshots.
58    ops_lifecycle_snapshots: HashMap<String, PersistedOpsSnapshot>,
59    /// Exact ops epochs retired by atomic unregister finalization. Tombstones
60    /// outlive row deletion so detached callbacks cannot resurrect them.
61    retired_ops_epochs: HashSet<(String, meerkat_core::RuntimeEpochId)>,
62    /// Runtime id -> transcript-rewrite-keyed compaction projection outbox.
63    compaction_projection_outbox:
64        HashMap<String, HashMap<meerkat_core::CompactionProjectionId, CompactionOutboxEntry>>,
65}
66
67/// In-memory runtime store. Thread-safe via `tokio::sync::Mutex`.
68#[derive(Debug, Clone)]
69pub struct InMemoryRuntimeStore {
70    inner: Arc<Mutex<Inner>>,
71    auth_oauth_flow_snapshot: Arc<StdMutex<Option<Vec<u8>>>>,
72}
73
74impl InMemoryRuntimeStore {
75    pub fn new() -> Self {
76        Self {
77            inner: Arc::new(Mutex::new(Inner::default())),
78            auth_oauth_flow_snapshot: Arc::new(StdMutex::new(None)),
79        }
80    }
81}
82
83impl Default for InMemoryRuntimeStore {
84    fn default() -> Self {
85        Self::new()
86    }
87}
88
89fn is_runtime_placeholder_session(session: &meerkat_core::Session) -> bool {
90    session.transcript_history_state().ok().flatten().is_none()
91        && matches!(
92            session.messages(),
93            [] | [meerkat_core::types::Message::System(_)]
94        )
95}
96
97/// Deserialize a persisted session-snapshot blob through typed serde, matching
98/// the SQLite runtime store read path. `Session::deserialize` validates the
99/// mandatory envelope version against the generated persistence version
100/// authority, so a missing or non-current (v0/v1) row fails closed instead of
101/// silently defaulting or upgrading on read.
102fn deserialize_persisted_session(bytes: &[u8]) -> Result<meerkat_core::Session, RuntimeStoreError> {
103    serde_json::from_slice(bytes).map_err(|err| RuntimeStoreError::ReadFailed(err.to_string()))
104}
105
106fn ensure_compaction_intents_already_outboxed(
107    inner: &Inner,
108    runtime_id: &LogicalRuntimeId,
109    session: &meerkat_core::Session,
110) -> Result<(), RuntimeStoreError> {
111    let intents = super::validated_compaction_projection_intents(session)?;
112    let existing = inner.compaction_projection_outbox.get(&runtime_id.0);
113    for intent in intents {
114        match existing.and_then(|entries| entries.get(&intent.projection)) {
115            Some(entry) if entry.finalized => {
116                return Err(RuntimeStoreError::WriteFailed(format!(
117                    "non-boundary snapshot replays finalized compaction intent {}",
118                    intent.projection.revision()
119                )));
120            }
121            Some(entry) if entry.intent == intent => {}
122            Some(_) => {
123                return Err(RuntimeStoreError::WriteFailed(format!(
124                    "non-boundary snapshot conflicts with compaction outbox rewrite {}",
125                    intent.projection.revision()
126                )));
127            }
128            None => {
129                return Err(RuntimeStoreError::WriteFailed(format!(
130                    "non-boundary snapshot introduces compaction intent {} without atomic outbox authority",
131                    intent.projection.revision()
132                )));
133            }
134        }
135    }
136    Ok(())
137}
138
139#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
140#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
141impl RuntimeStore for InMemoryRuntimeStore {
142    fn supports_compaction_projection_outbox(&self) -> bool {
143        true
144    }
145
146    fn persist_auth_oauth_flow_snapshot(
147        &self,
148        snapshot_json: &[u8],
149    ) -> Result<(), RuntimeStoreError> {
150        *self
151            .auth_oauth_flow_snapshot
152            .lock()
153            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))? =
154            Some(snapshot_json.to_vec());
155        Ok(())
156    }
157
158    fn load_auth_oauth_flow_snapshot(&self) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
159        self.auth_oauth_flow_snapshot
160            .lock()
161            .map(|snapshot| snapshot.clone())
162            .map_err(|err| RuntimeStoreError::ReadFailed(err.to_string()))
163    }
164
165    fn update_auth_oauth_flow_snapshot(
166        &self,
167        update: &mut AuthOAuthFlowSnapshotUpdate<'_>,
168    ) -> Result<(), RuntimeStoreError> {
169        let mut snapshot = self
170            .auth_oauth_flow_snapshot
171            .lock()
172            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
173        let next = update(snapshot.as_deref())?;
174        *snapshot = Some(next);
175        Ok(())
176    }
177
178    async fn commit_session_snapshot(
179        &self,
180        runtime_id: &LogicalRuntimeId,
181        session_delta: SessionDelta,
182    ) -> Result<(), RuntimeStoreError> {
183        let incoming: meerkat_core::Session =
184            serde_json::from_slice(&session_delta.session_snapshot)
185                .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
186        let mut inner = self.inner.lock().await;
187        ensure_compaction_intents_already_outboxed(&inner, runtime_id, &incoming)?;
188        let previous = inner
189            .sessions
190            .get(&runtime_id.0)
191            .map(|snapshot| deserialize_persisted_session(snapshot))
192            .transpose()?;
193        meerkat_core::session_store::run_boundary_snapshot_save_guard(&incoming, previous.as_ref())
194            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
195        inner
196            .sessions
197            .insert(runtime_id.0.clone(), session_delta.session_snapshot);
198        inner.projection_quarantine.remove(&runtime_id.0);
199        Ok(())
200    }
201
202    async fn commit_session_transcript_rewrite_snapshot(
203        &self,
204        runtime_id: &LogicalRuntimeId,
205        session_delta: SessionDelta,
206        commit: &meerkat_core::TranscriptRewriteCommit,
207    ) -> Result<(), RuntimeStoreError> {
208        let incoming: meerkat_core::Session =
209            serde_json::from_slice(&session_delta.session_snapshot)
210                .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
211        let mut inner = self.inner.lock().await;
212        ensure_compaction_intents_already_outboxed(&inner, runtime_id, &incoming)?;
213        let previous = inner
214            .sessions
215            .get(&runtime_id.0)
216            .map(|snapshot| deserialize_persisted_session(snapshot))
217            .transpose()?;
218        meerkat_core::session_store::transcript_rewrite_save_guard(
219            &incoming,
220            previous.as_ref(),
221            commit,
222        )
223        .map_err(|err| match err {
224            meerkat_core::SessionStoreError::TranscriptRevisionConflict {
225                expected,
226                actual,
227                ..
228            } => RuntimeStoreError::TranscriptRevisionConflict { expected, actual },
229            other => RuntimeStoreError::WriteFailed(other.to_string()),
230        })?;
231        inner
232            .sessions
233            .insert(runtime_id.0.clone(), session_delta.session_snapshot);
234        inner.projection_quarantine.remove(&runtime_id.0);
235        Ok(())
236    }
237
238    async fn atomic_apply(
239        &self,
240        runtime_id: &LogicalRuntimeId,
241        session_delta: Option<SessionDelta>,
242        receipt: RunBoundaryReceipt,
243        input_updates: Vec<InputStatePersistenceRecord>,
244        session_store_key: Option<meerkat_core::types::SessionId>,
245    ) -> Result<(), RuntimeStoreError> {
246        let mut inner = self.inner.lock().await;
247
248        // All writes in one lock acquisition (atomic for in-memory)
249        let rid = runtime_id.0.clone();
250
251        // Session delta. The supersession verdict computed here keys the
252        // entire commit: if the incoming session snapshot is classified as
253        // superseded (the persisted head is already a valid append-extension
254        // of it), the snapshot write is skipped AND so are the receipt + input
255        // writes, so receipt/input ordering identity never advances past the
256        // retained session truth.
257        let mut session_snapshot_superseded = false;
258        let mut compaction_intents = Vec::new();
259        if let Some(delta) = session_delta {
260            let incoming_session =
261                serde_json::from_slice::<meerkat_core::Session>(&delta.session_snapshot);
262            let mut persist_session_snapshot = true;
263            match (incoming_session, session_store_key) {
264                (Ok(incoming_session), session_store_key) => {
265                    compaction_intents =
266                        super::validated_compaction_projection_intents(&incoming_session)?;
267                    if let Some(existing) = inner.compaction_projection_outbox.get(&rid) {
268                        for intent in &compaction_intents {
269                            if let Some(entry) = existing.get(&intent.projection) {
270                                if entry.finalized {
271                                    return Err(RuntimeStoreError::WriteFailed(format!(
272                                        "atomic session snapshot replays finalized compaction intent {}",
273                                        intent.projection.revision()
274                                    )));
275                                }
276                                if entry.intent != *intent {
277                                    return Err(RuntimeStoreError::WriteFailed(format!(
278                                        "conflicting compaction outbox intent for rewrite {}",
279                                        intent.projection.revision()
280                                    )));
281                                }
282                            }
283                        }
284                    }
285                    if let Some(session_store_key) = session_store_key
286                        && incoming_session.id() != &session_store_key
287                    {
288                        return Err(RuntimeStoreError::SessionKeyMismatch {
289                            expected: session_store_key,
290                            actual: incoming_session.id().clone(),
291                        });
292                    }
293                    let previous_session = inner
294                        .sessions
295                        .get(&rid)
296                        .and_then(|snapshot| deserialize_persisted_session(snapshot).ok());
297                    if let Err(err) = meerkat_core::session_store::run_boundary_snapshot_save_guard(
298                        &incoming_session,
299                        previous_session.as_ref(),
300                    ) {
301                        if previous_session
302                            .as_ref()
303                            .is_some_and(is_runtime_placeholder_session)
304                        {
305                            persist_session_snapshot = true;
306                        } else if previous_session.as_ref().is_some_and(|previous_session| {
307                            meerkat_core::session_store::run_boundary_snapshot_save_guard(
308                                previous_session,
309                                Some(&incoming_session),
310                            )
311                            .is_ok()
312                        }) {
313                            persist_session_snapshot = false;
314                            session_snapshot_superseded = true;
315                        } else {
316                            return Err(RuntimeStoreError::WriteFailed(err.to_string()));
317                        }
318                    }
319                }
320                (Err(err), Some(session_store_key)) => {
321                    return Err(RuntimeStoreError::WriteFailed(format!(
322                        "session snapshot for {session_store_key} is not a Session: {err}"
323                    )));
324                }
325                (Err(err), None) => {
326                    return Err(RuntimeStoreError::WriteFailed(format!(
327                        "session snapshot is not a Session: {err}"
328                    )));
329                }
330            }
331            if persist_session_snapshot {
332                inner.sessions.insert(rid.clone(), delta.session_snapshot);
333                inner.projection_quarantine.remove(&rid);
334            }
335        }
336
337        // When the session snapshot was superseded and skipped, the boundary
338        // receipt and input-state updates for that boundary must also be
339        // skipped: advancing them against a retained (older) session snapshot
340        // would split receipt/input ordering identity from session truth.
341        if session_snapshot_superseded {
342            return Ok(());
343        }
344
345        let outbox = inner
346            .compaction_projection_outbox
347            .entry(rid.clone())
348            .or_default();
349        for intent in compaction_intents {
350            outbox
351                .entry(intent.projection.clone())
352                .or_insert(CompactionOutboxEntry {
353                    intent,
354                    finalized: false,
355                });
356        }
357
358        // Receipt
359        let key = ReceiptKey {
360            runtime_id: rid.clone(),
361            run_id: receipt.run_id.clone(),
362            sequence: receipt.sequence,
363        };
364        inner.receipts.insert(key, receipt);
365
366        // Input states
367        let states = inner.input_states.entry(rid).or_default();
368        for record in input_updates {
369            let bundle = record.into_stored();
370            states.insert(bundle.state.input_id.clone(), bundle);
371        }
372
373        Ok(())
374    }
375
376    async fn load_pending_compaction_projections(
377        &self,
378        runtime_id: &LogicalRuntimeId,
379    ) -> Result<Vec<meerkat_core::CompactionProjectionIntent>, RuntimeStoreError> {
380        let inner = self.inner.lock().await;
381        let mut pending = inner
382            .compaction_projection_outbox
383            .get(&runtime_id.0)
384            .into_iter()
385            .flat_map(HashMap::values)
386            .filter(|entry| !entry.finalized)
387            .map(|entry| entry.intent.clone())
388            .collect::<Vec<_>>();
389        pending.sort_by(|left, right| {
390            left.projection
391                .session_id()
392                .to_string()
393                .cmp(&right.projection.session_id().to_string())
394                .then_with(|| {
395                    left.projection
396                        .parent_revision()
397                        .cmp(right.projection.parent_revision())
398                })
399                .then_with(|| left.projection.revision().cmp(right.projection.revision()))
400                .then_with(|| {
401                    left.projection
402                        .commit_fingerprint()
403                        .cmp(right.projection.commit_fingerprint())
404                })
405        });
406        Ok(pending)
407    }
408
409    async fn mark_compaction_projection_finalized(
410        &self,
411        runtime_id: &LogicalRuntimeId,
412        projection: &meerkat_core::CompactionProjectionId,
413    ) -> Result<(), RuntimeStoreError> {
414        let mut inner = self.inner.lock().await;
415        let outbox_exists = inner
416            .compaction_projection_outbox
417            .get(&runtime_id.0)
418            .is_some_and(|entries| entries.contains_key(projection));
419        if !outbox_exists {
420            return Err(RuntimeStoreError::NotFound(format!(
421                "compaction outbox rewrite {}",
422                projection.revision()
423            )));
424        }
425        let cleaned_snapshot = inner
426            .sessions
427            .get(&runtime_id.0)
428            .map(|snapshot| {
429                let mut session = deserialize_persisted_session(snapshot)?;
430                session
431                    .complete_compaction_projection_intent(projection)
432                    .map_err(|error| RuntimeStoreError::WriteFailed(error.to_string()))?;
433                serde_json::to_vec(&session)
434                    .map_err(|error| RuntimeStoreError::WriteFailed(error.to_string()))
435            })
436            .transpose()?;
437        let entry = inner
438            .compaction_projection_outbox
439            .get_mut(&runtime_id.0)
440            .and_then(|entries| entries.get_mut(projection))
441            .ok_or_else(|| {
442                RuntimeStoreError::NotFound(format!(
443                    "compaction outbox rewrite {}",
444                    projection.revision()
445                ))
446            })?;
447        entry.finalized = true;
448        if let Some(cleaned_snapshot) = cleaned_snapshot {
449            inner
450                .sessions
451                .insert(runtime_id.0.clone(), cleaned_snapshot);
452        }
453        Ok(())
454    }
455
456    async fn load_input_states(
457        &self,
458        runtime_id: &LogicalRuntimeId,
459    ) -> Result<Vec<StoredInputState>, RuntimeStoreError> {
460        let inner = self.inner.lock().await;
461        let states = inner
462            .input_states
463            .get(&runtime_id.0)
464            .map(|m| m.values().cloned().collect())
465            .unwrap_or_default();
466        Ok(states)
467    }
468
469    async fn load_boundary_receipt(
470        &self,
471        runtime_id: &LogicalRuntimeId,
472        run_id: &RunId,
473        sequence: u64,
474    ) -> Result<Option<RunBoundaryReceipt>, RuntimeStoreError> {
475        let inner = self.inner.lock().await;
476        let key = ReceiptKey {
477            runtime_id: runtime_id.0.clone(),
478            run_id: run_id.clone(),
479            sequence,
480        };
481        Ok(inner.receipts.get(&key).cloned())
482    }
483
484    async fn load_session_snapshot(
485        &self,
486        runtime_id: &LogicalRuntimeId,
487    ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
488        let inner = self.inner.lock().await;
489        Ok(inner.sessions.get(&runtime_id.0).cloned())
490    }
491
492    async fn clear_session_snapshot(
493        &self,
494        runtime_id: &LogicalRuntimeId,
495    ) -> Result<(), RuntimeStoreError> {
496        let mut inner = self.inner.lock().await;
497        inner.sessions.remove(&runtime_id.0);
498        Ok(())
499    }
500
501    async fn replace_session_snapshot_if_current(
502        &self,
503        runtime_id: &LogicalRuntimeId,
504        expected_current: &[u8],
505        replacement: Vec<u8>,
506    ) -> Result<bool, RuntimeStoreError> {
507        let replacement_session: meerkat_core::Session = serde_json::from_slice(&replacement)
508            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
509        let mut inner = self.inner.lock().await;
510        let Some(current) = inner.sessions.get(&runtime_id.0) else {
511            return Ok(false);
512        };
513        if current.as_slice() != expected_current {
514            return Ok(false);
515        }
516        ensure_compaction_intents_already_outboxed(&inner, runtime_id, &replacement_session)?;
517        inner.sessions.insert(runtime_id.0.clone(), replacement);
518        inner.projection_quarantine.remove(&runtime_id.0);
519        Ok(true)
520    }
521
522    async fn clear_session_snapshot_if_current(
523        &self,
524        runtime_id: &LogicalRuntimeId,
525        expected_current: &[u8],
526    ) -> Result<bool, RuntimeStoreError> {
527        let mut inner = self.inner.lock().await;
528        let Some(current) = inner.sessions.get(&runtime_id.0) else {
529            return Ok(false);
530        };
531        if current.as_slice() != expected_current {
532            return Ok(false);
533        }
534        inner.sessions.remove(&runtime_id.0);
535        // Record the in-memory quarantine marker atomically with the snapshot
536        // removal, mirroring the durable SQLite path.
537        inner.projection_quarantine.insert(runtime_id.0.clone());
538        Ok(true)
539    }
540
541    async fn is_runtime_projection_quarantined(
542        &self,
543        runtime_id: &LogicalRuntimeId,
544    ) -> Result<bool, RuntimeStoreError> {
545        let inner = self.inner.lock().await;
546        Ok(inner.projection_quarantine.contains(&runtime_id.0))
547    }
548
549    async fn persist_input_state(
550        &self,
551        runtime_id: &LogicalRuntimeId,
552        state: &InputStatePersistenceRecord,
553    ) -> Result<(), RuntimeStoreError> {
554        let mut inner = self.inner.lock().await;
555        let states = inner.input_states.entry(runtime_id.0.clone()).or_default();
556        let bundle = state.as_stored();
557        states.insert(bundle.state.input_id.clone(), bundle.clone());
558        Ok(())
559    }
560
561    async fn load_input_state(
562        &self,
563        runtime_id: &LogicalRuntimeId,
564        input_id: &InputId,
565    ) -> Result<Option<StoredInputState>, RuntimeStoreError> {
566        let inner = self.inner.lock().await;
567        let state = inner
568            .input_states
569            .get(&runtime_id.0)
570            .and_then(|m| m.get(input_id).cloned());
571        Ok(state)
572    }
573
574    async fn load_machine_lifecycle_record(
575        &self,
576        runtime_id: &LogicalRuntimeId,
577    ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
578        let inner = self.inner.lock().await;
579        inner
580            .runtime_lifecycle
581            .get(&runtime_id.0)
582            .map(|snapshot| MachineLifecycleStoreRecord::from_snapshot(snapshot).encode())
583            .transpose()
584    }
585
586    async fn commit_machine_lifecycle(
587        &self,
588        runtime_id: &LogicalRuntimeId,
589        commit: MachineLifecycleCommit,
590        input_states: &[InputStatePersistenceRecord],
591    ) -> Result<(), RuntimeStoreError> {
592        let mut inner = self.inner.lock().await;
593        let rid = runtime_id.0.clone();
594
595        // Single lock acquisition — atomic for in-memory
596        inner
597            .runtime_lifecycle
598            .insert(rid.clone(), commit.into_snapshot());
599        let states = inner.input_states.entry(rid).or_default();
600        for record in input_states {
601            let bundle = record.as_stored();
602            states.insert(bundle.state.input_id.clone(), bundle.clone());
603        }
604
605        Ok(())
606    }
607
608    async fn commit_unregister_finalization(
609        &self,
610        runtime_id: &LogicalRuntimeId,
611        commit: MachineLifecycleCommit,
612        input_states: &[InputStatePersistenceRecord],
613    ) -> Result<(), RuntimeStoreError> {
614        let mut inner = self.inner.lock().await;
615        let rid = runtime_id.0.clone();
616        let retired_ops_epoch = commit.retired_ops_epoch().cloned().ok_or_else(|| {
617            RuntimeStoreError::WriteFailed(
618                "unregister finalization missing exact retired ops epoch".into(),
619            )
620        })?;
621
622        // One critical section publishes the terminal lifecycle and removes
623        // the old ops epoch while recording a deletion-wins tombstone. No
624        // observer can acquire the store between them.
625        inner
626            .runtime_lifecycle
627            .insert(rid.clone(), commit.into_snapshot());
628        let states = inner.input_states.entry(rid.clone()).or_default();
629        for record in input_states {
630            let bundle = record.as_stored();
631            states.insert(bundle.state.input_id.clone(), bundle.clone());
632        }
633        if inner
634            .ops_lifecycle_snapshots
635            .get(&rid)
636            .is_some_and(|snapshot| snapshot.epoch_id == retired_ops_epoch)
637        {
638            inner.ops_lifecycle_snapshots.remove(&rid);
639        }
640        inner.retired_ops_epochs.insert((rid, retired_ops_epoch));
641        Ok(())
642    }
643
644    async fn persist_ops_lifecycle(
645        &self,
646        runtime_id: &LogicalRuntimeId,
647        snapshot: &PersistedOpsSnapshot,
648    ) -> Result<(), RuntimeStoreError> {
649        let mut inner = self.inner.lock().await;
650        if inner
651            .retired_ops_epochs
652            .contains(&(runtime_id.0.clone(), snapshot.epoch_id.clone()))
653        {
654            return Err(RuntimeStoreError::OpsLifecycleEpochRetired {
655                runtime_id: runtime_id.0.clone(),
656                epoch_id: snapshot.epoch_id.clone(),
657            });
658        }
659        inner
660            .ops_lifecycle_snapshots
661            .insert(runtime_id.0.clone(), snapshot.clone());
662        Ok(())
663    }
664
665    async fn load_ops_lifecycle(
666        &self,
667        runtime_id: &LogicalRuntimeId,
668    ) -> Result<Option<PersistedOpsSnapshot>, RuntimeStoreError> {
669        let inner = self.inner.lock().await;
670        Ok(inner.ops_lifecycle_snapshots.get(&runtime_id.0).cloned())
671    }
672
673    async fn delete_ops_lifecycle(
674        &self,
675        runtime_id: &LogicalRuntimeId,
676    ) -> Result<(), RuntimeStoreError> {
677        let mut inner = self.inner.lock().await;
678        inner.ops_lifecycle_snapshots.remove(&runtime_id.0);
679        Ok(())
680    }
681}
682
683#[cfg(test)]
684#[allow(clippy::unwrap_used)]
685mod tests {
686    use super::*;
687    use crate::RuntimeState;
688    use crate::store::MachineLifecycleBindingFacts;
689    use meerkat_core::lifecycle::run_primitive::RunApplyBoundary;
690
691    fn make_receipt(run_id: RunId, seq: u64) -> RunBoundaryReceipt {
692        RunBoundaryReceipt {
693            run_id,
694            boundary: RunApplyBoundary::RunStart,
695            contributing_input_ids: vec![],
696            conversation_digest: None,
697            message_count: 0,
698            sequence: seq,
699        }
700    }
701
702    fn persistable(bundle: StoredInputState) -> InputStatePersistenceRecord {
703        InputStatePersistenceRecord::from_machine_snapshot(bundle).unwrap()
704    }
705
706    fn session_with_user(content: &str) -> meerkat_core::Session {
707        let mut session = meerkat_core::Session::new();
708        session.push(meerkat_core::types::Message::User(
709            meerkat_core::types::UserMessage::text(content.to_string()),
710        ));
711        session
712    }
713
714    fn session_with_compaction_intent() -> (
715        meerkat_core::Session,
716        meerkat_core::CompactionProjectionIntent,
717    ) {
718        let mut session = session_with_user("verbose context one");
719        session.push(meerkat_core::types::Message::User(
720            meerkat_core::types::UserMessage::text("verbose context two"),
721        ));
722        let parent = session.transcript_revision().unwrap();
723        session
724            .commit_transcript_rewrite(
725                meerkat_core::TranscriptRewriteSelection::MessageRange { start: 0, end: 2 },
726                vec![meerkat_core::types::Message::User(
727                    meerkat_core::types::UserMessage::compaction_summary("compacted context"),
728                )],
729                meerkat_core::TranscriptRewriteReason::new("compaction"),
730                Some("runtime-store-test".to_string()),
731                Some(parent),
732            )
733            .unwrap();
734        let mut encoded = serde_json::to_value(&session).unwrap();
735        encoded["metadata"][meerkat_core::SESSION_TRANSCRIPT_HISTORY_STATE_KEY]["commits"][0]["selection"] = serde_json::json!({
736            "type": "compaction_message_range",
737            "range": { "start": 0, "end": 2 }
738        });
739        let mut session: meerkat_core::Session = serde_json::from_value(encoded).unwrap();
740        let commit = session
741            .transcript_history_state()
742            .unwrap()
743            .unwrap()
744            .commits
745            .last()
746            .unwrap()
747            .clone();
748        let intent = meerkat_core::CompactionProjectionIntent {
749            projection: serde_json::from_value(serde_json::json!({
750                "session_id": session.id(),
751                "parent_revision": &commit.parent_revision,
752                "revision": &commit.revision,
753                "commit_fingerprint": "sha256:827d8ee5666e51b2ced4d303640740680d96151d92187fd6e981c29550072c62",
754            }))
755            .unwrap(),
756            summary_tokens: 5,
757            messages_before: 2,
758            messages_after: 1,
759        };
760        session
761            .add_compaction_projection_intent(intent.clone())
762            .unwrap();
763        (session, intent)
764    }
765
766    fn snapshot_with_raw_intents(
767        session: &meerkat_core::Session,
768        intents: &[meerkat_core::CompactionProjectionIntent],
769    ) -> Vec<u8> {
770        let mut value = serde_json::to_value(session).unwrap();
771        value["metadata"][meerkat_core::memory::SESSION_COMPACTION_PROJECTION_INTENTS_KEY] =
772            serde_json::to_value(intents).unwrap();
773        serde_json::to_vec(&value).unwrap()
774    }
775
776    fn unbacked_intent(
777        session_id: &meerkat_core::types::SessionId,
778    ) -> meerkat_core::CompactionProjectionIntent {
779        meerkat_core::CompactionProjectionIntent {
780            projection: serde_json::from_value(serde_json::json!({
781                "session_id": session_id,
782                "parent_revision": "missing-parent",
783                "revision": "missing-revision",
784                "commit_fingerprint": "sha256:unbacked-persisted-fixture",
785            }))
786            .unwrap(),
787            summary_tokens: 1,
788            messages_before: 2,
789            messages_after: 1,
790        }
791    }
792
793    #[tokio::test]
794    async fn atomic_apply_commits_rewrite_and_compaction_outbox_as_one_boundary() {
795        let store = InMemoryRuntimeStore::new();
796        let rid = LogicalRuntimeId::new("runtime-compaction-outbox");
797        let (session, intent) = session_with_compaction_intent();
798        let snapshot = serde_json::to_vec(&session).unwrap();
799        store
800            .atomic_apply(
801                &rid,
802                Some(SessionDelta {
803                    session_snapshot: snapshot.clone(),
804                }),
805                make_receipt(RunId::new(), 1),
806                vec![],
807                Some(session.id().clone()),
808            )
809            .await
810            .unwrap();
811        assert_eq!(
812            store.load_session_snapshot(&rid).await.unwrap(),
813            Some(snapshot)
814        );
815        assert_eq!(
816            store
817                .load_pending_compaction_projections(&rid)
818                .await
819                .unwrap(),
820            vec![intent.clone()]
821        );
822        store
823            .mark_compaction_projection_finalized(&rid, &intent.projection)
824            .await
825            .unwrap();
826        store
827            .mark_compaction_projection_finalized(&rid, &intent.projection)
828            .await
829            .unwrap();
830        assert!(
831            store
832                .load_pending_compaction_projections(&rid)
833                .await
834                .unwrap()
835                .is_empty()
836        );
837        let persisted: meerkat_core::Session =
838            serde_json::from_slice(&store.load_session_snapshot(&rid).await.unwrap().unwrap())
839                .unwrap();
840        assert!(
841            persisted
842                .compaction_projection_intents()
843                .unwrap()
844                .is_empty()
845        );
846    }
847
848    #[tokio::test]
849    async fn finalized_outbox_tombstone_rejects_atomic_and_non_boundary_snapshot_replay() {
850        let store = InMemoryRuntimeStore::new();
851        let rid = LogicalRuntimeId::new("runtime-finalized-compaction-replay");
852        let (session, intent) = session_with_compaction_intent();
853        let replay_snapshot = serde_json::to_vec(&session).unwrap();
854        let commit = session
855            .transcript_history_state()
856            .unwrap()
857            .unwrap()
858            .commits
859            .last()
860            .unwrap()
861            .clone();
862        store
863            .atomic_apply(
864                &rid,
865                Some(SessionDelta {
866                    session_snapshot: replay_snapshot.clone(),
867                }),
868                make_receipt(RunId::new(), 1),
869                vec![],
870                Some(session.id().clone()),
871            )
872            .await
873            .unwrap();
874        store
875            .mark_compaction_projection_finalized(&rid, &intent.projection)
876            .await
877            .unwrap();
878        let cleaned_snapshot = store.load_session_snapshot(&rid).await.unwrap().unwrap();
879
880        let replay_run_id = RunId::new();
881        let error = store
882            .atomic_apply(
883                &rid,
884                Some(SessionDelta {
885                    session_snapshot: replay_snapshot.clone(),
886                }),
887                make_receipt(replay_run_id.clone(), 2),
888                vec![],
889                Some(session.id().clone()),
890            )
891            .await
892            .unwrap_err();
893        assert!(error.to_string().contains("finalized compaction intent"));
894        assert!(
895            store
896                .load_boundary_receipt(&rid, &replay_run_id, 2)
897                .await
898                .unwrap()
899                .is_none(),
900            "finalized replay rejection must roll back the whole atomic boundary"
901        );
902
903        let error = store
904            .commit_session_snapshot(
905                &rid,
906                SessionDelta {
907                    session_snapshot: replay_snapshot.clone(),
908                },
909            )
910            .await
911            .unwrap_err();
912        assert!(error.to_string().contains("finalized compaction intent"));
913        let error = store
914            .commit_session_transcript_rewrite_snapshot(
915                &rid,
916                SessionDelta {
917                    session_snapshot: replay_snapshot.clone(),
918                },
919                &commit,
920            )
921            .await
922            .unwrap_err();
923        assert!(error.to_string().contains("finalized compaction intent"));
924        let error = store
925            .replace_session_snapshot_if_current(&rid, &cleaned_snapshot, replay_snapshot)
926            .await
927            .unwrap_err();
928        assert!(error.to_string().contains("finalized compaction intent"));
929
930        assert_eq!(
931            store.load_session_snapshot(&rid).await.unwrap(),
932            Some(cleaned_snapshot)
933        );
934        assert!(
935            store
936                .load_pending_compaction_projections(&rid)
937                .await
938                .unwrap()
939                .is_empty(),
940            "a finalized tombstone must never be silently revived or left untracked"
941        );
942    }
943
944    #[tokio::test]
945    async fn invalid_compaction_intent_leaves_snapshot_and_outbox_unmodified() {
946        let store = InMemoryRuntimeStore::new();
947        let rid = LogicalRuntimeId::new("runtime-invalid-compaction-outbox");
948        let (session, mut intent) = session_with_compaction_intent();
949        intent.summary_tokens += 1;
950        let conflicting = vec![
951            session.compaction_projection_intents().unwrap()[0].clone(),
952            intent,
953        ];
954        let error = store
955            .atomic_apply(
956                &rid,
957                Some(SessionDelta {
958                    session_snapshot: snapshot_with_raw_intents(&session, &conflicting),
959                }),
960                make_receipt(RunId::new(), 2),
961                vec![],
962                Some(session.id().clone()),
963            )
964            .await
965            .unwrap_err();
966        assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
967        assert_eq!(store.load_session_snapshot(&rid).await.unwrap(), None);
968        assert!(
969            store
970                .load_pending_compaction_projections(&rid)
971                .await
972                .unwrap()
973                .is_empty()
974        );
975
976        let foreign = session_with_compaction_intent().1;
977        for (sequence, invalid) in [foreign, unbacked_intent(session.id())]
978            .into_iter()
979            .enumerate()
980        {
981            let error = store
982                .atomic_apply(
983                    &rid,
984                    Some(SessionDelta {
985                        session_snapshot: snapshot_with_raw_intents(&session, &[invalid]),
986                    }),
987                    make_receipt(RunId::new(), 10 + sequence as u64),
988                    vec![],
989                    Some(session.id().clone()),
990                )
991                .await
992                .unwrap_err();
993            assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
994            assert_eq!(store.load_session_snapshot(&rid).await.unwrap(), None);
995            assert!(
996                store
997                    .load_pending_compaction_projections(&rid)
998                    .await
999                    .unwrap()
1000                    .is_empty()
1001            );
1002        }
1003    }
1004
1005    #[tokio::test]
1006    async fn superseded_snapshot_does_not_advance_compaction_outbox() {
1007        let store = InMemoryRuntimeStore::new();
1008        let rid = LogicalRuntimeId::new("runtime-superseded-compaction-outbox");
1009        let (incoming, intent) = session_with_compaction_intent();
1010        let mut current = incoming.clone();
1011        current
1012            .complete_compaction_projection_intent(&intent.projection)
1013            .unwrap();
1014        current.push(meerkat_core::types::Message::User(
1015            meerkat_core::types::UserMessage::text("already advanced"),
1016        ));
1017        let current_snapshot = serde_json::to_vec(&current).unwrap();
1018        store
1019            .commit_session_snapshot(
1020                &rid,
1021                SessionDelta {
1022                    session_snapshot: current_snapshot.clone(),
1023                },
1024            )
1025            .await
1026            .unwrap();
1027        store
1028            .atomic_apply(
1029                &rid,
1030                Some(SessionDelta {
1031                    session_snapshot: serde_json::to_vec(&incoming).unwrap(),
1032                }),
1033                make_receipt(RunId::new(), 3),
1034                vec![],
1035                Some(incoming.id().clone()),
1036            )
1037            .await
1038            .unwrap();
1039        assert_eq!(
1040            store.load_session_snapshot(&rid).await.unwrap(),
1041            Some(current_snapshot)
1042        );
1043        assert!(
1044            store
1045                .load_pending_compaction_projections(&rid)
1046                .await
1047                .unwrap()
1048                .is_empty()
1049        );
1050    }
1051
1052    #[tokio::test]
1053    async fn existing_outbox_rejects_changed_intent_without_advancing_snapshot() {
1054        let store = InMemoryRuntimeStore::new();
1055        let rid = LogicalRuntimeId::new("runtime-conflicting-compaction-outbox");
1056        let (session, intent) = session_with_compaction_intent();
1057        let original_snapshot = serde_json::to_vec(&session).unwrap();
1058        store
1059            .atomic_apply(
1060                &rid,
1061                Some(SessionDelta {
1062                    session_snapshot: original_snapshot.clone(),
1063                }),
1064                make_receipt(RunId::new(), 60),
1065                vec![],
1066                Some(session.id().clone()),
1067            )
1068            .await
1069            .unwrap();
1070
1071        let mut advanced = session.clone();
1072        advanced.push(meerkat_core::types::Message::User(
1073            meerkat_core::types::UserMessage::text("later turn"),
1074        ));
1075        let mut conflicting = intent.clone();
1076        conflicting.summary_tokens += 1;
1077        let error = store
1078            .atomic_apply(
1079                &rid,
1080                Some(SessionDelta {
1081                    session_snapshot: snapshot_with_raw_intents(&advanced, &[conflicting]),
1082                }),
1083                make_receipt(RunId::new(), 61),
1084                vec![],
1085                Some(session.id().clone()),
1086            )
1087            .await
1088            .unwrap_err();
1089        assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
1090        assert_eq!(
1091            store.load_session_snapshot(&rid).await.unwrap(),
1092            Some(original_snapshot)
1093        );
1094        assert_eq!(
1095            store
1096                .load_pending_compaction_projections(&rid)
1097                .await
1098                .unwrap(),
1099            vec![intent]
1100        );
1101    }
1102
1103    #[tokio::test]
1104    async fn non_boundary_snapshot_apis_cannot_bypass_compaction_outbox() {
1105        let store = InMemoryRuntimeStore::new();
1106        let rid = LogicalRuntimeId::new("runtime-compaction-bypass");
1107        let (session, _intent) = session_with_compaction_intent();
1108        let snapshot = serde_json::to_vec(&session).unwrap();
1109        let commit = session
1110            .transcript_history_state()
1111            .unwrap()
1112            .unwrap()
1113            .commits
1114            .last()
1115            .unwrap()
1116            .clone();
1117        assert!(
1118            store
1119                .commit_session_snapshot(
1120                    &rid,
1121                    SessionDelta {
1122                        session_snapshot: snapshot.clone(),
1123                    },
1124                )
1125                .await
1126                .is_err()
1127        );
1128        assert!(
1129            store
1130                .commit_session_transcript_rewrite_snapshot(
1131                    &rid,
1132                    SessionDelta {
1133                        session_snapshot: snapshot.clone(),
1134                    },
1135                    &commit,
1136                )
1137                .await
1138                .is_err()
1139        );
1140        assert_eq!(store.load_session_snapshot(&rid).await.unwrap(), None);
1141        let clean = meerkat_core::Session::with_id(session.id().clone());
1142        let clean_snapshot = serde_json::to_vec(&clean).unwrap();
1143        store
1144            .commit_session_snapshot(
1145                &rid,
1146                SessionDelta {
1147                    session_snapshot: clean_snapshot.clone(),
1148                },
1149            )
1150            .await
1151            .unwrap();
1152        assert!(
1153            store
1154                .replace_session_snapshot_if_current(&rid, &clean_snapshot, snapshot)
1155                .await
1156                .is_err()
1157        );
1158        assert_eq!(
1159            store.load_session_snapshot(&rid).await.unwrap(),
1160            Some(clean_snapshot)
1161        );
1162        assert!(
1163            store
1164                .load_pending_compaction_projections(&rid)
1165                .await
1166                .unwrap()
1167                .is_empty()
1168        );
1169    }
1170
1171    #[tokio::test]
1172    async fn atomic_apply_roundtrip() {
1173        let store = InMemoryRuntimeStore::new();
1174        let rid = LogicalRuntimeId::new("test-runtime");
1175        let run_id = RunId::new();
1176        let input_id = InputId::new();
1177
1178        let bundle = StoredInputState::new_accepted(input_id.clone());
1179        let receipt = make_receipt(run_id.clone(), 0);
1180
1181        let session = session_with_user("hello");
1182        let session_snapshot = serde_json::to_vec(&session).unwrap();
1183
1184        store
1185            .atomic_apply(
1186                &rid,
1187                Some(SessionDelta { session_snapshot }),
1188                receipt.clone(),
1189                vec![persistable(bundle)],
1190                None,
1191            )
1192            .await
1193            .unwrap();
1194
1195        // Load input states
1196        let states = store.load_input_states(&rid).await.unwrap();
1197        assert_eq!(states.len(), 1);
1198        assert_eq!(states[0].state.input_id, input_id);
1199
1200        // Load receipt
1201        let loaded = store.load_boundary_receipt(&rid, &run_id, 0).await.unwrap();
1202        assert!(loaded.is_some());
1203    }
1204
1205    #[tokio::test]
1206    async fn atomic_apply_rejects_non_session_snapshot_without_owner_context() {
1207        let store = InMemoryRuntimeStore::new();
1208        let rid = LogicalRuntimeId::new("test-runtime");
1209        let run_id = RunId::new();
1210        let input_id = InputId::new();
1211
1212        let bundle = StoredInputState::new_accepted(input_id);
1213        let receipt = make_receipt(run_id, 0);
1214
1215        // Owner-context absence is not a license to store arbitrary bytes as a
1216        // session snapshot: a non-deserializable snapshot must fail closed.
1217        let err = store
1218            .atomic_apply(
1219                &rid,
1220                Some(SessionDelta {
1221                    session_snapshot: b"session-data".to_vec(),
1222                }),
1223                receipt,
1224                vec![persistable(bundle)],
1225                None,
1226            )
1227            .await
1228            .expect_err("non-Session snapshot must be rejected");
1229
1230        match err {
1231            RuntimeStoreError::WriteFailed(message) => {
1232                assert!(
1233                    message.contains("not a Session"),
1234                    "unexpected WriteFailed message: {message}"
1235                );
1236            }
1237            other => panic!("expected WriteFailed, got {other:?}"),
1238        }
1239    }
1240
1241    #[tokio::test]
1242    async fn persist_and_load_single_state() {
1243        let store = InMemoryRuntimeStore::new();
1244        let rid = LogicalRuntimeId::new("test");
1245        let input_id = InputId::new();
1246        let bundle = StoredInputState::new_accepted(input_id.clone());
1247
1248        store
1249            .persist_input_state(&rid, &persistable(bundle))
1250            .await
1251            .unwrap();
1252
1253        let loaded = store.load_input_state(&rid, &input_id).await.unwrap();
1254        assert!(loaded.is_some());
1255        assert_eq!(loaded.unwrap().state.input_id, input_id);
1256    }
1257
1258    #[tokio::test]
1259    async fn load_nonexistent_returns_none() {
1260        let store = InMemoryRuntimeStore::new();
1261        let rid = LogicalRuntimeId::new("test");
1262
1263        let states = store.load_input_states(&rid).await.unwrap();
1264        assert!(states.is_empty());
1265
1266        let state = store.load_input_state(&rid, &InputId::new()).await.unwrap();
1267        assert!(state.is_none());
1268
1269        let receipt = store
1270            .load_boundary_receipt(&rid, &RunId::new(), 0)
1271            .await
1272            .unwrap();
1273        assert!(receipt.is_none());
1274    }
1275
1276    #[tokio::test]
1277    async fn atomic_apply_updates_existing() {
1278        let store = InMemoryRuntimeStore::new();
1279        let rid = LogicalRuntimeId::new("test");
1280        let input_id = InputId::new();
1281
1282        // First write
1283        let bundle1 = StoredInputState::new_accepted(input_id.clone());
1284        store
1285            .atomic_apply(
1286                &rid,
1287                None,
1288                make_receipt(RunId::new(), 0),
1289                vec![persistable(bundle1)],
1290                None,
1291            )
1292            .await
1293            .unwrap();
1294
1295        // Second write with updated seed phase
1296        let mut bundle2 = StoredInputState::new_accepted(input_id.clone());
1297        bundle2.seed.phase = crate::input_state::InputLifecycleState::Queued;
1298        store
1299            .atomic_apply(
1300                &rid,
1301                None,
1302                make_receipt(RunId::new(), 1),
1303                vec![persistable(bundle2)],
1304                None,
1305            )
1306            .await
1307            .unwrap();
1308
1309        let states = store.load_input_states(&rid).await.unwrap();
1310        assert_eq!(states.len(), 1);
1311        assert_eq!(
1312            states[0].seed.phase,
1313            crate::input_state::InputLifecycleState::Queued
1314        );
1315    }
1316
1317    #[tokio::test]
1318    async fn atomic_apply_validates_session_store_key_without_aliasing_snapshot() {
1319        let store = InMemoryRuntimeStore::new();
1320        let rid = LogicalRuntimeId::new("runtime-key");
1321        let session = meerkat_core::Session::new();
1322        let session_id = session.id().clone();
1323        let snapshot = serde_json::to_vec(&session).unwrap();
1324
1325        store
1326            .atomic_apply(
1327                &rid,
1328                Some(SessionDelta {
1329                    session_snapshot: snapshot.clone(),
1330                }),
1331                make_receipt(RunId::new(), 0),
1332                vec![],
1333                Some(session_id.clone()),
1334            )
1335            .await
1336            .unwrap();
1337
1338        assert_eq!(
1339            store.load_session_snapshot(&rid).await.unwrap(),
1340            Some(snapshot)
1341        );
1342        assert!(
1343            store
1344                .load_session_snapshot(&LogicalRuntimeId::legacy_session_uuid_alias(&session_id))
1345                .await
1346                .unwrap()
1347                .is_none(),
1348            "session_store_key must validate the snapshot identity, not create a raw UUID runtime alias"
1349        );
1350    }
1351
1352    #[tokio::test]
1353    async fn atomic_apply_rejects_mismatched_session_store_key() {
1354        let store = InMemoryRuntimeStore::new();
1355        let rid = LogicalRuntimeId::new("runtime-key");
1356        let session = meerkat_core::Session::new();
1357        let wrong_session_id = meerkat_core::Session::new().id().clone();
1358        let snapshot = serde_json::to_vec(&session).unwrap();
1359
1360        let err = store
1361            .atomic_apply(
1362                &rid,
1363                Some(SessionDelta {
1364                    session_snapshot: snapshot,
1365                }),
1366                make_receipt(RunId::new(), 0),
1367                vec![],
1368                Some(wrong_session_id),
1369            )
1370            .await
1371            .expect_err("mismatched session_store_key should fail");
1372
1373        assert!(matches!(err, RuntimeStoreError::SessionKeyMismatch { .. }));
1374        assert!(store.load_session_snapshot(&rid).await.unwrap().is_none());
1375    }
1376
1377    #[tokio::test]
1378    async fn atomic_apply_persists_machine_owned_receipt() {
1379        let store = InMemoryRuntimeStore::new();
1380        let rid = LogicalRuntimeId::new("test");
1381        let run_id = RunId::new();
1382        let input_id = InputId::new();
1383        let session = meerkat_core::Session::new();
1384        let snapshot = serde_json::to_vec(&session).unwrap();
1385        let receipt = RunBoundaryReceipt {
1386            run_id: run_id.clone(),
1387            boundary: RunApplyBoundary::Immediate,
1388            contributing_input_ids: vec![input_id.clone()],
1389            conversation_digest: Some("machine-owned-digest".to_string()),
1390            message_count: 42,
1391            sequence: 7,
1392        };
1393
1394        store
1395            .atomic_apply(
1396                &rid,
1397                Some(SessionDelta {
1398                    session_snapshot: snapshot,
1399                }),
1400                receipt.clone(),
1401                vec![persistable(StoredInputState::new_accepted(input_id))],
1402                None,
1403            )
1404            .await
1405            .unwrap();
1406
1407        assert_eq!(receipt.run_id, run_id);
1408        assert!(receipt.conversation_digest.is_some());
1409        let loaded = store
1410            .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1411            .await
1412            .unwrap();
1413        assert!(loaded.is_some(), "receipt should be persisted");
1414        let Some(loaded) = loaded else {
1415            unreachable!("asserted above");
1416        };
1417        assert_eq!(loaded, receipt);
1418    }
1419
1420    #[tokio::test]
1421    async fn multiple_runtimes_isolated() {
1422        let store = InMemoryRuntimeStore::new();
1423        let rid1 = LogicalRuntimeId::new("runtime-1");
1424        let rid2 = LogicalRuntimeId::new("runtime-2");
1425
1426        store
1427            .persist_input_state(
1428                &rid1,
1429                &persistable(StoredInputState::new_accepted(InputId::new())),
1430            )
1431            .await
1432            .unwrap();
1433        store
1434            .persist_input_state(
1435                &rid2,
1436                &persistable(StoredInputState::new_accepted(InputId::new())),
1437            )
1438            .await
1439            .unwrap();
1440        store
1441            .persist_input_state(
1442                &rid2,
1443                &persistable(StoredInputState::new_accepted(InputId::new())),
1444            )
1445            .await
1446            .unwrap();
1447
1448        let s1 = store.load_input_states(&rid1).await.unwrap();
1449        let s2 = store.load_input_states(&rid2).await.unwrap();
1450        assert_eq!(s1.len(), 1);
1451        assert_eq!(s2.len(), 2);
1452    }
1453
1454    #[tokio::test]
1455    async fn load_session_snapshot_roundtrip() {
1456        let store = InMemoryRuntimeStore::new();
1457        let rid = LogicalRuntimeId::new("runtime");
1458        let snapshot = serde_json::to_vec(&meerkat_core::Session::new()).unwrap();
1459
1460        store
1461            .atomic_apply(
1462                &rid,
1463                Some(SessionDelta {
1464                    session_snapshot: snapshot.clone(),
1465                }),
1466                make_receipt(RunId::new(), 0),
1467                vec![],
1468                None,
1469            )
1470            .await
1471            .unwrap();
1472
1473        let loaded = store.load_session_snapshot(&rid).await.unwrap();
1474        assert_eq!(loaded, Some(snapshot));
1475    }
1476
1477    #[tokio::test]
1478    async fn commit_session_snapshot_rejects_stale_runtime_parent() {
1479        let store = InMemoryRuntimeStore::new();
1480        let rid = LogicalRuntimeId::new("runtime-stale-parent");
1481        let accepted = session_with_user("accepted runtime turn");
1482        let mut stale = meerkat_core::Session::with_id(accepted.id().clone());
1483        stale.push(meerkat_core::types::Message::User(
1484            meerkat_core::types::UserMessage::text("stale runtime turn".to_string()),
1485        ));
1486        let accepted_snapshot = serde_json::to_vec(&accepted).unwrap();
1487
1488        store
1489            .commit_session_snapshot(
1490                &rid,
1491                SessionDelta {
1492                    session_snapshot: accepted_snapshot.clone(),
1493                },
1494            )
1495            .await
1496            .unwrap();
1497
1498        let err = store
1499            .commit_session_snapshot(
1500                &rid,
1501                SessionDelta {
1502                    session_snapshot: serde_json::to_vec(&stale).unwrap(),
1503                },
1504            )
1505            .await
1506            .expect_err("stale non-continuation must not overwrite runtime snapshot");
1507
1508        assert!(matches!(err, RuntimeStoreError::WriteFailed(_)));
1509        assert_eq!(
1510            store.load_session_snapshot(&rid).await.unwrap(),
1511            Some(accepted_snapshot)
1512        );
1513    }
1514
1515    #[tokio::test]
1516    async fn atomic_apply_keeps_current_snapshot_when_incoming_is_superseded() {
1517        let store = InMemoryRuntimeStore::new();
1518        let rid = LogicalRuntimeId::new("runtime-superseded-terminal");
1519        let incoming = session_with_user("turn input");
1520        let mut current = incoming.clone();
1521        current.push(meerkat_core::types::Message::BlockAssistant(
1522            meerkat_core::types::BlockAssistantMessage {
1523                blocks: vec![meerkat_core::types::AssistantBlock::Text {
1524                    text: "peer response already applied".to_string(),
1525                    meta: None,
1526                }],
1527                stop_reason: meerkat_core::types::StopReason::EndTurn,
1528                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1529                created_at: meerkat_core::types::message_timestamp_now(),
1530            },
1531        ));
1532        let current_snapshot = serde_json::to_vec(&current).unwrap();
1533        let receipt = make_receipt(RunId::new(), 11);
1534
1535        store
1536            .commit_session_snapshot(
1537                &rid,
1538                SessionDelta {
1539                    session_snapshot: current_snapshot.clone(),
1540                },
1541            )
1542            .await
1543            .unwrap();
1544
1545        store
1546            .atomic_apply(
1547                &rid,
1548                Some(SessionDelta {
1549                    session_snapshot: serde_json::to_vec(&incoming).unwrap(),
1550                }),
1551                receipt.clone(),
1552                vec![],
1553                Some(incoming.id().clone()),
1554            )
1555            .await
1556            .unwrap();
1557
1558        assert_eq!(
1559            store.load_session_snapshot(&rid).await.unwrap(),
1560            Some(current_snapshot)
1561        );
1562        // The session snapshot was classified superseded and skipped, so the
1563        // boundary receipt for that boundary must NOT advance against the
1564        // retained (more-advanced) session snapshot.
1565        assert_eq!(
1566            store
1567                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1568                .await
1569                .unwrap(),
1570            None
1571        );
1572    }
1573
1574    #[tokio::test]
1575    async fn atomic_apply_skips_inputs_when_session_snapshot_superseded() {
1576        let store = InMemoryRuntimeStore::new();
1577        let rid = LogicalRuntimeId::new("runtime-superseded-inputs");
1578        let incoming = session_with_user("turn input");
1579        let mut current = incoming.clone();
1580        current.push(meerkat_core::types::Message::BlockAssistant(
1581            meerkat_core::types::BlockAssistantMessage {
1582                blocks: vec![meerkat_core::types::AssistantBlock::Text {
1583                    text: "peer response already applied".to_string(),
1584                    meta: None,
1585                }],
1586                stop_reason: meerkat_core::types::StopReason::EndTurn,
1587                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1588                created_at: meerkat_core::types::message_timestamp_now(),
1589            },
1590        ));
1591        let current_snapshot = serde_json::to_vec(&current).unwrap();
1592        let receipt = make_receipt(RunId::new(), 21);
1593        let input_id = InputId::new();
1594        let bundle = StoredInputState::new_accepted(input_id.clone());
1595
1596        store
1597            .commit_session_snapshot(
1598                &rid,
1599                SessionDelta {
1600                    session_snapshot: current_snapshot.clone(),
1601                },
1602            )
1603            .await
1604            .unwrap();
1605
1606        store
1607            .atomic_apply(
1608                &rid,
1609                Some(SessionDelta {
1610                    session_snapshot: serde_json::to_vec(&incoming).unwrap(),
1611                }),
1612                receipt.clone(),
1613                vec![persistable(bundle)],
1614                Some(incoming.id().clone()),
1615            )
1616            .await
1617            .unwrap();
1618
1619        // Snapshot retained, receipt + input-state writes skipped as a unit.
1620        assert_eq!(
1621            store.load_session_snapshot(&rid).await.unwrap(),
1622            Some(current_snapshot)
1623        );
1624        assert_eq!(
1625            store
1626                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1627                .await
1628                .unwrap(),
1629            None
1630        );
1631        assert!(store.load_input_states(&rid).await.unwrap().is_empty());
1632    }
1633
1634    #[tokio::test]
1635    async fn atomic_apply_allows_first_generated_snapshot_after_placeholder() {
1636        let store = InMemoryRuntimeStore::new();
1637        let rid = LogicalRuntimeId::new("runtime-placeholder");
1638        let mut placeholder = meerkat_core::Session::new();
1639        placeholder.set_system_prompt("base system".to_string());
1640        let mut incoming = meerkat_core::Session::with_id(placeholder.id().clone());
1641        incoming.set_system_prompt("base system".to_string());
1642        incoming.push(meerkat_core::types::Message::User(
1643            meerkat_core::types::UserMessage::text("verbose first turn".to_string()),
1644        ));
1645        let parent_revision = incoming.transcript_revision().unwrap();
1646        incoming
1647            .commit_transcript_rewrite(
1648                meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
1649                vec![meerkat_core::types::Message::User(
1650                    meerkat_core::types::UserMessage::compaction_summary(
1651                        "[Context compacted] first turn",
1652                    ),
1653                )],
1654                meerkat_core::TranscriptRewriteReason::new("compaction"),
1655                Some("meerkat-core".to_string()),
1656                Some(parent_revision),
1657            )
1658            .unwrap();
1659        let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
1660        let receipt = make_receipt(RunId::new(), 12);
1661
1662        store
1663            .commit_session_snapshot(
1664                &rid,
1665                SessionDelta {
1666                    session_snapshot: serde_json::to_vec(&placeholder).unwrap(),
1667                },
1668            )
1669            .await
1670            .unwrap();
1671
1672        store
1673            .atomic_apply(
1674                &rid,
1675                Some(SessionDelta {
1676                    session_snapshot: incoming_snapshot.clone(),
1677                }),
1678                receipt.clone(),
1679                vec![],
1680                Some(incoming.id().clone()),
1681            )
1682            .await
1683            .unwrap();
1684
1685        assert_eq!(
1686            store.load_session_snapshot(&rid).await.unwrap(),
1687            Some(incoming_snapshot)
1688        );
1689        assert_eq!(
1690            store
1691                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1692                .await
1693                .unwrap(),
1694            Some(receipt)
1695        );
1696    }
1697
1698    #[tokio::test]
1699    async fn atomic_apply_allows_generated_compaction_before_retained_tail() {
1700        let store = InMemoryRuntimeStore::new();
1701        let rid = LogicalRuntimeId::new("runtime-compaction-tail");
1702        let mut previous = meerkat_core::Session::new();
1703        previous.set_system_prompt("runtime system before context refresh".to_string());
1704        previous.push(meerkat_core::types::Message::User(
1705            meerkat_core::types::UserMessage::text("Turn 1 request".to_string()),
1706        ));
1707        previous.push(meerkat_core::types::Message::BlockAssistant(
1708            meerkat_core::types::BlockAssistantMessage {
1709                blocks: vec![meerkat_core::types::AssistantBlock::Text {
1710                    text: "Turn 1 answer".to_string(),
1711                    meta: None,
1712                }],
1713                stop_reason: meerkat_core::types::StopReason::EndTurn,
1714                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1715                created_at: meerkat_core::types::message_timestamp_now(),
1716            },
1717        ));
1718
1719        let mut incoming = meerkat_core::Session::with_id(previous.id().clone());
1720        incoming.set_system_prompt("runtime system after context refresh".to_string());
1721        incoming.push(meerkat_core::types::Message::User(
1722            meerkat_core::types::UserMessage::text(
1723                "Verbose context that will be compacted".to_string(),
1724            ),
1725        ));
1726        for message in previous.messages()[1..].iter().cloned() {
1727            incoming.push(message);
1728        }
1729        incoming.push(meerkat_core::types::Message::BlockAssistant(
1730            meerkat_core::types::BlockAssistantMessage {
1731                blocks: vec![meerkat_core::types::AssistantBlock::Text {
1732                    text: "Turn 2 generated answer".to_string(),
1733                    meta: None,
1734                }],
1735                stop_reason: meerkat_core::types::StopReason::EndTurn,
1736                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1737                created_at: meerkat_core::types::message_timestamp_now(),
1738            },
1739        ));
1740        let parent_revision = incoming.transcript_revision().unwrap();
1741        incoming
1742            .commit_transcript_rewrite(
1743                meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
1744                vec![meerkat_core::types::Message::User(
1745                    meerkat_core::types::UserMessage::compaction_summary(
1746                        "[Context compacted] Earlier runtime context".to_string(),
1747                    ),
1748                )],
1749                meerkat_core::TranscriptRewriteReason::new("compaction"),
1750                Some("meerkat-core".to_string()),
1751                Some(parent_revision),
1752            )
1753            .unwrap();
1754        let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
1755        let receipt = make_receipt(RunId::new(), 13);
1756
1757        store
1758            .commit_session_snapshot(
1759                &rid,
1760                SessionDelta {
1761                    session_snapshot: serde_json::to_vec(&previous).unwrap(),
1762                },
1763            )
1764            .await
1765            .unwrap();
1766
1767        store
1768            .atomic_apply(
1769                &rid,
1770                Some(SessionDelta {
1771                    session_snapshot: incoming_snapshot.clone(),
1772                }),
1773                receipt.clone(),
1774                vec![],
1775                Some(incoming.id().clone()),
1776            )
1777            .await
1778            .unwrap();
1779
1780        assert_eq!(
1781            store.load_session_snapshot(&rid).await.unwrap(),
1782            Some(incoming_snapshot)
1783        );
1784        assert_eq!(
1785            store
1786                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1787                .await
1788                .unwrap(),
1789            Some(receipt)
1790        );
1791    }
1792
1793    #[tokio::test]
1794    async fn commit_machine_lifecycle_persists_binding_facts() {
1795        use crate::runtime_state::RuntimeState;
1796
1797        let store = InMemoryRuntimeStore::new();
1798        let rid = LogicalRuntimeId::new("runtime-binding");
1799        let binding = MachineLifecycleBindingFacts::new(
1800            Some("rt:session:abc".to_string()),
1801            Some(7),
1802            Some(3),
1803            Some("epoch-1".to_string()),
1804        );
1805
1806        store
1807            .commit_machine_lifecycle(
1808                &rid,
1809                MachineLifecycleCommit::new_with_binding(
1810                    RuntimeState::Retired,
1811                    binding.clone(),
1812                    crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1813                ),
1814                &[],
1815            )
1816            .await
1817            .unwrap();
1818
1819        let lifecycle = crate::store::load_machine_lifecycle(&store, &rid)
1820            .await
1821            .unwrap()
1822            .expect("machine lifecycle snapshot");
1823        assert_eq!(lifecycle.runtime_state(), RuntimeState::Retired);
1824        assert_eq!(lifecycle.binding(), &binding);
1825        assert_eq!(
1826            crate::store::load_runtime_state(&store, &rid)
1827                .await
1828                .unwrap(),
1829            Some(RuntimeState::Retired)
1830        );
1831    }
1832
1833    #[tokio::test]
1834    async fn unregister_finalization_atomically_retires_ops_epoch_and_is_idempotent() {
1835        let store = InMemoryRuntimeStore::new();
1836        let reopened = store.clone();
1837        let runtime_id = LogicalRuntimeId::new("runtime-unregister-finalization");
1838        let stale_ops = crate::ops_lifecycle::RuntimeOpsLifecycleRegistry::new()
1839            .capture_persistence_snapshot(
1840                meerkat_core::RuntimeEpochId::new(),
1841                &meerkat_core::EpochCursorState::new(),
1842            )
1843            .unwrap();
1844        store
1845            .persist_ops_lifecycle(&runtime_id, &stale_ops)
1846            .await
1847            .unwrap();
1848        let retired_ops_epoch = stale_ops.epoch_id.clone();
1849
1850        for _ in 0..2 {
1851            store
1852                .commit_unregister_finalization(
1853                    &runtime_id,
1854                    MachineLifecycleCommit::new_with_binding(
1855                        RuntimeState::Stopped,
1856                        MachineLifecycleBindingFacts::new(None, None, None, None),
1857                        crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1858                    )
1859                    .for_unregister_finalization(retired_ops_epoch.clone()),
1860                    &[],
1861                )
1862                .await
1863                .unwrap();
1864        }
1865
1866        assert_eq!(
1867            crate::store::load_runtime_state(&reopened, &runtime_id)
1868                .await
1869                .unwrap(),
1870            Some(RuntimeState::Stopped)
1871        );
1872        assert!(
1873            reopened
1874                .load_ops_lifecycle(&runtime_id)
1875                .await
1876                .unwrap()
1877                .is_none(),
1878            "the same critical section that publishes terminal lifecycle must remove the ops epoch"
1879        );
1880        let late_error = reopened
1881            .persist_ops_lifecycle(&runtime_id, &stale_ops)
1882            .await
1883            .expect_err("a detached callback must not resurrect its retired ops epoch");
1884        assert!(matches!(
1885            late_error,
1886            RuntimeStoreError::OpsLifecycleEpochRetired { epoch_id, .. }
1887                if epoch_id == retired_ops_epoch
1888        ));
1889        assert!(
1890            reopened
1891                .load_ops_lifecycle(&runtime_id)
1892                .await
1893                .unwrap()
1894                .is_none()
1895        );
1896    }
1897
1898    #[tokio::test]
1899    async fn delayed_old_epoch_finalizer_cannot_delete_or_overwrite_new_ops_epoch() {
1900        let store = InMemoryRuntimeStore::new();
1901        let runtime_id = LogicalRuntimeId::new("runtime-old-finalizer-new-epoch");
1902        let registry = crate::ops_lifecycle::RuntimeOpsLifecycleRegistry::new();
1903        let old_ops = registry
1904            .capture_persistence_snapshot(
1905                meerkat_core::RuntimeEpochId::new(),
1906                &meerkat_core::EpochCursorState::new(),
1907            )
1908            .unwrap();
1909        let new_ops = registry
1910            .capture_persistence_snapshot(
1911                meerkat_core::RuntimeEpochId::new(),
1912                &meerkat_core::EpochCursorState::new(),
1913            )
1914            .unwrap();
1915        store
1916            .persist_ops_lifecycle(&runtime_id, &old_ops)
1917            .await
1918            .unwrap();
1919        store
1920            .persist_ops_lifecycle(&runtime_id, &new_ops)
1921            .await
1922            .unwrap();
1923
1924        store
1925            .commit_unregister_finalization(
1926                &runtime_id,
1927                MachineLifecycleCommit::new_with_binding(
1928                    RuntimeState::Stopped,
1929                    MachineLifecycleBindingFacts::new(None, None, None, None),
1930                    crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1931                )
1932                .for_unregister_finalization(old_ops.epoch_id.clone()),
1933                &[],
1934            )
1935            .await
1936            .unwrap();
1937
1938        assert_eq!(
1939            store
1940                .load_ops_lifecycle(&runtime_id)
1941                .await
1942                .unwrap()
1943                .expect("new epoch row must survive delayed old finalization")
1944                .epoch_id,
1945            new_ops.epoch_id
1946        );
1947        assert!(matches!(
1948            store
1949                .persist_ops_lifecycle(&runtime_id, &old_ops)
1950                .await
1951                .expect_err("retired old epoch stays fenced"),
1952            RuntimeStoreError::OpsLifecycleEpochRetired { .. }
1953        ));
1954        store
1955            .persist_ops_lifecycle(&runtime_id, &new_ops)
1956            .await
1957            .unwrap();
1958    }
1959
1960    #[tokio::test]
1961    async fn clear_session_snapshot_if_current_sets_quarantine_marker_cleared_on_write() {
1962        let store = InMemoryRuntimeStore::new();
1963        let rid = LogicalRuntimeId::new("runtime-quarantine");
1964        let rejected = serde_json::to_vec(&session_with_user("rejected")).unwrap();
1965
1966        assert!(!store.is_runtime_projection_quarantined(&rid).await.unwrap());
1967        store
1968            .commit_session_snapshot(
1969                &rid,
1970                SessionDelta {
1971                    session_snapshot: rejected.clone(),
1972                },
1973            )
1974            .await
1975            .unwrap();
1976        assert!(
1977            store
1978                .clear_session_snapshot_if_current(&rid, &rejected)
1979                .await
1980                .unwrap()
1981        );
1982        assert!(
1983            store.is_runtime_projection_quarantined(&rid).await.unwrap(),
1984            "clearing the rejected snapshot must record the in-memory quarantine marker"
1985        );
1986
1987        // A live snapshot write reclaims runtime authority and clears the marker.
1988        store
1989            .commit_session_snapshot(
1990                &rid,
1991                SessionDelta {
1992                    session_snapshot: serde_json::to_vec(&session_with_user("revived")).unwrap(),
1993                },
1994            )
1995            .await
1996            .unwrap();
1997        assert!(
1998            !store.is_runtime_projection_quarantined(&rid).await.unwrap(),
1999            "a live snapshot write must clear the in-memory quarantine marker"
2000        );
2001    }
2002}