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/// Inner state protected by the mutex.
34#[derive(Debug, Default)]
35struct Inner {
36    /// runtime_id → (input_id → StoredInputState). IndexMap for deterministic iteration order.
37    input_states: HashMap<String, IndexMap<InputId, StoredInputState>>,
38    /// Receipt storage.
39    receipts: HashMap<ReceiptKey, RunBoundaryReceipt>,
40    /// Runtime session snapshots keyed by canonical runtime id.
41    sessions: HashMap<String, Vec<u8>>,
42    /// Canonical runtime ids whose projection fallback is quarantined.
43    ///
44    /// Mirrors the durable SQLite `runtime_projection_quarantine` table: set
45    /// when a rejected runtime snapshot is cleared via
46    /// `clear_session_snapshot_if_current`, cleared whenever a live snapshot is
47    /// written for the runtime.
48    projection_quarantine: HashSet<String>,
49    /// Persisted machine lifecycle snapshots.
50    runtime_lifecycle: HashMap<String, MachineLifecycleSnapshot>,
51    /// Persisted ops lifecycle snapshots.
52    ops_lifecycle_snapshots: HashMap<String, PersistedOpsSnapshot>,
53}
54
55/// In-memory runtime store. Thread-safe via `tokio::sync::Mutex`.
56#[derive(Debug, Clone)]
57pub struct InMemoryRuntimeStore {
58    inner: Arc<Mutex<Inner>>,
59    auth_oauth_flow_snapshot: Arc<StdMutex<Option<Vec<u8>>>>,
60}
61
62impl InMemoryRuntimeStore {
63    pub fn new() -> Self {
64        Self {
65            inner: Arc::new(Mutex::new(Inner::default())),
66            auth_oauth_flow_snapshot: Arc::new(StdMutex::new(None)),
67        }
68    }
69}
70
71impl Default for InMemoryRuntimeStore {
72    fn default() -> Self {
73        Self::new()
74    }
75}
76
77fn is_runtime_placeholder_session(session: &meerkat_core::Session) -> bool {
78    session.transcript_history_state().ok().flatten().is_none()
79        && matches!(
80            session.messages(),
81            [] | [meerkat_core::types::Message::System(_)]
82        )
83}
84
85/// Deserialize a persisted session-snapshot blob through typed serde, matching
86/// the SQLite runtime store read path. `Session::deserialize` validates the
87/// mandatory envelope version against the generated persistence version
88/// authority, so a missing or non-current (v0/v1) row fails closed instead of
89/// silently defaulting or upgrading on read.
90fn deserialize_persisted_session(bytes: &[u8]) -> Result<meerkat_core::Session, RuntimeStoreError> {
91    serde_json::from_slice(bytes).map_err(|err| RuntimeStoreError::ReadFailed(err.to_string()))
92}
93
94#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
95#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
96impl RuntimeStore for InMemoryRuntimeStore {
97    fn persist_auth_oauth_flow_snapshot(
98        &self,
99        snapshot_json: &[u8],
100    ) -> Result<(), RuntimeStoreError> {
101        *self
102            .auth_oauth_flow_snapshot
103            .lock()
104            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))? =
105            Some(snapshot_json.to_vec());
106        Ok(())
107    }
108
109    fn load_auth_oauth_flow_snapshot(&self) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
110        self.auth_oauth_flow_snapshot
111            .lock()
112            .map(|snapshot| snapshot.clone())
113            .map_err(|err| RuntimeStoreError::ReadFailed(err.to_string()))
114    }
115
116    fn update_auth_oauth_flow_snapshot(
117        &self,
118        update: &mut AuthOAuthFlowSnapshotUpdate<'_>,
119    ) -> Result<(), RuntimeStoreError> {
120        let mut snapshot = self
121            .auth_oauth_flow_snapshot
122            .lock()
123            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
124        let next = update(snapshot.as_deref())?;
125        *snapshot = Some(next);
126        Ok(())
127    }
128
129    async fn commit_session_snapshot(
130        &self,
131        runtime_id: &LogicalRuntimeId,
132        session_delta: SessionDelta,
133    ) -> Result<(), RuntimeStoreError> {
134        let incoming: meerkat_core::Session =
135            serde_json::from_slice(&session_delta.session_snapshot)
136                .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
137        let mut inner = self.inner.lock().await;
138        let previous = inner
139            .sessions
140            .get(&runtime_id.0)
141            .map(|snapshot| deserialize_persisted_session(snapshot))
142            .transpose()?;
143        meerkat_core::session_store::run_boundary_snapshot_save_guard(&incoming, previous.as_ref())
144            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
145        inner
146            .sessions
147            .insert(runtime_id.0.clone(), session_delta.session_snapshot);
148        inner.projection_quarantine.remove(&runtime_id.0);
149        Ok(())
150    }
151
152    async fn commit_session_transcript_rewrite_snapshot(
153        &self,
154        runtime_id: &LogicalRuntimeId,
155        session_delta: SessionDelta,
156        commit: &meerkat_core::TranscriptRewriteCommit,
157    ) -> Result<(), RuntimeStoreError> {
158        let incoming: meerkat_core::Session =
159            serde_json::from_slice(&session_delta.session_snapshot)
160                .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
161        let mut inner = self.inner.lock().await;
162        let previous = inner
163            .sessions
164            .get(&runtime_id.0)
165            .map(|snapshot| deserialize_persisted_session(snapshot))
166            .transpose()?;
167        meerkat_core::session_store::transcript_rewrite_save_guard(
168            &incoming,
169            previous.as_ref(),
170            commit,
171        )
172        .map_err(|err| match err {
173            meerkat_core::SessionStoreError::TranscriptRevisionConflict {
174                expected,
175                actual,
176                ..
177            } => RuntimeStoreError::TranscriptRevisionConflict { expected, actual },
178            other => RuntimeStoreError::WriteFailed(other.to_string()),
179        })?;
180        inner
181            .sessions
182            .insert(runtime_id.0.clone(), session_delta.session_snapshot);
183        inner.projection_quarantine.remove(&runtime_id.0);
184        Ok(())
185    }
186
187    async fn atomic_apply(
188        &self,
189        runtime_id: &LogicalRuntimeId,
190        session_delta: Option<SessionDelta>,
191        receipt: RunBoundaryReceipt,
192        input_updates: Vec<InputStatePersistenceRecord>,
193        session_store_key: Option<meerkat_core::types::SessionId>,
194    ) -> Result<(), RuntimeStoreError> {
195        let mut inner = self.inner.lock().await;
196
197        // All writes in one lock acquisition (atomic for in-memory)
198        let rid = runtime_id.0.clone();
199
200        // Session delta. The supersession verdict computed here keys the
201        // entire commit: if the incoming session snapshot is classified as
202        // superseded (the persisted head is already a valid append-extension
203        // of it), the snapshot write is skipped AND so are the receipt + input
204        // writes, so receipt/input ordering identity never advances past the
205        // retained session truth.
206        let mut session_snapshot_superseded = false;
207        if let Some(delta) = session_delta {
208            let incoming_session =
209                serde_json::from_slice::<meerkat_core::Session>(&delta.session_snapshot);
210            let mut persist_session_snapshot = true;
211            match (incoming_session, session_store_key) {
212                (Ok(incoming_session), session_store_key) => {
213                    if let Some(session_store_key) = session_store_key
214                        && incoming_session.id() != &session_store_key
215                    {
216                        return Err(RuntimeStoreError::SessionKeyMismatch {
217                            expected: session_store_key,
218                            actual: incoming_session.id().clone(),
219                        });
220                    }
221                    let previous_session = inner
222                        .sessions
223                        .get(&rid)
224                        .and_then(|snapshot| deserialize_persisted_session(snapshot).ok());
225                    if let Err(err) = meerkat_core::session_store::run_boundary_snapshot_save_guard(
226                        &incoming_session,
227                        previous_session.as_ref(),
228                    ) {
229                        if previous_session
230                            .as_ref()
231                            .is_some_and(is_runtime_placeholder_session)
232                        {
233                            persist_session_snapshot = true;
234                        } else if previous_session.as_ref().is_some_and(|previous_session| {
235                            meerkat_core::session_store::run_boundary_snapshot_save_guard(
236                                previous_session,
237                                Some(&incoming_session),
238                            )
239                            .is_ok()
240                        }) {
241                            persist_session_snapshot = false;
242                            session_snapshot_superseded = true;
243                        } else {
244                            return Err(RuntimeStoreError::WriteFailed(err.to_string()));
245                        }
246                    }
247                }
248                (Err(err), Some(session_store_key)) => {
249                    return Err(RuntimeStoreError::WriteFailed(format!(
250                        "session snapshot for {session_store_key} is not a Session: {err}"
251                    )));
252                }
253                (Err(err), None) => {
254                    return Err(RuntimeStoreError::WriteFailed(format!(
255                        "session snapshot is not a Session: {err}"
256                    )));
257                }
258            }
259            if persist_session_snapshot {
260                inner.sessions.insert(rid.clone(), delta.session_snapshot);
261                inner.projection_quarantine.remove(&rid);
262            }
263        }
264
265        // When the session snapshot was superseded and skipped, the boundary
266        // receipt and input-state updates for that boundary must also be
267        // skipped: advancing them against a retained (older) session snapshot
268        // would split receipt/input ordering identity from session truth.
269        if session_snapshot_superseded {
270            return Ok(());
271        }
272
273        // Receipt
274        let key = ReceiptKey {
275            runtime_id: rid.clone(),
276            run_id: receipt.run_id.clone(),
277            sequence: receipt.sequence,
278        };
279        inner.receipts.insert(key, receipt);
280
281        // Input states
282        let states = inner.input_states.entry(rid).or_default();
283        for record in input_updates {
284            let bundle = record.into_stored();
285            states.insert(bundle.state.input_id.clone(), bundle);
286        }
287
288        Ok(())
289    }
290
291    async fn load_input_states(
292        &self,
293        runtime_id: &LogicalRuntimeId,
294    ) -> Result<Vec<StoredInputState>, RuntimeStoreError> {
295        let inner = self.inner.lock().await;
296        let states = inner
297            .input_states
298            .get(&runtime_id.0)
299            .map(|m| m.values().cloned().collect())
300            .unwrap_or_default();
301        Ok(states)
302    }
303
304    async fn load_boundary_receipt(
305        &self,
306        runtime_id: &LogicalRuntimeId,
307        run_id: &RunId,
308        sequence: u64,
309    ) -> Result<Option<RunBoundaryReceipt>, RuntimeStoreError> {
310        let inner = self.inner.lock().await;
311        let key = ReceiptKey {
312            runtime_id: runtime_id.0.clone(),
313            run_id: run_id.clone(),
314            sequence,
315        };
316        Ok(inner.receipts.get(&key).cloned())
317    }
318
319    async fn load_session_snapshot(
320        &self,
321        runtime_id: &LogicalRuntimeId,
322    ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
323        let inner = self.inner.lock().await;
324        Ok(inner.sessions.get(&runtime_id.0).cloned())
325    }
326
327    async fn clear_session_snapshot(
328        &self,
329        runtime_id: &LogicalRuntimeId,
330    ) -> Result<(), RuntimeStoreError> {
331        let mut inner = self.inner.lock().await;
332        inner.sessions.remove(&runtime_id.0);
333        Ok(())
334    }
335
336    async fn replace_session_snapshot_if_current(
337        &self,
338        runtime_id: &LogicalRuntimeId,
339        expected_current: &[u8],
340        replacement: Vec<u8>,
341    ) -> Result<bool, RuntimeStoreError> {
342        let _: meerkat_core::Session = serde_json::from_slice(&replacement)
343            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
344        let mut inner = self.inner.lock().await;
345        let Some(current) = inner.sessions.get_mut(&runtime_id.0) else {
346            return Ok(false);
347        };
348        if current.as_slice() != expected_current {
349            return Ok(false);
350        }
351        *current = replacement;
352        inner.projection_quarantine.remove(&runtime_id.0);
353        Ok(true)
354    }
355
356    async fn clear_session_snapshot_if_current(
357        &self,
358        runtime_id: &LogicalRuntimeId,
359        expected_current: &[u8],
360    ) -> Result<bool, RuntimeStoreError> {
361        let mut inner = self.inner.lock().await;
362        let Some(current) = inner.sessions.get(&runtime_id.0) else {
363            return Ok(false);
364        };
365        if current.as_slice() != expected_current {
366            return Ok(false);
367        }
368        inner.sessions.remove(&runtime_id.0);
369        // Record the in-memory quarantine marker atomically with the snapshot
370        // removal, mirroring the durable SQLite path.
371        inner.projection_quarantine.insert(runtime_id.0.clone());
372        Ok(true)
373    }
374
375    async fn is_runtime_projection_quarantined(
376        &self,
377        runtime_id: &LogicalRuntimeId,
378    ) -> Result<bool, RuntimeStoreError> {
379        let inner = self.inner.lock().await;
380        Ok(inner.projection_quarantine.contains(&runtime_id.0))
381    }
382
383    async fn persist_input_state(
384        &self,
385        runtime_id: &LogicalRuntimeId,
386        state: &InputStatePersistenceRecord,
387    ) -> Result<(), RuntimeStoreError> {
388        let mut inner = self.inner.lock().await;
389        let states = inner.input_states.entry(runtime_id.0.clone()).or_default();
390        let bundle = state.as_stored();
391        states.insert(bundle.state.input_id.clone(), bundle.clone());
392        Ok(())
393    }
394
395    async fn load_input_state(
396        &self,
397        runtime_id: &LogicalRuntimeId,
398        input_id: &InputId,
399    ) -> Result<Option<StoredInputState>, RuntimeStoreError> {
400        let inner = self.inner.lock().await;
401        let state = inner
402            .input_states
403            .get(&runtime_id.0)
404            .and_then(|m| m.get(input_id).cloned());
405        Ok(state)
406    }
407
408    async fn load_machine_lifecycle_record(
409        &self,
410        runtime_id: &LogicalRuntimeId,
411    ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
412        let inner = self.inner.lock().await;
413        inner
414            .runtime_lifecycle
415            .get(&runtime_id.0)
416            .map(|snapshot| MachineLifecycleStoreRecord::from_snapshot(snapshot).encode())
417            .transpose()
418    }
419
420    async fn commit_machine_lifecycle(
421        &self,
422        runtime_id: &LogicalRuntimeId,
423        commit: MachineLifecycleCommit,
424        input_states: &[InputStatePersistenceRecord],
425    ) -> Result<(), RuntimeStoreError> {
426        let mut inner = self.inner.lock().await;
427        let rid = runtime_id.0.clone();
428
429        // Single lock acquisition — atomic for in-memory
430        inner
431            .runtime_lifecycle
432            .insert(rid.clone(), commit.into_snapshot());
433        let states = inner.input_states.entry(rid).or_default();
434        for record in input_states {
435            let bundle = record.as_stored();
436            states.insert(bundle.state.input_id.clone(), bundle.clone());
437        }
438
439        Ok(())
440    }
441
442    async fn persist_ops_lifecycle(
443        &self,
444        runtime_id: &LogicalRuntimeId,
445        snapshot: &PersistedOpsSnapshot,
446    ) -> Result<(), RuntimeStoreError> {
447        let mut inner = self.inner.lock().await;
448        inner
449            .ops_lifecycle_snapshots
450            .insert(runtime_id.0.clone(), snapshot.clone());
451        Ok(())
452    }
453
454    async fn load_ops_lifecycle(
455        &self,
456        runtime_id: &LogicalRuntimeId,
457    ) -> Result<Option<PersistedOpsSnapshot>, RuntimeStoreError> {
458        let inner = self.inner.lock().await;
459        Ok(inner.ops_lifecycle_snapshots.get(&runtime_id.0).cloned())
460    }
461
462    async fn delete_ops_lifecycle(
463        &self,
464        runtime_id: &LogicalRuntimeId,
465    ) -> Result<(), RuntimeStoreError> {
466        let mut inner = self.inner.lock().await;
467        inner.ops_lifecycle_snapshots.remove(&runtime_id.0);
468        Ok(())
469    }
470}
471
472#[cfg(test)]
473#[allow(clippy::unwrap_used)]
474mod tests {
475    use super::*;
476    use crate::store::MachineLifecycleBindingFacts;
477    use meerkat_core::lifecycle::run_primitive::RunApplyBoundary;
478
479    fn make_receipt(run_id: RunId, seq: u64) -> RunBoundaryReceipt {
480        RunBoundaryReceipt {
481            run_id,
482            boundary: RunApplyBoundary::RunStart,
483            contributing_input_ids: vec![],
484            conversation_digest: None,
485            message_count: 0,
486            sequence: seq,
487        }
488    }
489
490    fn persistable(bundle: StoredInputState) -> InputStatePersistenceRecord {
491        InputStatePersistenceRecord::from_machine_snapshot(bundle).unwrap()
492    }
493
494    fn session_with_user(content: &str) -> meerkat_core::Session {
495        let mut session = meerkat_core::Session::new();
496        session.push(meerkat_core::types::Message::User(
497            meerkat_core::types::UserMessage::text(content.to_string()),
498        ));
499        session
500    }
501
502    #[tokio::test]
503    async fn atomic_apply_roundtrip() {
504        let store = InMemoryRuntimeStore::new();
505        let rid = LogicalRuntimeId::new("test-runtime");
506        let run_id = RunId::new();
507        let input_id = InputId::new();
508
509        let bundle = StoredInputState::new_accepted(input_id.clone());
510        let receipt = make_receipt(run_id.clone(), 0);
511
512        let session = session_with_user("hello");
513        let session_snapshot = serde_json::to_vec(&session).unwrap();
514
515        store
516            .atomic_apply(
517                &rid,
518                Some(SessionDelta { session_snapshot }),
519                receipt.clone(),
520                vec![persistable(bundle)],
521                None,
522            )
523            .await
524            .unwrap();
525
526        // Load input states
527        let states = store.load_input_states(&rid).await.unwrap();
528        assert_eq!(states.len(), 1);
529        assert_eq!(states[0].state.input_id, input_id);
530
531        // Load receipt
532        let loaded = store.load_boundary_receipt(&rid, &run_id, 0).await.unwrap();
533        assert!(loaded.is_some());
534    }
535
536    #[tokio::test]
537    async fn atomic_apply_rejects_non_session_snapshot_without_owner_context() {
538        let store = InMemoryRuntimeStore::new();
539        let rid = LogicalRuntimeId::new("test-runtime");
540        let run_id = RunId::new();
541        let input_id = InputId::new();
542
543        let bundle = StoredInputState::new_accepted(input_id);
544        let receipt = make_receipt(run_id, 0);
545
546        // Owner-context absence is not a license to store arbitrary bytes as a
547        // session snapshot: a non-deserializable snapshot must fail closed.
548        let err = store
549            .atomic_apply(
550                &rid,
551                Some(SessionDelta {
552                    session_snapshot: b"session-data".to_vec(),
553                }),
554                receipt,
555                vec![persistable(bundle)],
556                None,
557            )
558            .await
559            .expect_err("non-Session snapshot must be rejected");
560
561        match err {
562            RuntimeStoreError::WriteFailed(message) => {
563                assert!(
564                    message.contains("not a Session"),
565                    "unexpected WriteFailed message: {message}"
566                );
567            }
568            other => panic!("expected WriteFailed, got {other:?}"),
569        }
570    }
571
572    #[tokio::test]
573    async fn persist_and_load_single_state() {
574        let store = InMemoryRuntimeStore::new();
575        let rid = LogicalRuntimeId::new("test");
576        let input_id = InputId::new();
577        let bundle = StoredInputState::new_accepted(input_id.clone());
578
579        store
580            .persist_input_state(&rid, &persistable(bundle))
581            .await
582            .unwrap();
583
584        let loaded = store.load_input_state(&rid, &input_id).await.unwrap();
585        assert!(loaded.is_some());
586        assert_eq!(loaded.unwrap().state.input_id, input_id);
587    }
588
589    #[tokio::test]
590    async fn load_nonexistent_returns_none() {
591        let store = InMemoryRuntimeStore::new();
592        let rid = LogicalRuntimeId::new("test");
593
594        let states = store.load_input_states(&rid).await.unwrap();
595        assert!(states.is_empty());
596
597        let state = store.load_input_state(&rid, &InputId::new()).await.unwrap();
598        assert!(state.is_none());
599
600        let receipt = store
601            .load_boundary_receipt(&rid, &RunId::new(), 0)
602            .await
603            .unwrap();
604        assert!(receipt.is_none());
605    }
606
607    #[tokio::test]
608    async fn atomic_apply_updates_existing() {
609        let store = InMemoryRuntimeStore::new();
610        let rid = LogicalRuntimeId::new("test");
611        let input_id = InputId::new();
612
613        // First write
614        let bundle1 = StoredInputState::new_accepted(input_id.clone());
615        store
616            .atomic_apply(
617                &rid,
618                None,
619                make_receipt(RunId::new(), 0),
620                vec![persistable(bundle1)],
621                None,
622            )
623            .await
624            .unwrap();
625
626        // Second write with updated seed phase
627        let mut bundle2 = StoredInputState::new_accepted(input_id.clone());
628        bundle2.seed.phase = crate::input_state::InputLifecycleState::Queued;
629        store
630            .atomic_apply(
631                &rid,
632                None,
633                make_receipt(RunId::new(), 1),
634                vec![persistable(bundle2)],
635                None,
636            )
637            .await
638            .unwrap();
639
640        let states = store.load_input_states(&rid).await.unwrap();
641        assert_eq!(states.len(), 1);
642        assert_eq!(
643            states[0].seed.phase,
644            crate::input_state::InputLifecycleState::Queued
645        );
646    }
647
648    #[tokio::test]
649    async fn atomic_apply_validates_session_store_key_without_aliasing_snapshot() {
650        let store = InMemoryRuntimeStore::new();
651        let rid = LogicalRuntimeId::new("runtime-key");
652        let session = meerkat_core::Session::new();
653        let session_id = session.id().clone();
654        let snapshot = serde_json::to_vec(&session).unwrap();
655
656        store
657            .atomic_apply(
658                &rid,
659                Some(SessionDelta {
660                    session_snapshot: snapshot.clone(),
661                }),
662                make_receipt(RunId::new(), 0),
663                vec![],
664                Some(session_id.clone()),
665            )
666            .await
667            .unwrap();
668
669        assert_eq!(
670            store.load_session_snapshot(&rid).await.unwrap(),
671            Some(snapshot)
672        );
673        assert!(
674            store
675                .load_session_snapshot(&LogicalRuntimeId::legacy_session_uuid_alias(&session_id))
676                .await
677                .unwrap()
678                .is_none(),
679            "session_store_key must validate the snapshot identity, not create a raw UUID runtime alias"
680        );
681    }
682
683    #[tokio::test]
684    async fn atomic_apply_rejects_mismatched_session_store_key() {
685        let store = InMemoryRuntimeStore::new();
686        let rid = LogicalRuntimeId::new("runtime-key");
687        let session = meerkat_core::Session::new();
688        let wrong_session_id = meerkat_core::Session::new().id().clone();
689        let snapshot = serde_json::to_vec(&session).unwrap();
690
691        let err = store
692            .atomic_apply(
693                &rid,
694                Some(SessionDelta {
695                    session_snapshot: snapshot,
696                }),
697                make_receipt(RunId::new(), 0),
698                vec![],
699                Some(wrong_session_id),
700            )
701            .await
702            .expect_err("mismatched session_store_key should fail");
703
704        assert!(matches!(err, RuntimeStoreError::SessionKeyMismatch { .. }));
705        assert!(store.load_session_snapshot(&rid).await.unwrap().is_none());
706    }
707
708    #[tokio::test]
709    async fn atomic_apply_persists_machine_owned_receipt() {
710        let store = InMemoryRuntimeStore::new();
711        let rid = LogicalRuntimeId::new("test");
712        let run_id = RunId::new();
713        let input_id = InputId::new();
714        let session = meerkat_core::Session::new();
715        let snapshot = serde_json::to_vec(&session).unwrap();
716        let receipt = RunBoundaryReceipt {
717            run_id: run_id.clone(),
718            boundary: RunApplyBoundary::Immediate,
719            contributing_input_ids: vec![input_id.clone()],
720            conversation_digest: Some("machine-owned-digest".to_string()),
721            message_count: 42,
722            sequence: 7,
723        };
724
725        store
726            .atomic_apply(
727                &rid,
728                Some(SessionDelta {
729                    session_snapshot: snapshot,
730                }),
731                receipt.clone(),
732                vec![persistable(StoredInputState::new_accepted(input_id))],
733                None,
734            )
735            .await
736            .unwrap();
737
738        assert_eq!(receipt.run_id, run_id);
739        assert!(receipt.conversation_digest.is_some());
740        let loaded = store
741            .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
742            .await
743            .unwrap();
744        assert!(loaded.is_some(), "receipt should be persisted");
745        let Some(loaded) = loaded else {
746            unreachable!("asserted above");
747        };
748        assert_eq!(loaded, receipt);
749    }
750
751    #[tokio::test]
752    async fn multiple_runtimes_isolated() {
753        let store = InMemoryRuntimeStore::new();
754        let rid1 = LogicalRuntimeId::new("runtime-1");
755        let rid2 = LogicalRuntimeId::new("runtime-2");
756
757        store
758            .persist_input_state(
759                &rid1,
760                &persistable(StoredInputState::new_accepted(InputId::new())),
761            )
762            .await
763            .unwrap();
764        store
765            .persist_input_state(
766                &rid2,
767                &persistable(StoredInputState::new_accepted(InputId::new())),
768            )
769            .await
770            .unwrap();
771        store
772            .persist_input_state(
773                &rid2,
774                &persistable(StoredInputState::new_accepted(InputId::new())),
775            )
776            .await
777            .unwrap();
778
779        let s1 = store.load_input_states(&rid1).await.unwrap();
780        let s2 = store.load_input_states(&rid2).await.unwrap();
781        assert_eq!(s1.len(), 1);
782        assert_eq!(s2.len(), 2);
783    }
784
785    #[tokio::test]
786    async fn load_session_snapshot_roundtrip() {
787        let store = InMemoryRuntimeStore::new();
788        let rid = LogicalRuntimeId::new("runtime");
789        let snapshot = serde_json::to_vec(&meerkat_core::Session::new()).unwrap();
790
791        store
792            .atomic_apply(
793                &rid,
794                Some(SessionDelta {
795                    session_snapshot: snapshot.clone(),
796                }),
797                make_receipt(RunId::new(), 0),
798                vec![],
799                None,
800            )
801            .await
802            .unwrap();
803
804        let loaded = store.load_session_snapshot(&rid).await.unwrap();
805        assert_eq!(loaded, Some(snapshot));
806    }
807
808    #[tokio::test]
809    async fn commit_session_snapshot_rejects_stale_runtime_parent() {
810        let store = InMemoryRuntimeStore::new();
811        let rid = LogicalRuntimeId::new("runtime-stale-parent");
812        let accepted = session_with_user("accepted runtime turn");
813        let mut stale = meerkat_core::Session::with_id(accepted.id().clone());
814        stale.push(meerkat_core::types::Message::User(
815            meerkat_core::types::UserMessage::text("stale runtime turn".to_string()),
816        ));
817        let accepted_snapshot = serde_json::to_vec(&accepted).unwrap();
818
819        store
820            .commit_session_snapshot(
821                &rid,
822                SessionDelta {
823                    session_snapshot: accepted_snapshot.clone(),
824                },
825            )
826            .await
827            .unwrap();
828
829        let err = store
830            .commit_session_snapshot(
831                &rid,
832                SessionDelta {
833                    session_snapshot: serde_json::to_vec(&stale).unwrap(),
834                },
835            )
836            .await
837            .expect_err("stale non-continuation must not overwrite runtime snapshot");
838
839        assert!(matches!(err, RuntimeStoreError::WriteFailed(_)));
840        assert_eq!(
841            store.load_session_snapshot(&rid).await.unwrap(),
842            Some(accepted_snapshot)
843        );
844    }
845
846    #[tokio::test]
847    async fn atomic_apply_keeps_current_snapshot_when_incoming_is_superseded() {
848        let store = InMemoryRuntimeStore::new();
849        let rid = LogicalRuntimeId::new("runtime-superseded-terminal");
850        let incoming = session_with_user("turn input");
851        let mut current = incoming.clone();
852        current.push(meerkat_core::types::Message::BlockAssistant(
853            meerkat_core::types::BlockAssistantMessage {
854                blocks: vec![meerkat_core::types::AssistantBlock::Text {
855                    text: "peer response already applied".to_string(),
856                    meta: None,
857                }],
858                stop_reason: meerkat_core::types::StopReason::EndTurn,
859                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
860                created_at: meerkat_core::types::message_timestamp_now(),
861            },
862        ));
863        let current_snapshot = serde_json::to_vec(&current).unwrap();
864        let receipt = make_receipt(RunId::new(), 11);
865
866        store
867            .commit_session_snapshot(
868                &rid,
869                SessionDelta {
870                    session_snapshot: current_snapshot.clone(),
871                },
872            )
873            .await
874            .unwrap();
875
876        store
877            .atomic_apply(
878                &rid,
879                Some(SessionDelta {
880                    session_snapshot: serde_json::to_vec(&incoming).unwrap(),
881                }),
882                receipt.clone(),
883                vec![],
884                Some(incoming.id().clone()),
885            )
886            .await
887            .unwrap();
888
889        assert_eq!(
890            store.load_session_snapshot(&rid).await.unwrap(),
891            Some(current_snapshot)
892        );
893        // The session snapshot was classified superseded and skipped, so the
894        // boundary receipt for that boundary must NOT advance against the
895        // retained (more-advanced) session snapshot.
896        assert_eq!(
897            store
898                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
899                .await
900                .unwrap(),
901            None
902        );
903    }
904
905    #[tokio::test]
906    async fn atomic_apply_skips_inputs_when_session_snapshot_superseded() {
907        let store = InMemoryRuntimeStore::new();
908        let rid = LogicalRuntimeId::new("runtime-superseded-inputs");
909        let incoming = session_with_user("turn input");
910        let mut current = incoming.clone();
911        current.push(meerkat_core::types::Message::BlockAssistant(
912            meerkat_core::types::BlockAssistantMessage {
913                blocks: vec![meerkat_core::types::AssistantBlock::Text {
914                    text: "peer response already applied".to_string(),
915                    meta: None,
916                }],
917                stop_reason: meerkat_core::types::StopReason::EndTurn,
918                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
919                created_at: meerkat_core::types::message_timestamp_now(),
920            },
921        ));
922        let current_snapshot = serde_json::to_vec(&current).unwrap();
923        let receipt = make_receipt(RunId::new(), 21);
924        let input_id = InputId::new();
925        let bundle = StoredInputState::new_accepted(input_id.clone());
926
927        store
928            .commit_session_snapshot(
929                &rid,
930                SessionDelta {
931                    session_snapshot: current_snapshot.clone(),
932                },
933            )
934            .await
935            .unwrap();
936
937        store
938            .atomic_apply(
939                &rid,
940                Some(SessionDelta {
941                    session_snapshot: serde_json::to_vec(&incoming).unwrap(),
942                }),
943                receipt.clone(),
944                vec![persistable(bundle)],
945                Some(incoming.id().clone()),
946            )
947            .await
948            .unwrap();
949
950        // Snapshot retained, receipt + input-state writes skipped as a unit.
951        assert_eq!(
952            store.load_session_snapshot(&rid).await.unwrap(),
953            Some(current_snapshot)
954        );
955        assert_eq!(
956            store
957                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
958                .await
959                .unwrap(),
960            None
961        );
962        assert!(store.load_input_states(&rid).await.unwrap().is_empty());
963    }
964
965    #[tokio::test]
966    async fn atomic_apply_allows_first_generated_snapshot_after_placeholder() {
967        let store = InMemoryRuntimeStore::new();
968        let rid = LogicalRuntimeId::new("runtime-placeholder");
969        let mut placeholder = meerkat_core::Session::new();
970        placeholder.set_system_prompt("base system".to_string());
971        let mut incoming = meerkat_core::Session::with_id(placeholder.id().clone());
972        incoming.set_system_prompt("base system".to_string());
973        incoming.push(meerkat_core::types::Message::User(
974            meerkat_core::types::UserMessage::text("verbose first turn".to_string()),
975        ));
976        let parent_revision = incoming.transcript_revision().unwrap();
977        incoming
978            .commit_transcript_rewrite(
979                meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
980                vec![meerkat_core::types::Message::User(
981                    meerkat_core::types::UserMessage::compaction_summary(
982                        "[Context compacted] first turn",
983                    ),
984                )],
985                meerkat_core::TranscriptRewriteReason::new("compaction"),
986                Some("meerkat-core".to_string()),
987                Some(parent_revision),
988            )
989            .unwrap();
990        let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
991        let receipt = make_receipt(RunId::new(), 12);
992
993        store
994            .commit_session_snapshot(
995                &rid,
996                SessionDelta {
997                    session_snapshot: serde_json::to_vec(&placeholder).unwrap(),
998                },
999            )
1000            .await
1001            .unwrap();
1002
1003        store
1004            .atomic_apply(
1005                &rid,
1006                Some(SessionDelta {
1007                    session_snapshot: incoming_snapshot.clone(),
1008                }),
1009                receipt.clone(),
1010                vec![],
1011                Some(incoming.id().clone()),
1012            )
1013            .await
1014            .unwrap();
1015
1016        assert_eq!(
1017            store.load_session_snapshot(&rid).await.unwrap(),
1018            Some(incoming_snapshot)
1019        );
1020        assert_eq!(
1021            store
1022                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1023                .await
1024                .unwrap(),
1025            Some(receipt)
1026        );
1027    }
1028
1029    #[tokio::test]
1030    async fn atomic_apply_allows_generated_compaction_before_retained_tail() {
1031        let store = InMemoryRuntimeStore::new();
1032        let rid = LogicalRuntimeId::new("runtime-compaction-tail");
1033        let mut previous = meerkat_core::Session::new();
1034        previous.set_system_prompt("runtime system before context refresh".to_string());
1035        previous.push(meerkat_core::types::Message::User(
1036            meerkat_core::types::UserMessage::text("Turn 1 request".to_string()),
1037        ));
1038        previous.push(meerkat_core::types::Message::BlockAssistant(
1039            meerkat_core::types::BlockAssistantMessage {
1040                blocks: vec![meerkat_core::types::AssistantBlock::Text {
1041                    text: "Turn 1 answer".to_string(),
1042                    meta: None,
1043                }],
1044                stop_reason: meerkat_core::types::StopReason::EndTurn,
1045                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1046                created_at: meerkat_core::types::message_timestamp_now(),
1047            },
1048        ));
1049
1050        let mut incoming = meerkat_core::Session::with_id(previous.id().clone());
1051        incoming.set_system_prompt("runtime system after context refresh".to_string());
1052        incoming.push(meerkat_core::types::Message::User(
1053            meerkat_core::types::UserMessage::text(
1054                "Verbose context that will be compacted".to_string(),
1055            ),
1056        ));
1057        for message in previous.messages()[1..].iter().cloned() {
1058            incoming.push(message);
1059        }
1060        incoming.push(meerkat_core::types::Message::BlockAssistant(
1061            meerkat_core::types::BlockAssistantMessage {
1062                blocks: vec![meerkat_core::types::AssistantBlock::Text {
1063                    text: "Turn 2 generated answer".to_string(),
1064                    meta: None,
1065                }],
1066                stop_reason: meerkat_core::types::StopReason::EndTurn,
1067                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1068                created_at: meerkat_core::types::message_timestamp_now(),
1069            },
1070        ));
1071        let parent_revision = incoming.transcript_revision().unwrap();
1072        incoming
1073            .commit_transcript_rewrite(
1074                meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
1075                vec![meerkat_core::types::Message::User(
1076                    meerkat_core::types::UserMessage::compaction_summary(
1077                        "[Context compacted] Earlier runtime context".to_string(),
1078                    ),
1079                )],
1080                meerkat_core::TranscriptRewriteReason::new("compaction"),
1081                Some("meerkat-core".to_string()),
1082                Some(parent_revision),
1083            )
1084            .unwrap();
1085        let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
1086        let receipt = make_receipt(RunId::new(), 13);
1087
1088        store
1089            .commit_session_snapshot(
1090                &rid,
1091                SessionDelta {
1092                    session_snapshot: serde_json::to_vec(&previous).unwrap(),
1093                },
1094            )
1095            .await
1096            .unwrap();
1097
1098        store
1099            .atomic_apply(
1100                &rid,
1101                Some(SessionDelta {
1102                    session_snapshot: incoming_snapshot.clone(),
1103                }),
1104                receipt.clone(),
1105                vec![],
1106                Some(incoming.id().clone()),
1107            )
1108            .await
1109            .unwrap();
1110
1111        assert_eq!(
1112            store.load_session_snapshot(&rid).await.unwrap(),
1113            Some(incoming_snapshot)
1114        );
1115        assert_eq!(
1116            store
1117                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1118                .await
1119                .unwrap(),
1120            Some(receipt)
1121        );
1122    }
1123
1124    #[tokio::test]
1125    async fn commit_machine_lifecycle_persists_binding_facts() {
1126        use crate::runtime_state::RuntimeState;
1127
1128        let store = InMemoryRuntimeStore::new();
1129        let rid = LogicalRuntimeId::new("runtime-binding");
1130        let binding = MachineLifecycleBindingFacts::new(
1131            Some("rt:session:abc".to_string()),
1132            Some(7),
1133            Some(3),
1134            Some("epoch-1".to_string()),
1135        );
1136
1137        store
1138            .commit_machine_lifecycle(
1139                &rid,
1140                MachineLifecycleCommit::new_with_binding(RuntimeState::Retired, binding.clone()),
1141                &[],
1142            )
1143            .await
1144            .unwrap();
1145
1146        let lifecycle = crate::store::load_machine_lifecycle(&store, &rid)
1147            .await
1148            .unwrap()
1149            .expect("machine lifecycle snapshot");
1150        assert_eq!(lifecycle.runtime_state(), RuntimeState::Retired);
1151        assert_eq!(lifecycle.binding(), &binding);
1152        assert_eq!(
1153            crate::store::load_runtime_state(&store, &rid)
1154                .await
1155                .unwrap(),
1156            Some(RuntimeState::Retired)
1157        );
1158    }
1159
1160    #[tokio::test]
1161    async fn clear_session_snapshot_if_current_sets_quarantine_marker_cleared_on_write() {
1162        let store = InMemoryRuntimeStore::new();
1163        let rid = LogicalRuntimeId::new("runtime-quarantine");
1164        let rejected = serde_json::to_vec(&session_with_user("rejected")).unwrap();
1165
1166        assert!(!store.is_runtime_projection_quarantined(&rid).await.unwrap());
1167        store
1168            .commit_session_snapshot(
1169                &rid,
1170                SessionDelta {
1171                    session_snapshot: rejected.clone(),
1172                },
1173            )
1174            .await
1175            .unwrap();
1176        assert!(
1177            store
1178                .clear_session_snapshot_if_current(&rid, &rejected)
1179                .await
1180                .unwrap()
1181        );
1182        assert!(
1183            store.is_runtime_projection_quarantined(&rid).await.unwrap(),
1184            "clearing the rejected snapshot must record the in-memory quarantine marker"
1185        );
1186
1187        // A live snapshot write reclaims runtime authority and clears the marker.
1188        store
1189            .commit_session_snapshot(
1190                &rid,
1191                SessionDelta {
1192                    session_snapshot: serde_json::to_vec(&session_with_user("revived")).unwrap(),
1193                },
1194            )
1195            .await
1196            .unwrap();
1197        assert!(
1198            !store.is_runtime_projection_quarantined(&rid).await.unwrap(),
1199            "a live snapshot write must clear the in-memory quarantine marker"
1200        );
1201    }
1202}