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#[cfg(test)]
10use std::sync::atomic::{AtomicUsize, Ordering};
11
12use indexmap::IndexMap;
13use meerkat_core::lifecycle::{InputId, RunBoundaryReceipt, RunId};
14#[cfg(not(target_arch = "wasm32"))]
15use tokio::sync::Mutex;
16#[cfg(target_arch = "wasm32")]
17use tokio_with_wasm::alias::sync::Mutex;
18
19use super::{
20    AuthOAuthFlowSnapshotUpdate, FencedInputStateBatchCasOutcome, FencedMachineLifecycleCasOutcome,
21    InputStateBatchCasOutcome, MachineLifecycleCasOutcome, MachineLifecycleCommit,
22    MachineLifecycleExpectedVersion, MachineLifecycleObservation, MachineLifecycleStoreRecord,
23    RuntimeStore, RuntimeStoreError, RuntimeStoreWriteFence, RuntimeStoreWriteFenceOutcome,
24    SessionDelta, classify_machine_lifecycle_record, complete_compaction_projection_checkpoint,
25    decoded_prepared_machine_lifecycle_replacement, execute_runtime_store_write_fence,
26    prepare_input_state_batch_cas, prepare_machine_lifecycle_replacement,
27    validate_machine_lifecycle_replacement,
28};
29use crate::identifiers::LogicalRuntimeId;
30use crate::input_state::{InputStatePersistenceRecord, StoredInputState};
31use crate::ops_lifecycle::PersistedOpsSnapshot;
32
33/// Receipt key: (runtime_id, run_id, sequence).
34#[derive(Debug, Clone, PartialEq, Eq, Hash)]
35struct ReceiptKey {
36    runtime_id: String,
37    run_id: RunId,
38    sequence: u64,
39}
40
41#[derive(Debug, Clone)]
42struct CompactionOutboxEntry {
43    intent: meerkat_core::CompactionProjectionIntent,
44    finalized: bool,
45}
46
47#[cfg(test)]
48type InputStateBatchCasTestBlock = (
49    Arc<crate::tokio::sync::Notify>,
50    Arc<crate::tokio::sync::Notify>,
51);
52
53/// Inner state protected by the mutex.
54#[derive(Debug, Default)]
55struct Inner {
56    /// runtime_id → (input_id → StoredInputState). IndexMap for deterministic iteration order.
57    input_states: HashMap<String, IndexMap<InputId, StoredInputState>>,
58    /// Receipt storage.
59    receipts: HashMap<ReceiptKey, RunBoundaryReceipt>,
60    /// Runtime session snapshots keyed by canonical runtime id.
61    sessions: HashMap<String, Vec<u8>>,
62    /// Canonical runtime ids whose projection fallback is quarantined.
63    ///
64    /// Mirrors the durable SQLite `runtime_projection_quarantine` table: set
65    /// when a rejected runtime snapshot is cleared via
66    /// `clear_session_snapshot_if_current`, cleared whenever a live snapshot is
67    /// written for the runtime.
68    projection_quarantine: HashSet<String>,
69    /// Exact persisted machine-lifecycle bytes. Raw storage is required so
70    /// malformed and unsupported rows remain observable instead of being
71    /// normalized by an eager typed decode.
72    runtime_lifecycle: HashMap<String, Vec<u8>>,
73    /// Persisted ops lifecycle snapshots.
74    ops_lifecycle_snapshots: HashMap<String, PersistedOpsSnapshot>,
75    /// Exact ops epochs retired by atomic unregister finalization. Tombstones
76    /// outlive row deletion so detached callbacks cannot resurrect them.
77    retired_ops_epochs: HashSet<(String, meerkat_core::RuntimeEpochId)>,
78    /// Runtime id -> transcript-rewrite-keyed compaction projection outbox.
79    compaction_projection_outbox:
80        HashMap<String, HashMap<meerkat_core::CompactionProjectionId, CompactionOutboxEntry>>,
81}
82
83/// In-memory runtime store. Thread-safe via `tokio::sync::Mutex`.
84#[derive(Debug, Clone)]
85pub struct InMemoryRuntimeStore {
86    inner: Arc<Mutex<Inner>>,
87    auth_oauth_flow_snapshot: Arc<StdMutex<Option<Vec<u8>>>>,
88    #[cfg(test)]
89    input_state_batch_cas_before: Arc<StdMutex<Option<InputStateBatchCasTestBlock>>>,
90    #[cfg(test)]
91    input_state_batch_cas_after_commit: Arc<StdMutex<Option<InputStateBatchCasTestBlock>>>,
92    #[cfg(test)]
93    machine_lifecycle_cas_conflicts_remaining: Arc<AtomicUsize>,
94    #[cfg(test)]
95    machine_lifecycle_observe_errors_remaining: Arc<AtomicUsize>,
96}
97
98impl InMemoryRuntimeStore {
99    pub fn new() -> Self {
100        Self {
101            inner: Arc::new(Mutex::new(Inner::default())),
102            auth_oauth_flow_snapshot: Arc::new(StdMutex::new(None)),
103            #[cfg(test)]
104            input_state_batch_cas_before: Arc::new(StdMutex::new(None)),
105            #[cfg(test)]
106            input_state_batch_cas_after_commit: Arc::new(StdMutex::new(None)),
107            #[cfg(test)]
108            machine_lifecycle_cas_conflicts_remaining: Arc::new(AtomicUsize::new(0)),
109            #[cfg(test)]
110            machine_lifecycle_observe_errors_remaining: Arc::new(AtomicUsize::new(0)),
111        }
112    }
113
114    #[cfg(test)]
115    pub(crate) fn block_next_input_state_batch_cas_before_mutation(
116        &self,
117        entered: Arc<crate::tokio::sync::Notify>,
118        release: Arc<crate::tokio::sync::Notify>,
119    ) {
120        *self
121            .input_state_batch_cas_before
122            .lock()
123            .unwrap_or_else(std::sync::PoisonError::into_inner) = Some((entered, release));
124    }
125
126    #[cfg(test)]
127    pub(crate) fn block_next_input_state_batch_cas_after_commit(
128        &self,
129        entered: Arc<crate::tokio::sync::Notify>,
130        release: Arc<crate::tokio::sync::Notify>,
131    ) {
132        *self
133            .input_state_batch_cas_after_commit
134            .lock()
135            .unwrap_or_else(std::sync::PoisonError::into_inner) = Some((entered, release));
136    }
137
138    #[cfg(test)]
139    pub(crate) async fn seed_machine_lifecycle_raw(
140        &self,
141        runtime_id: &LogicalRuntimeId,
142        bytes: Vec<u8>,
143    ) {
144        self.inner
145            .lock()
146            .await
147            .runtime_lifecycle
148            .insert(runtime_id.0.clone(), bytes);
149    }
150
151    #[cfg(test)]
152    pub(crate) fn conflict_next_machine_lifecycle_cas(&self) {
153        self.machine_lifecycle_cas_conflicts_remaining
154            .fetch_add(1, Ordering::SeqCst);
155    }
156
157    #[cfg(test)]
158    pub(crate) fn fail_next_machine_lifecycle_observation(&self) {
159        self.machine_lifecycle_observe_errors_remaining
160            .fetch_add(1, Ordering::SeqCst);
161    }
162}
163
164impl Default for InMemoryRuntimeStore {
165    fn default() -> Self {
166        Self::new()
167    }
168}
169
170fn is_runtime_placeholder_session(session: &meerkat_core::Session) -> bool {
171    session.transcript_history_state().ok().flatten().is_none()
172        && matches!(
173            session.messages(),
174            [] | [meerkat_core::types::Message::System(_)]
175        )
176}
177
178/// Deserialize a persisted session-snapshot blob through typed serde, matching
179/// the SQLite runtime store read path. `Session::deserialize` validates the
180/// mandatory envelope version against the generated persistence version
181/// authority, so a missing or non-current (v0/v1) row fails closed instead of
182/// silently defaulting or upgrading on read.
183fn deserialize_persisted_session(bytes: &[u8]) -> Result<meerkat_core::Session, RuntimeStoreError> {
184    serde_json::from_slice(bytes).map_err(|err| RuntimeStoreError::ReadFailed(err.to_string()))
185}
186
187fn ensure_compaction_intents_already_outboxed(
188    inner: &Inner,
189    runtime_id: &LogicalRuntimeId,
190    session: &meerkat_core::Session,
191) -> Result<(), RuntimeStoreError> {
192    let intents = super::validated_compaction_projection_intents(session)?;
193    let existing = inner.compaction_projection_outbox.get(&runtime_id.0);
194    for intent in intents {
195        match existing.and_then(|entries| entries.get(&intent.projection)) {
196            Some(entry) if entry.finalized => {
197                return Err(RuntimeStoreError::WriteFailed(format!(
198                    "non-boundary snapshot replays finalized compaction intent {}",
199                    intent.projection.revision()
200                )));
201            }
202            Some(entry) if entry.intent == intent => {}
203            Some(_) => {
204                return Err(RuntimeStoreError::WriteFailed(format!(
205                    "non-boundary snapshot conflicts with compaction outbox rewrite {}",
206                    intent.projection.revision()
207                )));
208            }
209            None => {
210                return Err(RuntimeStoreError::WriteFailed(format!(
211                    "non-boundary snapshot introduces compaction intent {} without atomic outbox authority",
212                    intent.projection.revision()
213                )));
214            }
215        }
216    }
217    Ok(())
218}
219
220#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
221#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
222impl RuntimeStore for InMemoryRuntimeStore {
223    fn supports_compaction_projection_outbox(&self) -> bool {
224        true
225    }
226
227    fn persist_auth_oauth_flow_snapshot(
228        &self,
229        snapshot_json: &[u8],
230    ) -> Result<(), RuntimeStoreError> {
231        *self
232            .auth_oauth_flow_snapshot
233            .lock()
234            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))? =
235            Some(snapshot_json.to_vec());
236        Ok(())
237    }
238
239    fn load_auth_oauth_flow_snapshot(&self) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
240        self.auth_oauth_flow_snapshot
241            .lock()
242            .map(|snapshot| snapshot.clone())
243            .map_err(|err| RuntimeStoreError::ReadFailed(err.to_string()))
244    }
245
246    fn update_auth_oauth_flow_snapshot(
247        &self,
248        update: &mut AuthOAuthFlowSnapshotUpdate<'_>,
249    ) -> Result<(), RuntimeStoreError> {
250        let mut snapshot = self
251            .auth_oauth_flow_snapshot
252            .lock()
253            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
254        let next = update(snapshot.as_deref())?;
255        *snapshot = Some(next);
256        Ok(())
257    }
258
259    async fn commit_session_snapshot(
260        &self,
261        runtime_id: &LogicalRuntimeId,
262        session_delta: SessionDelta,
263    ) -> Result<(), RuntimeStoreError> {
264        let incoming: meerkat_core::Session =
265            serde_json::from_slice(&session_delta.session_snapshot)
266                .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
267        let mut inner = self.inner.lock().await;
268        ensure_compaction_intents_already_outboxed(&inner, runtime_id, &incoming)?;
269        if inner.sessions.get(&runtime_id.0).is_some_and(|snapshot| {
270            snapshot.as_slice() == session_delta.session_snapshot.as_slice()
271        }) {
272            // The incoming document still crossed typed Session validation
273            // and compaction-intent authority above. Exact byte identity now
274            // proves there is no prior Session to parse or map value to
275            // replace. The self-guard preserves live-head coherence and every
276            // fail-closed save invariant before the fast return.
277            meerkat_core::session_store::run_boundary_snapshot_head_coherence_guard(&incoming)
278                .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
279            inner.projection_quarantine.remove(&runtime_id.0);
280            return Ok(());
281        }
282        let previous = inner
283            .sessions
284            .get(&runtime_id.0)
285            .map(|snapshot| deserialize_persisted_session(snapshot))
286            .transpose()?;
287        meerkat_core::session_store::run_boundary_snapshot_save_guard(&incoming, previous.as_ref())
288            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
289        inner
290            .sessions
291            .insert(runtime_id.0.clone(), session_delta.session_snapshot);
292        inner.projection_quarantine.remove(&runtime_id.0);
293        Ok(())
294    }
295
296    async fn commit_session_transcript_rewrite_snapshot(
297        &self,
298        runtime_id: &LogicalRuntimeId,
299        session_delta: SessionDelta,
300        commit: &meerkat_core::TranscriptRewriteCommit,
301    ) -> Result<(), RuntimeStoreError> {
302        let incoming: meerkat_core::Session =
303            serde_json::from_slice(&session_delta.session_snapshot)
304                .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
305        let mut inner = self.inner.lock().await;
306        ensure_compaction_intents_already_outboxed(&inner, runtime_id, &incoming)?;
307        let previous = inner
308            .sessions
309            .get(&runtime_id.0)
310            .map(|snapshot| deserialize_persisted_session(snapshot))
311            .transpose()?;
312        meerkat_core::session_store::transcript_rewrite_save_guard(
313            &incoming,
314            previous.as_ref(),
315            commit,
316        )
317        .map_err(|err| match err {
318            meerkat_core::SessionStoreError::TranscriptRevisionConflict {
319                expected,
320                actual,
321                ..
322            } => RuntimeStoreError::TranscriptRevisionConflict { expected, actual },
323            other => RuntimeStoreError::WriteFailed(other.to_string()),
324        })?;
325        inner
326            .sessions
327            .insert(runtime_id.0.clone(), session_delta.session_snapshot);
328        inner.projection_quarantine.remove(&runtime_id.0);
329        Ok(())
330    }
331
332    async fn atomic_apply(
333        &self,
334        runtime_id: &LogicalRuntimeId,
335        session_delta: Option<SessionDelta>,
336        receipt: RunBoundaryReceipt,
337        input_updates: Vec<InputStatePersistenceRecord>,
338        session_store_key: Option<meerkat_core::types::SessionId>,
339    ) -> Result<(), RuntimeStoreError> {
340        let mut inner = self.inner.lock().await;
341
342        // All writes in one lock acquisition (atomic for in-memory)
343        let rid = runtime_id.0.clone();
344
345        // Session delta. The supersession verdict computed here keys the
346        // entire commit: if the incoming session snapshot is classified as
347        // superseded (the persisted head is already a valid append-extension
348        // of it), the snapshot write is skipped AND so are the receipt + input
349        // writes, so receipt/input ordering identity never advances past the
350        // retained session truth.
351        let mut session_snapshot_superseded = false;
352        let mut compaction_intents = Vec::new();
353        let mut session_snapshot_to_persist = None;
354        if let Some(delta) = session_delta {
355            let incoming_session =
356                serde_json::from_slice::<meerkat_core::Session>(&delta.session_snapshot);
357            let mut persist_session_snapshot = true;
358            match (incoming_session, session_store_key) {
359                (Ok(incoming_session), session_store_key) => {
360                    compaction_intents =
361                        super::validated_compaction_projection_intents(&incoming_session)?;
362                    if let Some(existing) = inner.compaction_projection_outbox.get(&rid) {
363                        for intent in &compaction_intents {
364                            if let Some(entry) = existing.get(&intent.projection) {
365                                if entry.finalized {
366                                    return Err(RuntimeStoreError::WriteFailed(format!(
367                                        "atomic session snapshot replays finalized compaction intent {}",
368                                        intent.projection.revision()
369                                    )));
370                                }
371                                if entry.intent != *intent {
372                                    return Err(RuntimeStoreError::WriteFailed(format!(
373                                        "conflicting compaction outbox intent for rewrite {}",
374                                        intent.projection.revision()
375                                    )));
376                                }
377                            }
378                        }
379                    }
380                    if let Some(session_store_key) = session_store_key
381                        && incoming_session.id() != &session_store_key
382                    {
383                        return Err(RuntimeStoreError::SessionKeyMismatch {
384                            expected: session_store_key,
385                            actual: incoming_session.id().clone(),
386                        });
387                    }
388                    let previous_session = inner
389                        .sessions
390                        .get(&rid)
391                        .map(|snapshot| deserialize_persisted_session(snapshot))
392                        .transpose()?;
393                    if let Err(err) = meerkat_core::session_store::run_boundary_snapshot_save_guard(
394                        &incoming_session,
395                        previous_session.as_ref(),
396                    ) {
397                        if previous_session
398                            .as_ref()
399                            .is_some_and(is_runtime_placeholder_session)
400                        {
401                            persist_session_snapshot = true;
402                        } else if previous_session.as_ref().is_some_and(|previous_session| {
403                            meerkat_core::session_store::run_boundary_snapshot_save_guard(
404                                previous_session,
405                                Some(&incoming_session),
406                            )
407                            .is_ok()
408                        }) {
409                            persist_session_snapshot = false;
410                            session_snapshot_superseded = true;
411                        } else {
412                            return Err(RuntimeStoreError::WriteFailed(err.to_string()));
413                        }
414                    }
415                }
416                (Err(err), Some(session_store_key)) => {
417                    return Err(RuntimeStoreError::WriteFailed(format!(
418                        "session snapshot for {session_store_key} is not a Session: {err}"
419                    )));
420                }
421                (Err(err), None) => {
422                    return Err(RuntimeStoreError::WriteFailed(format!(
423                        "session snapshot is not a Session: {err}"
424                    )));
425                }
426            }
427            if persist_session_snapshot {
428                session_snapshot_to_persist = Some(delta.session_snapshot);
429            }
430        }
431
432        // When the session snapshot was superseded and skipped, the boundary
433        // receipt and input-state updates for that boundary must also be
434        // skipped: advancing them against a retained (older) session snapshot
435        // would split receipt/input ordering identity from session truth.
436        if session_snapshot_superseded {
437            return Err(RuntimeStoreError::SessionSnapshotSuperseded { runtime_id: rid });
438        }
439
440        // Receipt immutability is validated before any mutation so a stale or
441        // duplicate writer cannot partially advance the session snapshot and
442        // then overwrite a prior boundary identity. This mirrors SQLite's
443        // primary-key INSERT behavior.
444        let key = ReceiptKey {
445            runtime_id: rid.clone(),
446            run_id: receipt.run_id.clone(),
447            sequence: receipt.sequence,
448        };
449        if inner.receipts.contains_key(&key) {
450            return Err(RuntimeStoreError::WriteFailed(format!(
451                "boundary receipt already exists for runtime '{}' run {} sequence {}",
452                runtime_id, receipt.run_id, receipt.sequence
453            )));
454        }
455
456        let outbox = inner
457            .compaction_projection_outbox
458            .entry(rid.clone())
459            .or_default();
460        for intent in compaction_intents {
461            outbox
462                .entry(intent.projection.clone())
463                .or_insert(CompactionOutboxEntry {
464                    intent,
465                    finalized: false,
466                });
467        }
468
469        if let Some(session_snapshot) = session_snapshot_to_persist {
470            inner.sessions.insert(rid.clone(), session_snapshot);
471            inner.projection_quarantine.remove(&rid);
472        }
473        inner.receipts.insert(key, receipt);
474
475        // Input states
476        let states = inner.input_states.entry(rid).or_default();
477        for record in input_updates {
478            let bundle = record.into_stored();
479            states.insert(bundle.state.input_id.clone(), bundle);
480        }
481
482        Ok(())
483    }
484
485    async fn load_pending_compaction_projections(
486        &self,
487        runtime_id: &LogicalRuntimeId,
488    ) -> Result<Vec<meerkat_core::CompactionProjectionIntent>, RuntimeStoreError> {
489        let inner = self.inner.lock().await;
490        let mut pending = inner
491            .compaction_projection_outbox
492            .get(&runtime_id.0)
493            .into_iter()
494            .flat_map(HashMap::values)
495            .filter(|entry| !entry.finalized)
496            .map(|entry| entry.intent.clone())
497            .collect::<Vec<_>>();
498        pending.sort_by(|left, right| {
499            left.projection
500                .session_id()
501                .to_string()
502                .cmp(&right.projection.session_id().to_string())
503                .then_with(|| {
504                    left.projection
505                        .parent_revision()
506                        .cmp(right.projection.parent_revision())
507                })
508                .then_with(|| left.projection.revision().cmp(right.projection.revision()))
509                .then_with(|| {
510                    left.projection
511                        .commit_fingerprint()
512                        .cmp(right.projection.commit_fingerprint())
513                })
514        });
515        Ok(pending)
516    }
517
518    async fn mark_compaction_projection_finalized(
519        &self,
520        runtime_id: &LogicalRuntimeId,
521        projection: &meerkat_core::CompactionProjectionId,
522    ) -> Result<(), RuntimeStoreError> {
523        let mut inner = self.inner.lock().await;
524        let outbox_exists = inner
525            .compaction_projection_outbox
526            .get(&runtime_id.0)
527            .is_some_and(|entries| entries.contains_key(projection));
528        if !outbox_exists {
529            return Err(RuntimeStoreError::NotFound(format!(
530                "compaction outbox rewrite {}",
531                projection.revision()
532            )));
533        }
534        let cleaned_snapshot = inner
535            .sessions
536            .get(&runtime_id.0)
537            .map(|snapshot| {
538                let mut session = deserialize_persisted_session(snapshot)?;
539                complete_compaction_projection_checkpoint(&mut session, projection)?;
540                serde_json::to_vec(&session)
541                    .map_err(|error| RuntimeStoreError::WriteFailed(error.to_string()))
542            })
543            .transpose()?;
544        let entry = inner
545            .compaction_projection_outbox
546            .get_mut(&runtime_id.0)
547            .and_then(|entries| entries.get_mut(projection))
548            .ok_or_else(|| {
549                RuntimeStoreError::NotFound(format!(
550                    "compaction outbox rewrite {}",
551                    projection.revision()
552                ))
553            })?;
554        entry.finalized = true;
555        if let Some(cleaned_snapshot) = cleaned_snapshot {
556            inner
557                .sessions
558                .insert(runtime_id.0.clone(), cleaned_snapshot);
559        }
560        Ok(())
561    }
562
563    async fn atomic_apply_with_machine_lifecycle(
564        &self,
565        runtime_id: &LogicalRuntimeId,
566        session_delta: SessionDelta,
567        receipt: RunBoundaryReceipt,
568        machine_lifecycle: MachineLifecycleCommit,
569        input_updates: Vec<InputStatePersistenceRecord>,
570        session_store_key: meerkat_core::types::SessionId,
571    ) -> Result<(), RuntimeStoreError> {
572        let machine_lifecycle_record = machine_lifecycle.store_record().encode()?;
573        let mut inner = self.inner.lock().await;
574        let rid = runtime_id.0.clone();
575        let incoming_session =
576            serde_json::from_slice::<meerkat_core::Session>(&session_delta.session_snapshot)
577                .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
578        let compaction_intents = super::validated_compaction_projection_intents(&incoming_session)?;
579        if let Some(existing) = inner.compaction_projection_outbox.get(&rid) {
580            for intent in &compaction_intents {
581                if let Some(entry) = existing.get(&intent.projection) {
582                    if entry.finalized {
583                        return Err(RuntimeStoreError::WriteFailed(format!(
584                            "atomic session snapshot replays finalized compaction intent {}",
585                            intent.projection.revision()
586                        )));
587                    }
588                    if entry.intent != *intent {
589                        return Err(RuntimeStoreError::WriteFailed(format!(
590                            "conflicting compaction outbox intent for rewrite {}",
591                            intent.projection.revision()
592                        )));
593                    }
594                }
595            }
596        }
597        if incoming_session.id() != &session_store_key {
598            return Err(RuntimeStoreError::SessionKeyMismatch {
599                expected: session_store_key,
600                actual: incoming_session.id().clone(),
601            });
602        }
603        let previous_session = inner
604            .sessions
605            .get(&rid)
606            .map(|snapshot| deserialize_persisted_session(snapshot))
607            .transpose()?;
608        if let Err(err) = meerkat_core::session_store::run_boundary_snapshot_save_guard(
609            &incoming_session,
610            previous_session.as_ref(),
611        ) {
612            if previous_session
613                .as_ref()
614                .is_some_and(is_runtime_placeholder_session)
615            {
616                // The first generated transcript replaces its runtime
617                // placeholder atomically with all machine-terminal state.
618            } else if previous_session.as_ref().is_some_and(|previous_session| {
619                meerkat_core::session_store::run_boundary_snapshot_save_guard(
620                    previous_session,
621                    Some(&incoming_session),
622                )
623                .is_ok()
624            }) {
625                return Err(RuntimeStoreError::SessionSnapshotSuperseded { runtime_id: rid });
626            } else {
627                return Err(RuntimeStoreError::WriteFailed(err.to_string()));
628            }
629        }
630
631        let key = ReceiptKey {
632            runtime_id: rid.clone(),
633            run_id: receipt.run_id.clone(),
634            sequence: receipt.sequence,
635        };
636        if inner.receipts.contains_key(&key) {
637            return Err(RuntimeStoreError::WriteFailed(format!(
638                "boundary receipt already exists for runtime '{}' run {} sequence {}",
639                runtime_id, receipt.run_id, receipt.sequence
640            )));
641        }
642
643        let outbox = inner
644            .compaction_projection_outbox
645            .entry(rid.clone())
646            .or_default();
647        for intent in compaction_intents {
648            outbox
649                .entry(intent.projection.clone())
650                .or_insert(CompactionOutboxEntry {
651                    intent,
652                    finalized: false,
653                });
654        }
655        inner
656            .sessions
657            .insert(rid.clone(), session_delta.session_snapshot);
658        inner.projection_quarantine.remove(&rid);
659        inner
660            .runtime_lifecycle
661            .insert(rid.clone(), machine_lifecycle_record);
662        inner.receipts.insert(key, receipt);
663        let states = inner.input_states.entry(rid).or_default();
664        for record in input_updates {
665            let bundle = record.into_stored();
666            states.insert(bundle.state.input_id.clone(), bundle);
667        }
668        Ok(())
669    }
670
671    async fn load_input_states(
672        &self,
673        runtime_id: &LogicalRuntimeId,
674    ) -> Result<Vec<StoredInputState>, RuntimeStoreError> {
675        let inner = self.inner.lock().await;
676        let states = inner
677            .input_states
678            .get(&runtime_id.0)
679            .map(|m| m.values().cloned().collect())
680            .unwrap_or_default();
681        Ok(states)
682    }
683
684    async fn load_boundary_receipt(
685        &self,
686        runtime_id: &LogicalRuntimeId,
687        run_id: &RunId,
688        sequence: u64,
689    ) -> Result<Option<RunBoundaryReceipt>, RuntimeStoreError> {
690        let inner = self.inner.lock().await;
691        let key = ReceiptKey {
692            runtime_id: runtime_id.0.clone(),
693            run_id: run_id.clone(),
694            sequence,
695        };
696        Ok(inner.receipts.get(&key).cloned())
697    }
698
699    async fn load_session_snapshot(
700        &self,
701        runtime_id: &LogicalRuntimeId,
702    ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
703        let inner = self.inner.lock().await;
704        Ok(inner.sessions.get(&runtime_id.0).cloned())
705    }
706
707    async fn clear_session_snapshot(
708        &self,
709        runtime_id: &LogicalRuntimeId,
710    ) -> Result<(), RuntimeStoreError> {
711        let mut inner = self.inner.lock().await;
712        inner.sessions.remove(&runtime_id.0);
713        Ok(())
714    }
715
716    async fn replace_session_snapshot_if_current(
717        &self,
718        runtime_id: &LogicalRuntimeId,
719        expected_current: &[u8],
720        replacement: Vec<u8>,
721    ) -> Result<bool, RuntimeStoreError> {
722        let replacement_session: meerkat_core::Session = serde_json::from_slice(&replacement)
723            .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
724        let mut inner = self.inner.lock().await;
725        let Some(current) = inner.sessions.get(&runtime_id.0) else {
726            return Ok(false);
727        };
728        if current.as_slice() != expected_current {
729            return Ok(false);
730        }
731        ensure_compaction_intents_already_outboxed(&inner, runtime_id, &replacement_session)?;
732        inner.sessions.insert(runtime_id.0.clone(), replacement);
733        inner.projection_quarantine.remove(&runtime_id.0);
734        Ok(true)
735    }
736
737    async fn clear_session_snapshot_if_current(
738        &self,
739        runtime_id: &LogicalRuntimeId,
740        expected_current: &[u8],
741    ) -> Result<bool, RuntimeStoreError> {
742        let mut inner = self.inner.lock().await;
743        let Some(current) = inner.sessions.get(&runtime_id.0) else {
744            return Ok(false);
745        };
746        if current.as_slice() != expected_current {
747            return Ok(false);
748        }
749        inner.sessions.remove(&runtime_id.0);
750        // Record the in-memory quarantine marker atomically with the snapshot
751        // removal, mirroring the durable SQLite path.
752        inner.projection_quarantine.insert(runtime_id.0.clone());
753        Ok(true)
754    }
755
756    async fn is_runtime_projection_quarantined(
757        &self,
758        runtime_id: &LogicalRuntimeId,
759    ) -> Result<bool, RuntimeStoreError> {
760        let inner = self.inner.lock().await;
761        Ok(inner.projection_quarantine.contains(&runtime_id.0))
762    }
763
764    async fn persist_input_state(
765        &self,
766        runtime_id: &LogicalRuntimeId,
767        state: &InputStatePersistenceRecord,
768    ) -> Result<(), RuntimeStoreError> {
769        let mut inner = self.inner.lock().await;
770        let states = inner.input_states.entry(runtime_id.0.clone()).or_default();
771        let bundle = state.as_stored();
772        states.insert(bundle.state.input_id.clone(), bundle.clone());
773        Ok(())
774    }
775
776    async fn persist_input_states_atomically(
777        &self,
778        runtime_id: &LogicalRuntimeId,
779        records: &[InputStatePersistenceRecord],
780    ) -> Result<(), RuntimeStoreError> {
781        let mut inner = self.inner.lock().await;
782        let states = inner.input_states.entry(runtime_id.0.clone()).or_default();
783        for record in records {
784            let bundle = record.as_stored();
785            states.insert(bundle.state.input_id.clone(), bundle.clone());
786        }
787        Ok(())
788    }
789
790    async fn compare_and_swap_input_states_atomically(
791        &self,
792        runtime_id: &LogicalRuntimeId,
793        expected: &[StoredInputState],
794        replacements: &[InputStatePersistenceRecord],
795    ) -> Result<InputStateBatchCasOutcome, RuntimeStoreError> {
796        // Serialize and validate the full request before taking the mutation
797        // lock, so no fallible request preparation can occur after writes.
798        let prepared = prepare_input_state_batch_cas(expected, replacements)?;
799        if prepared.is_empty() {
800            return Ok(InputStateBatchCasOutcome::Swapped);
801        }
802
803        #[cfg(test)]
804        let before_block = {
805            self.input_state_batch_cas_before
806                .lock()
807                .unwrap_or_else(std::sync::PoisonError::into_inner)
808                .take()
809        };
810        #[cfg(test)]
811        if let Some((entered, release)) = before_block {
812            entered.notify_one();
813            release.notified().await;
814        }
815
816        let mut inner = self.inner.lock().await;
817        let Some(states) = inner.input_states.get_mut(&runtime_id.0) else {
818            return Ok(InputStateBatchCasOutcome::Stale);
819        };
820        let mut all_expected = true;
821        let mut all_replacements = true;
822        for row in &prepared {
823            let Some(current) = states.get(&row.input_id) else {
824                return Ok(InputStateBatchCasOutcome::Stale);
825            };
826            let current_json = serde_json::to_vec(current)
827                .map_err(|error| RuntimeStoreError::ReadFailed(error.to_string()))?;
828            if current_json != row.expected_json {
829                all_expected = false;
830            }
831            if current_json != row.replacement_json {
832                all_replacements = false;
833            }
834        }
835        if all_replacements {
836            return Ok(InputStateBatchCasOutcome::Swapped);
837        }
838        if !all_expected {
839            return Ok(InputStateBatchCasOutcome::Stale);
840        }
841        for row in prepared {
842            states.insert(row.input_id, row.replacement);
843        }
844        drop(inner);
845
846        #[cfg(test)]
847        let after_commit_block = {
848            self.input_state_batch_cas_after_commit
849                .lock()
850                .unwrap_or_else(std::sync::PoisonError::into_inner)
851                .take()
852        };
853        #[cfg(test)]
854        if let Some((entered, release)) = after_commit_block {
855            entered.notify_one();
856            release.notified().await;
857        }
858        Ok(InputStateBatchCasOutcome::Swapped)
859    }
860
861    async fn compare_and_swap_input_states_atomically_with_fence(
862        &self,
863        runtime_id: &LogicalRuntimeId,
864        expected: &[StoredInputState],
865        replacements: &[InputStatePersistenceRecord],
866        write_fence: Arc<dyn RuntimeStoreWriteFence>,
867    ) -> Result<FencedInputStateBatchCasOutcome, RuntimeStoreError> {
868        let prepared = prepare_input_state_batch_cas(expected, replacements)?;
869        if prepared.is_empty() {
870            return Ok(FencedInputStateBatchCasOutcome::Swapped);
871        }
872
873        let mut inner = self.inner.lock().await;
874        let Some(states) = inner.input_states.get_mut(&runtime_id.0) else {
875            return Ok(FencedInputStateBatchCasOutcome::Stale);
876        };
877        let mut all_expected = true;
878        let mut all_replacements = true;
879        for row in &prepared {
880            let Some(current) = states.get(&row.input_id) else {
881                return Ok(FencedInputStateBatchCasOutcome::Stale);
882            };
883            let current_json = serde_json::to_vec(current)
884                .map_err(|error| RuntimeStoreError::ReadFailed(error.to_string()))?;
885            if current_json != row.expected_json {
886                all_expected = false;
887            }
888            if current_json != row.replacement_json {
889                all_replacements = false;
890            }
891        }
892        if !all_replacements && !all_expected {
893            return Ok(FencedInputStateBatchCasOutcome::Stale);
894        }
895
896        let fence_outcome = execute_runtime_store_write_fence(write_fence.as_ref(), || {
897            if !all_replacements {
898                for row in &prepared {
899                    states.insert(row.input_id.clone(), row.replacement.clone());
900                }
901            }
902            Ok(())
903        })?;
904        match fence_outcome {
905            RuntimeStoreWriteFenceOutcome::Applied => Ok(FencedInputStateBatchCasOutcome::Swapped),
906            RuntimeStoreWriteFenceOutcome::Conflict { reason } => {
907                Ok(FencedInputStateBatchCasOutcome::FenceConflict { reason })
908            }
909            RuntimeStoreWriteFenceOutcome::Backoff { reason } => {
910                Ok(FencedInputStateBatchCasOutcome::FenceBackoff { reason })
911            }
912        }
913    }
914
915    async fn load_input_state(
916        &self,
917        runtime_id: &LogicalRuntimeId,
918        input_id: &InputId,
919    ) -> Result<Option<StoredInputState>, RuntimeStoreError> {
920        let inner = self.inner.lock().await;
921        let state = inner
922            .input_states
923            .get(&runtime_id.0)
924            .and_then(|m| m.get(input_id).cloned());
925        Ok(state)
926    }
927
928    async fn observe_machine_lifecycle(
929        &self,
930        runtime_id: &LogicalRuntimeId,
931    ) -> Result<MachineLifecycleObservation, RuntimeStoreError> {
932        #[cfg(test)]
933        if self
934            .machine_lifecycle_observe_errors_remaining
935            .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |remaining| {
936                remaining.checked_sub(1)
937            })
938            .is_ok()
939        {
940            return Err(RuntimeStoreError::ReadFailed(
941                "synthetic machine lifecycle transport failure".to_string(),
942            ));
943        }
944        let inner = self.inner.lock().await;
945        Ok(inner
946            .runtime_lifecycle
947            .get(&runtime_id.0)
948            .map_or(MachineLifecycleObservation::Missing, |bytes| {
949                classify_machine_lifecycle_record(bytes)
950            }))
951    }
952
953    async fn compare_and_swap_machine_lifecycle(
954        &self,
955        runtime_id: &LogicalRuntimeId,
956        expected: MachineLifecycleExpectedVersion,
957        replacement: MachineLifecycleCommit,
958    ) -> Result<MachineLifecycleCasOutcome, RuntimeStoreError> {
959        let replacement = prepare_machine_lifecycle_replacement(replacement)?;
960        let mut inner = self.inner.lock().await;
961        let current_raw = inner.runtime_lifecycle.get(&runtime_id.0).cloned();
962        let current = current_raw.as_deref().map_or(
963            MachineLifecycleObservation::Missing,
964            classify_machine_lifecycle_record,
965        );
966        #[cfg(test)]
967        if self
968            .machine_lifecycle_cas_conflicts_remaining
969            .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |remaining| {
970                remaining.checked_sub(1)
971            })
972            .is_ok()
973        {
974            return Ok(MachineLifecycleCasOutcome::Conflict { current });
975        }
976        let matches = match (&expected, &current) {
977            (MachineLifecycleExpectedVersion::Missing, MachineLifecycleObservation::Missing) => {
978                true
979            }
980            (MachineLifecycleExpectedVersion::Version(expected), current) => {
981                current.version().is_some_and(|actual| actual == expected)
982            }
983            _ => false,
984        };
985        if !matches {
986            return Ok(MachineLifecycleCasOutcome::Conflict { current });
987        }
988        let replacement = replacement.preserve_observed_custody(&current)?;
989        validate_machine_lifecycle_replacement(
990            &current,
991            current_raw.as_deref(),
992            &replacement.snapshot,
993        )?;
994        inner
995            .runtime_lifecycle
996            .insert(runtime_id.0.clone(), replacement.bytes);
997        Ok(MachineLifecycleCasOutcome::Applied {
998            version: replacement.version,
999        })
1000    }
1001
1002    async fn compare_and_swap_machine_lifecycle_with_fence(
1003        &self,
1004        runtime_id: &LogicalRuntimeId,
1005        expected: MachineLifecycleExpectedVersion,
1006        replacement: MachineLifecycleCommit,
1007        write_fence: Arc<dyn RuntimeStoreWriteFence>,
1008    ) -> Result<FencedMachineLifecycleCasOutcome, RuntimeStoreError> {
1009        let replacement = prepare_machine_lifecycle_replacement(replacement)?;
1010        let mut inner = self.inner.lock().await;
1011        let current_raw = inner.runtime_lifecycle.get(&runtime_id.0).cloned();
1012        let current = current_raw.as_deref().map_or(
1013            MachineLifecycleObservation::Missing,
1014            classify_machine_lifecycle_record,
1015        );
1016        let matches = match (&expected, &current) {
1017            (MachineLifecycleExpectedVersion::Missing, MachineLifecycleObservation::Missing) => {
1018                true
1019            }
1020            (MachineLifecycleExpectedVersion::Version(expected), current) => {
1021                current.version().is_some_and(|actual| actual == expected)
1022            }
1023            _ => false,
1024        };
1025        if !matches {
1026            return Ok(FencedMachineLifecycleCasOutcome::Conflict { current });
1027        }
1028        let replacement = replacement.preserve_observed_custody(&current)?;
1029        validate_machine_lifecycle_replacement(
1030            &current,
1031            current_raw.as_deref(),
1032            &replacement.snapshot,
1033        )?;
1034        let already_exact = current_raw.as_deref() == Some(replacement.bytes.as_slice());
1035        let record = decoded_prepared_machine_lifecycle_replacement(&replacement)?;
1036        let version = replacement.version.clone();
1037        let fence_outcome = execute_runtime_store_write_fence(write_fence.as_ref(), || {
1038            if !already_exact {
1039                inner
1040                    .runtime_lifecycle
1041                    .insert(runtime_id.0.clone(), replacement.bytes.clone());
1042            }
1043            Ok(())
1044        })?;
1045        match fence_outcome {
1046            RuntimeStoreWriteFenceOutcome::Applied if already_exact => {
1047                Ok(FencedMachineLifecycleCasOutcome::AlreadyExact { record, version })
1048            }
1049            RuntimeStoreWriteFenceOutcome::Applied => {
1050                Ok(FencedMachineLifecycleCasOutcome::Applied { record, version })
1051            }
1052            RuntimeStoreWriteFenceOutcome::Conflict { reason } => {
1053                Ok(FencedMachineLifecycleCasOutcome::FenceConflict { reason })
1054            }
1055            RuntimeStoreWriteFenceOutcome::Backoff { reason } => {
1056                Ok(FencedMachineLifecycleCasOutcome::FenceBackoff { reason })
1057            }
1058        }
1059    }
1060
1061    async fn load_machine_lifecycle_record(
1062        &self,
1063        runtime_id: &LogicalRuntimeId,
1064    ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
1065        let inner = self.inner.lock().await;
1066        Ok(inner.runtime_lifecycle.get(&runtime_id.0).cloned())
1067    }
1068
1069    async fn commit_machine_lifecycle(
1070        &self,
1071        runtime_id: &LogicalRuntimeId,
1072        commit: MachineLifecycleCommit,
1073        input_states: &[InputStatePersistenceRecord],
1074    ) -> Result<(), RuntimeStoreError> {
1075        let record = commit.store_record().encode()?;
1076        let mut inner = self.inner.lock().await;
1077        let rid = runtime_id.0.clone();
1078
1079        // Single lock acquisition — atomic for in-memory
1080        inner.runtime_lifecycle.insert(rid.clone(), record);
1081        let states = inner.input_states.entry(rid).or_default();
1082        for record in input_states {
1083            let bundle = record.as_stored();
1084            states.insert(bundle.state.input_id.clone(), bundle.clone());
1085        }
1086
1087        Ok(())
1088    }
1089
1090    async fn commit_unregister_finalization(
1091        &self,
1092        runtime_id: &LogicalRuntimeId,
1093        finalization: crate::store::UnregisterFinalizationCommit,
1094    ) -> Result<(), RuntimeStoreError> {
1095        let (snapshot, input_states, retired_ops_epoch) = finalization.into_parts();
1096        let lifecycle_record = MachineLifecycleStoreRecord::from_snapshot(&snapshot).encode()?;
1097        let mut inner = self.inner.lock().await;
1098        let rid = runtime_id.0.clone();
1099
1100        // One lock acquisition is the in-memory transaction boundary. The
1101        // finalization token prepared every owned value before this method, so
1102        // no fallible request preparation remains after the first mutation.
1103        inner
1104            .runtime_lifecycle
1105            .insert(rid.clone(), lifecycle_record);
1106        let states = inner.input_states.entry(rid.clone()).or_default();
1107        for record in input_states {
1108            let bundle = record.clone_stored();
1109            states.insert(bundle.state.input_id.clone(), bundle);
1110        }
1111        if inner
1112            .ops_lifecycle_snapshots
1113            .get(&rid)
1114            .is_some_and(|snapshot| snapshot.epoch_id == retired_ops_epoch)
1115        {
1116            inner.ops_lifecycle_snapshots.remove(&rid);
1117        }
1118        inner.retired_ops_epochs.insert((rid, retired_ops_epoch));
1119        Ok(())
1120    }
1121
1122    async fn persist_ops_lifecycle(
1123        &self,
1124        runtime_id: &LogicalRuntimeId,
1125        snapshot: &PersistedOpsSnapshot,
1126    ) -> Result<(), RuntimeStoreError> {
1127        let mut inner = self.inner.lock().await;
1128        if inner
1129            .retired_ops_epochs
1130            .contains(&(runtime_id.0.clone(), snapshot.epoch_id.clone()))
1131        {
1132            return Err(RuntimeStoreError::OpsLifecycleEpochRetired {
1133                runtime_id: runtime_id.0.clone(),
1134                epoch_id: snapshot.epoch_id.clone(),
1135            });
1136        }
1137        inner
1138            .ops_lifecycle_snapshots
1139            .insert(runtime_id.0.clone(), snapshot.clone());
1140        Ok(())
1141    }
1142
1143    async fn initialize_ops_lifecycle_if_absent(
1144        &self,
1145        runtime_id: &LogicalRuntimeId,
1146        candidate: &PersistedOpsSnapshot,
1147    ) -> Result<PersistedOpsSnapshot, RuntimeStoreError> {
1148        let mut inner = self.inner.lock().await;
1149        let key = runtime_id.0.clone();
1150        if inner
1151            .retired_ops_epochs
1152            .contains(&(key.clone(), candidate.epoch_id.clone()))
1153        {
1154            return Err(RuntimeStoreError::OpsLifecycleEpochRetired {
1155                runtime_id: key,
1156                epoch_id: candidate.epoch_id.clone(),
1157            });
1158        }
1159        let canonical = inner
1160            .ops_lifecycle_snapshots
1161            .entry(key)
1162            .or_insert_with(|| candidate.clone())
1163            .clone();
1164        if inner
1165            .retired_ops_epochs
1166            .contains(&(runtime_id.0.clone(), canonical.epoch_id.clone()))
1167        {
1168            return Err(RuntimeStoreError::OpsLifecycleEpochRetired {
1169                runtime_id: runtime_id.0.clone(),
1170                epoch_id: canonical.epoch_id,
1171            });
1172        }
1173        Ok(canonical)
1174    }
1175
1176    async fn load_ops_lifecycle(
1177        &self,
1178        runtime_id: &LogicalRuntimeId,
1179    ) -> Result<Option<PersistedOpsSnapshot>, RuntimeStoreError> {
1180        let inner = self.inner.lock().await;
1181        Ok(inner.ops_lifecycle_snapshots.get(&runtime_id.0).cloned())
1182    }
1183
1184    async fn delete_ops_lifecycle(
1185        &self,
1186        runtime_id: &LogicalRuntimeId,
1187    ) -> Result<(), RuntimeStoreError> {
1188        let mut inner = self.inner.lock().await;
1189        inner.ops_lifecycle_snapshots.remove(&runtime_id.0);
1190        Ok(())
1191    }
1192}
1193
1194#[cfg(test)]
1195#[allow(clippy::unwrap_used)]
1196mod tests {
1197    use super::*;
1198    use crate::RuntimeState;
1199    use crate::store::MachineLifecycleBindingFacts;
1200    use meerkat_core::lifecycle::run_primitive::RunApplyBoundary;
1201
1202    fn make_receipt(run_id: RunId, seq: u64) -> RunBoundaryReceipt {
1203        RunBoundaryReceipt {
1204            run_id,
1205            boundary: RunApplyBoundary::RunStart,
1206            contributing_input_ids: vec![],
1207            conversation_digest: None,
1208            message_count: 0,
1209            sequence: seq,
1210        }
1211    }
1212
1213    fn lifecycle_commit(
1214        runtime_id: &LogicalRuntimeId,
1215        state: RuntimeState,
1216        fence_token: u64,
1217        runtime_generation: u64,
1218    ) -> MachineLifecycleCommit {
1219        MachineLifecycleCommit::new_with_binding(
1220            state,
1221            MachineLifecycleBindingFacts::new(
1222                Some(runtime_id.0.clone()),
1223                Some(fence_token),
1224                Some(runtime_generation),
1225                Some(format!("epoch-{runtime_generation}")),
1226            ),
1227            crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1228        )
1229    }
1230
1231    fn persistable(bundle: StoredInputState) -> InputStatePersistenceRecord {
1232        InputStatePersistenceRecord::from_machine_snapshot(bundle).unwrap()
1233    }
1234
1235    fn session_with_user(content: &str) -> meerkat_core::Session {
1236        let mut session = meerkat_core::Session::new();
1237        session.push(meerkat_core::types::Message::User(
1238            meerkat_core::types::UserMessage::text(content.to_string()),
1239        ));
1240        session
1241    }
1242
1243    fn session_with_compaction_intent() -> (
1244        meerkat_core::Session,
1245        meerkat_core::CompactionProjectionIntent,
1246    ) {
1247        let mut session = session_with_user("verbose context one");
1248        session.push(meerkat_core::types::Message::User(
1249            meerkat_core::types::UserMessage::text("verbose context two"),
1250        ));
1251        let parent = session.transcript_revision().unwrap();
1252        session
1253            .commit_transcript_rewrite(
1254                meerkat_core::TranscriptRewriteSelection::MessageRange { start: 0, end: 2 },
1255                vec![meerkat_core::types::Message::User(
1256                    meerkat_core::types::UserMessage::compaction_summary("compacted context"),
1257                )],
1258                meerkat_core::TranscriptRewriteReason::new("compaction"),
1259                Some("runtime-store-test".to_string()),
1260                Some(parent),
1261            )
1262            .unwrap();
1263        let mut encoded = serde_json::to_value(&session).unwrap();
1264        encoded["metadata"][meerkat_core::SESSION_TRANSCRIPT_HISTORY_STATE_KEY]["commits"][0]["selection"] = serde_json::json!({
1265            "type": "compaction_message_range",
1266            "range": { "start": 0, "end": 2 }
1267        });
1268        let mut session: meerkat_core::Session = serde_json::from_value(encoded).unwrap();
1269        let commit = session
1270            .transcript_history_state()
1271            .unwrap()
1272            .unwrap()
1273            .commits
1274            .last()
1275            .unwrap()
1276            .clone();
1277        let intent = meerkat_core::CompactionProjectionIntent {
1278            projection: serde_json::from_value(serde_json::json!({
1279                "session_id": session.id(),
1280                "parent_revision": &commit.parent_revision,
1281                "revision": &commit.revision,
1282                "commit_fingerprint": "sha256:827d8ee5666e51b2ced4d303640740680d96151d92187fd6e981c29550072c62",
1283            }))
1284            .unwrap(),
1285            summary_tokens: 5,
1286            messages_before: 2,
1287            messages_after: 1,
1288        };
1289        session
1290            .add_compaction_projection_intent(intent.clone())
1291            .unwrap();
1292        (session, intent)
1293    }
1294
1295    fn snapshot_with_raw_intents(
1296        session: &meerkat_core::Session,
1297        intents: &[meerkat_core::CompactionProjectionIntent],
1298    ) -> Vec<u8> {
1299        let mut value = serde_json::to_value(session).unwrap();
1300        value["metadata"][meerkat_core::memory::SESSION_COMPACTION_PROJECTION_INTENTS_KEY] =
1301            serde_json::to_value(intents).unwrap();
1302        serde_json::to_vec(&value).unwrap()
1303    }
1304
1305    fn unbacked_intent(
1306        session_id: &meerkat_core::types::SessionId,
1307    ) -> meerkat_core::CompactionProjectionIntent {
1308        meerkat_core::CompactionProjectionIntent {
1309            projection: serde_json::from_value(serde_json::json!({
1310                "session_id": session_id,
1311                "parent_revision": "missing-parent",
1312                "revision": "missing-revision",
1313                "commit_fingerprint": "sha256:unbacked-persisted-fixture",
1314            }))
1315            .unwrap(),
1316            summary_tokens: 1,
1317            messages_before: 2,
1318            messages_after: 1,
1319        }
1320    }
1321
1322    #[tokio::test]
1323    async fn atomic_apply_commits_rewrite_and_compaction_outbox_as_one_boundary() {
1324        let store = InMemoryRuntimeStore::new();
1325        let rid = LogicalRuntimeId::new("runtime-compaction-outbox");
1326        let (session, intent) = session_with_compaction_intent();
1327        let snapshot = serde_json::to_vec(&session).unwrap();
1328        store
1329            .atomic_apply(
1330                &rid,
1331                Some(SessionDelta {
1332                    session_snapshot: snapshot.clone(),
1333                }),
1334                make_receipt(RunId::new(), 1),
1335                vec![],
1336                Some(session.id().clone()),
1337            )
1338            .await
1339            .unwrap();
1340        assert_eq!(
1341            store.load_session_snapshot(&rid).await.unwrap(),
1342            Some(snapshot)
1343        );
1344        assert_eq!(
1345            store
1346                .load_pending_compaction_projections(&rid)
1347                .await
1348                .unwrap(),
1349            vec![intent.clone()]
1350        );
1351        store
1352            .mark_compaction_projection_finalized(&rid, &intent.projection)
1353            .await
1354            .unwrap();
1355        store
1356            .mark_compaction_projection_finalized(&rid, &intent.projection)
1357            .await
1358            .unwrap();
1359        assert!(
1360            store
1361                .load_pending_compaction_projections(&rid)
1362                .await
1363                .unwrap()
1364                .is_empty()
1365        );
1366        let persisted: meerkat_core::Session =
1367            serde_json::from_slice(&store.load_session_snapshot(&rid).await.unwrap().unwrap())
1368                .unwrap();
1369        assert!(
1370            persisted
1371                .compaction_projection_intents()
1372                .unwrap()
1373                .is_empty()
1374        );
1375    }
1376
1377    #[tokio::test]
1378    async fn finalized_outbox_tombstone_rejects_atomic_and_non_boundary_snapshot_replay() {
1379        let store = InMemoryRuntimeStore::new();
1380        let rid = LogicalRuntimeId::new("runtime-finalized-compaction-replay");
1381        let (session, intent) = session_with_compaction_intent();
1382        let replay_snapshot = serde_json::to_vec(&session).unwrap();
1383        let commit = session
1384            .transcript_history_state()
1385            .unwrap()
1386            .unwrap()
1387            .commits
1388            .last()
1389            .unwrap()
1390            .clone();
1391        store
1392            .atomic_apply(
1393                &rid,
1394                Some(SessionDelta {
1395                    session_snapshot: replay_snapshot.clone(),
1396                }),
1397                make_receipt(RunId::new(), 1),
1398                vec![],
1399                Some(session.id().clone()),
1400            )
1401            .await
1402            .unwrap();
1403        store
1404            .mark_compaction_projection_finalized(&rid, &intent.projection)
1405            .await
1406            .unwrap();
1407        let cleaned_snapshot = store.load_session_snapshot(&rid).await.unwrap().unwrap();
1408
1409        let replay_run_id = RunId::new();
1410        let error = store
1411            .atomic_apply(
1412                &rid,
1413                Some(SessionDelta {
1414                    session_snapshot: replay_snapshot.clone(),
1415                }),
1416                make_receipt(replay_run_id.clone(), 2),
1417                vec![],
1418                Some(session.id().clone()),
1419            )
1420            .await
1421            .unwrap_err();
1422        assert!(error.to_string().contains("finalized compaction intent"));
1423        assert!(
1424            store
1425                .load_boundary_receipt(&rid, &replay_run_id, 2)
1426                .await
1427                .unwrap()
1428                .is_none(),
1429            "finalized replay rejection must roll back the whole atomic boundary"
1430        );
1431
1432        let error = store
1433            .commit_session_snapshot(
1434                &rid,
1435                SessionDelta {
1436                    session_snapshot: replay_snapshot.clone(),
1437                },
1438            )
1439            .await
1440            .unwrap_err();
1441        assert!(error.to_string().contains("finalized compaction intent"));
1442        let error = store
1443            .commit_session_transcript_rewrite_snapshot(
1444                &rid,
1445                SessionDelta {
1446                    session_snapshot: replay_snapshot.clone(),
1447                },
1448                &commit,
1449            )
1450            .await
1451            .unwrap_err();
1452        assert!(error.to_string().contains("finalized compaction intent"));
1453        let error = store
1454            .replace_session_snapshot_if_current(&rid, &cleaned_snapshot, replay_snapshot)
1455            .await
1456            .unwrap_err();
1457        assert!(error.to_string().contains("finalized compaction intent"));
1458
1459        assert_eq!(
1460            store.load_session_snapshot(&rid).await.unwrap(),
1461            Some(cleaned_snapshot)
1462        );
1463        assert!(
1464            store
1465                .load_pending_compaction_projections(&rid)
1466                .await
1467                .unwrap()
1468                .is_empty(),
1469            "a finalized tombstone must never be silently revived or left untracked"
1470        );
1471    }
1472
1473    #[tokio::test]
1474    async fn invalid_compaction_intent_leaves_snapshot_and_outbox_unmodified() {
1475        let store = InMemoryRuntimeStore::new();
1476        let rid = LogicalRuntimeId::new("runtime-invalid-compaction-outbox");
1477        let (session, mut intent) = session_with_compaction_intent();
1478        intent.summary_tokens += 1;
1479        let conflicting = vec![
1480            session.compaction_projection_intents().unwrap()[0].clone(),
1481            intent,
1482        ];
1483        let error = store
1484            .atomic_apply(
1485                &rid,
1486                Some(SessionDelta {
1487                    session_snapshot: snapshot_with_raw_intents(&session, &conflicting),
1488                }),
1489                make_receipt(RunId::new(), 2),
1490                vec![],
1491                Some(session.id().clone()),
1492            )
1493            .await
1494            .unwrap_err();
1495        assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
1496        assert_eq!(store.load_session_snapshot(&rid).await.unwrap(), None);
1497        assert!(
1498            store
1499                .load_pending_compaction_projections(&rid)
1500                .await
1501                .unwrap()
1502                .is_empty()
1503        );
1504
1505        let foreign = session_with_compaction_intent().1;
1506        for (sequence, invalid) in [foreign, unbacked_intent(session.id())]
1507            .into_iter()
1508            .enumerate()
1509        {
1510            let error = store
1511                .atomic_apply(
1512                    &rid,
1513                    Some(SessionDelta {
1514                        session_snapshot: snapshot_with_raw_intents(&session, &[invalid]),
1515                    }),
1516                    make_receipt(RunId::new(), 10 + sequence as u64),
1517                    vec![],
1518                    Some(session.id().clone()),
1519                )
1520                .await
1521                .unwrap_err();
1522            assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
1523            assert_eq!(store.load_session_snapshot(&rid).await.unwrap(), None);
1524            assert!(
1525                store
1526                    .load_pending_compaction_projections(&rid)
1527                    .await
1528                    .unwrap()
1529                    .is_empty()
1530            );
1531        }
1532    }
1533
1534    #[tokio::test]
1535    async fn superseded_snapshot_rejects_without_advancing_compaction_outbox() {
1536        let store = InMemoryRuntimeStore::new();
1537        let rid = LogicalRuntimeId::new("runtime-superseded-compaction-outbox");
1538        let (incoming, intent) = session_with_compaction_intent();
1539        let mut current = incoming.clone();
1540        current
1541            .complete_compaction_projection_intent(&intent.projection)
1542            .unwrap();
1543        current.push(meerkat_core::types::Message::User(
1544            meerkat_core::types::UserMessage::text("already advanced"),
1545        ));
1546        let current_snapshot = serde_json::to_vec(&current).unwrap();
1547        store
1548            .commit_session_snapshot(
1549                &rid,
1550                SessionDelta {
1551                    session_snapshot: current_snapshot.clone(),
1552                },
1553            )
1554            .await
1555            .unwrap();
1556        let error = store
1557            .atomic_apply(
1558                &rid,
1559                Some(SessionDelta {
1560                    session_snapshot: serde_json::to_vec(&incoming).unwrap(),
1561                }),
1562                make_receipt(RunId::new(), 3),
1563                vec![],
1564                Some(incoming.id().clone()),
1565            )
1566            .await
1567            .expect_err("superseded compaction boundary must be explicitly rejected");
1568        assert!(matches!(
1569            error,
1570            RuntimeStoreError::SessionSnapshotSuperseded { .. }
1571        ));
1572        assert_eq!(
1573            store.load_session_snapshot(&rid).await.unwrap(),
1574            Some(current_snapshot)
1575        );
1576        assert!(
1577            store
1578                .load_pending_compaction_projections(&rid)
1579                .await
1580                .unwrap()
1581                .is_empty()
1582        );
1583    }
1584
1585    #[tokio::test]
1586    async fn existing_outbox_rejects_changed_intent_without_advancing_snapshot() {
1587        let store = InMemoryRuntimeStore::new();
1588        let rid = LogicalRuntimeId::new("runtime-conflicting-compaction-outbox");
1589        let (session, intent) = session_with_compaction_intent();
1590        let original_snapshot = serde_json::to_vec(&session).unwrap();
1591        store
1592            .atomic_apply(
1593                &rid,
1594                Some(SessionDelta {
1595                    session_snapshot: original_snapshot.clone(),
1596                }),
1597                make_receipt(RunId::new(), 60),
1598                vec![],
1599                Some(session.id().clone()),
1600            )
1601            .await
1602            .unwrap();
1603
1604        let mut advanced = session.clone();
1605        advanced.push(meerkat_core::types::Message::User(
1606            meerkat_core::types::UserMessage::text("later turn"),
1607        ));
1608        let mut conflicting = intent.clone();
1609        conflicting.summary_tokens += 1;
1610        let error = store
1611            .atomic_apply(
1612                &rid,
1613                Some(SessionDelta {
1614                    session_snapshot: snapshot_with_raw_intents(&advanced, &[conflicting]),
1615                }),
1616                make_receipt(RunId::new(), 61),
1617                vec![],
1618                Some(session.id().clone()),
1619            )
1620            .await
1621            .unwrap_err();
1622        assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
1623        assert_eq!(
1624            store.load_session_snapshot(&rid).await.unwrap(),
1625            Some(original_snapshot)
1626        );
1627        assert_eq!(
1628            store
1629                .load_pending_compaction_projections(&rid)
1630                .await
1631                .unwrap(),
1632            vec![intent]
1633        );
1634    }
1635
1636    #[tokio::test]
1637    async fn non_boundary_snapshot_apis_cannot_bypass_compaction_outbox() {
1638        let store = InMemoryRuntimeStore::new();
1639        let rid = LogicalRuntimeId::new("runtime-compaction-bypass");
1640        let (session, _intent) = session_with_compaction_intent();
1641        let snapshot = serde_json::to_vec(&session).unwrap();
1642        let commit = session
1643            .transcript_history_state()
1644            .unwrap()
1645            .unwrap()
1646            .commits
1647            .last()
1648            .unwrap()
1649            .clone();
1650        assert!(
1651            store
1652                .commit_session_snapshot(
1653                    &rid,
1654                    SessionDelta {
1655                        session_snapshot: snapshot.clone(),
1656                    },
1657                )
1658                .await
1659                .is_err()
1660        );
1661        assert!(
1662            store
1663                .commit_session_transcript_rewrite_snapshot(
1664                    &rid,
1665                    SessionDelta {
1666                        session_snapshot: snapshot.clone(),
1667                    },
1668                    &commit,
1669                )
1670                .await
1671                .is_err()
1672        );
1673        assert_eq!(store.load_session_snapshot(&rid).await.unwrap(), None);
1674        let clean = meerkat_core::Session::with_id(session.id().clone());
1675        let clean_snapshot = serde_json::to_vec(&clean).unwrap();
1676        store
1677            .commit_session_snapshot(
1678                &rid,
1679                SessionDelta {
1680                    session_snapshot: clean_snapshot.clone(),
1681                },
1682            )
1683            .await
1684            .unwrap();
1685        assert!(
1686            store
1687                .replace_session_snapshot_if_current(&rid, &clean_snapshot, snapshot)
1688                .await
1689                .is_err()
1690        );
1691        assert_eq!(
1692            store.load_session_snapshot(&rid).await.unwrap(),
1693            Some(clean_snapshot)
1694        );
1695        assert!(
1696            store
1697                .load_pending_compaction_projections(&rid)
1698                .await
1699                .unwrap()
1700                .is_empty()
1701        );
1702    }
1703
1704    #[tokio::test]
1705    async fn atomic_apply_roundtrip() {
1706        let store = InMemoryRuntimeStore::new();
1707        let rid = LogicalRuntimeId::new("test-runtime");
1708        let run_id = RunId::new();
1709        let input_id = InputId::new();
1710
1711        let bundle = StoredInputState::new_accepted(input_id.clone());
1712        let receipt = make_receipt(run_id.clone(), 0);
1713
1714        let session = session_with_user("hello");
1715        let session_snapshot = serde_json::to_vec(&session).unwrap();
1716
1717        store
1718            .atomic_apply(
1719                &rid,
1720                Some(SessionDelta { session_snapshot }),
1721                receipt.clone(),
1722                vec![persistable(bundle)],
1723                None,
1724            )
1725            .await
1726            .unwrap();
1727
1728        // Load input states
1729        let states = store.load_input_states(&rid).await.unwrap();
1730        assert_eq!(states.len(), 1);
1731        assert_eq!(states[0].state.input_id, input_id);
1732
1733        // Load receipt
1734        let loaded = store.load_boundary_receipt(&rid, &run_id, 0).await.unwrap();
1735        assert!(loaded.is_some());
1736    }
1737
1738    #[tokio::test]
1739    async fn machine_terminal_atomic_apply_rolls_back_all_maps_on_receipt_conflict() {
1740        let store = InMemoryRuntimeStore::new();
1741        let rid = LogicalRuntimeId::new("terminal-receipt-conflict");
1742        let receipt = make_receipt(RunId::new(), 0);
1743        let seeded_input = StoredInputState::new_accepted(InputId::new());
1744        store
1745            .atomic_apply(
1746                &rid,
1747                None,
1748                receipt.clone(),
1749                vec![persistable(seeded_input.clone())],
1750                None,
1751            )
1752            .await
1753            .unwrap();
1754
1755        let session = session_with_user("must roll back");
1756        let replacement_input = StoredInputState::new_accepted(InputId::new());
1757        let error = store
1758            .atomic_apply_with_machine_lifecycle(
1759                &rid,
1760                SessionDelta {
1761                    session_snapshot: serde_json::to_vec(&session).unwrap(),
1762                },
1763                receipt,
1764                MachineLifecycleCommit::new_with_binding(
1765                    crate::RuntimeState::Idle,
1766                    crate::store::MachineLifecycleBindingFacts::default(),
1767                    crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1768                ),
1769                vec![persistable(replacement_input)],
1770                session.id().clone(),
1771            )
1772            .await
1773            .expect_err("duplicate receipt must reject the entire terminal transaction");
1774        assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
1775        assert!(store.load_session_snapshot(&rid).await.unwrap().is_none());
1776        assert_eq!(
1777            crate::store::load_runtime_state(&store, &rid)
1778                .await
1779                .unwrap(),
1780            None
1781        );
1782        let inputs = store.load_input_states(&rid).await.unwrap();
1783        assert_eq!(inputs.len(), 1);
1784        assert_eq!(inputs[0].state.input_id, seeded_input.state.input_id);
1785    }
1786
1787    #[tokio::test]
1788    async fn machine_terminal_atomic_apply_tracks_and_tombstones_compaction_intents() {
1789        let store = InMemoryRuntimeStore::new();
1790        let rid = LogicalRuntimeId::new("terminal-compaction-outbox");
1791        let (session, intent) = session_with_compaction_intent();
1792        let encoded = serde_json::to_vec(&session).unwrap();
1793
1794        store
1795            .atomic_apply_with_machine_lifecycle(
1796                &rid,
1797                SessionDelta {
1798                    session_snapshot: encoded.clone(),
1799                },
1800                make_receipt(RunId::new(), 0),
1801                MachineLifecycleCommit::new_with_binding(
1802                    crate::RuntimeState::Idle,
1803                    crate::store::MachineLifecycleBindingFacts::default(),
1804                    crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1805                ),
1806                Vec::new(),
1807                session.id().clone(),
1808            )
1809            .await
1810            .unwrap();
1811        assert_eq!(
1812            store
1813                .load_pending_compaction_projections(&rid)
1814                .await
1815                .unwrap(),
1816            vec![intent.clone()]
1817        );
1818
1819        store
1820            .mark_compaction_projection_finalized(&rid, &intent.projection)
1821            .await
1822            .unwrap();
1823        let error = store
1824            .atomic_apply_with_machine_lifecycle(
1825                &rid,
1826                SessionDelta {
1827                    session_snapshot: encoded,
1828                },
1829                make_receipt(RunId::new(), 1),
1830                MachineLifecycleCommit::new_with_binding(
1831                    crate::RuntimeState::Idle,
1832                    crate::store::MachineLifecycleBindingFacts::default(),
1833                    crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1834                ),
1835                Vec::new(),
1836                session.id().clone(),
1837            )
1838            .await
1839            .expect_err("a finalized compaction tombstone must reject stale terminal replay");
1840        assert!(
1841            error
1842                .to_string()
1843                .contains("replays finalized compaction intent")
1844        );
1845    }
1846
1847    #[tokio::test]
1848    async fn machine_terminal_atomic_apply_rejects_corrupt_previous_snapshot_before_mutation() {
1849        let store = InMemoryRuntimeStore::new();
1850        let rid = LogicalRuntimeId::new("terminal-corrupt-head");
1851        let corrupt = b"{not-a-session".to_vec();
1852        store
1853            .inner
1854            .lock()
1855            .await
1856            .sessions
1857            .insert(rid.0.clone(), corrupt.clone());
1858        let session = session_with_user("incoming terminal transcript");
1859        let receipt = make_receipt(RunId::new(), 0);
1860        let error = store
1861            .atomic_apply_with_machine_lifecycle(
1862                &rid,
1863                SessionDelta {
1864                    session_snapshot: serde_json::to_vec(&session).unwrap(),
1865                },
1866                receipt.clone(),
1867                MachineLifecycleCommit::new_with_binding(
1868                    crate::RuntimeState::Idle,
1869                    crate::store::MachineLifecycleBindingFacts::default(),
1870                    crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1871                ),
1872                vec![persistable(StoredInputState::new_accepted(InputId::new()))],
1873                session.id().clone(),
1874            )
1875            .await
1876            .expect_err("corrupt durable head must fail before every terminal mutation");
1877        assert!(matches!(error, RuntimeStoreError::ReadFailed(_)));
1878        assert_eq!(
1879            store.load_session_snapshot(&rid).await.unwrap(),
1880            Some(corrupt)
1881        );
1882        assert_eq!(
1883            crate::store::load_runtime_state(&store, &rid)
1884                .await
1885                .unwrap(),
1886            None
1887        );
1888        assert!(store.load_input_states(&rid).await.unwrap().is_empty());
1889        assert!(
1890            store
1891                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1892                .await
1893                .unwrap()
1894                .is_none()
1895        );
1896    }
1897
1898    #[tokio::test]
1899    async fn machine_terminal_atomic_apply_rejects_superseded_snapshot_without_publication() {
1900        let store = InMemoryRuntimeStore::new();
1901        let rid = LogicalRuntimeId::new("terminal-superseded-head");
1902        let incoming = session_with_user("failed turn input");
1903        let mut durable_head = incoming.clone();
1904        durable_head.push(meerkat_core::types::Message::User(
1905            meerkat_core::types::UserMessage::text("already advanced"),
1906        ));
1907        let durable_snapshot = serde_json::to_vec(&durable_head).unwrap();
1908        store
1909            .commit_session_snapshot(
1910                &rid,
1911                SessionDelta {
1912                    session_snapshot: durable_snapshot.clone(),
1913                },
1914            )
1915            .await
1916            .unwrap();
1917
1918        let receipt = make_receipt(RunId::new(), 0);
1919        let error = store
1920            .atomic_apply_with_machine_lifecycle(
1921                &rid,
1922                SessionDelta {
1923                    session_snapshot: serde_json::to_vec(&incoming).unwrap(),
1924                },
1925                receipt.clone(),
1926                MachineLifecycleCommit::new_with_binding(
1927                    crate::RuntimeState::Idle,
1928                    MachineLifecycleBindingFacts::default(),
1929                    crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1930                ),
1931                vec![persistable(StoredInputState::new_accepted(InputId::new()))],
1932                incoming.id().clone(),
1933            )
1934            .await
1935            .expect_err("superseded terminal snapshot must reject the entire transaction");
1936        assert!(matches!(
1937            error,
1938            RuntimeStoreError::SessionSnapshotSuperseded { .. }
1939        ));
1940        assert_eq!(
1941            store.load_session_snapshot(&rid).await.unwrap(),
1942            Some(durable_snapshot)
1943        );
1944        assert_eq!(
1945            crate::store::load_runtime_state(&store, &rid)
1946                .await
1947                .unwrap(),
1948            None
1949        );
1950        assert!(store.load_input_states(&rid).await.unwrap().is_empty());
1951        assert!(
1952            store
1953                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1954                .await
1955                .unwrap()
1956                .is_none()
1957        );
1958    }
1959
1960    #[tokio::test]
1961    async fn legacy_atomic_apply_rejects_corrupt_previous_snapshot_before_mutation() {
1962        let store = InMemoryRuntimeStore::new();
1963        let rid = LogicalRuntimeId::new("legacy-corrupt-head");
1964        let corrupt = b"{not-a-session".to_vec();
1965        store
1966            .inner
1967            .lock()
1968            .await
1969            .sessions
1970            .insert(rid.0.clone(), corrupt.clone());
1971        let session = session_with_user("incoming transcript");
1972        let receipt = make_receipt(RunId::new(), 0);
1973        let error = store
1974            .atomic_apply(
1975                &rid,
1976                Some(SessionDelta {
1977                    session_snapshot: serde_json::to_vec(&session).unwrap(),
1978                }),
1979                receipt.clone(),
1980                vec![persistable(StoredInputState::new_accepted(InputId::new()))],
1981                Some(session.id().clone()),
1982            )
1983            .await
1984            .expect_err("corrupt durable head must fail before every boundary mutation");
1985        assert!(matches!(error, RuntimeStoreError::ReadFailed(_)));
1986        assert_eq!(
1987            store.load_session_snapshot(&rid).await.unwrap(),
1988            Some(corrupt)
1989        );
1990        assert!(store.load_input_states(&rid).await.unwrap().is_empty());
1991        assert!(
1992            store
1993                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1994                .await
1995                .unwrap()
1996                .is_none()
1997        );
1998    }
1999
2000    #[tokio::test]
2001    async fn atomic_apply_rejects_non_session_snapshot_without_owner_context() {
2002        let store = InMemoryRuntimeStore::new();
2003        let rid = LogicalRuntimeId::new("test-runtime");
2004        let run_id = RunId::new();
2005        let input_id = InputId::new();
2006
2007        let bundle = StoredInputState::new_accepted(input_id);
2008        let receipt = make_receipt(run_id, 0);
2009
2010        // Owner-context absence is not a license to store arbitrary bytes as a
2011        // session snapshot: a non-deserializable snapshot must fail closed.
2012        let err = store
2013            .atomic_apply(
2014                &rid,
2015                Some(SessionDelta {
2016                    session_snapshot: b"session-data".to_vec(),
2017                }),
2018                receipt,
2019                vec![persistable(bundle)],
2020                None,
2021            )
2022            .await
2023            .expect_err("non-Session snapshot must be rejected");
2024
2025        match err {
2026            RuntimeStoreError::WriteFailed(message) => {
2027                assert!(
2028                    message.contains("not a Session"),
2029                    "unexpected WriteFailed message: {message}"
2030                );
2031            }
2032            other => panic!("expected WriteFailed, got {other:?}"),
2033        }
2034    }
2035
2036    #[tokio::test]
2037    async fn persist_and_load_single_state() {
2038        let store = InMemoryRuntimeStore::new();
2039        let rid = LogicalRuntimeId::new("test");
2040        let input_id = InputId::new();
2041        let bundle = StoredInputState::new_accepted(input_id.clone());
2042
2043        store
2044            .persist_input_state(&rid, &persistable(bundle))
2045            .await
2046            .unwrap();
2047
2048        let loaded = store.load_input_state(&rid, &input_id).await.unwrap();
2049        assert!(loaded.is_some());
2050        assert_eq!(loaded.unwrap().state.input_id, input_id);
2051    }
2052
2053    fn replacement_records(
2054        expected: &[StoredInputState],
2055        recovery_count: u32,
2056    ) -> Vec<InputStatePersistenceRecord> {
2057        expected
2058            .iter()
2059            .cloned()
2060            .map(|mut row| {
2061                row.state.recovery_count = recovery_count;
2062                persistable(row)
2063            })
2064            .collect()
2065    }
2066
2067    #[tokio::test]
2068    async fn input_state_batch_cas_memory_swaps_once_and_stale_is_noop() {
2069        let store = InMemoryRuntimeStore::new();
2070        let rid = LogicalRuntimeId::new("input-cas-memory");
2071        let expected: Vec<_> = (0..3)
2072            .map(|_| StoredInputState::new_accepted(InputId::new()))
2073            .collect();
2074        let initial: Vec<_> = expected.iter().cloned().map(persistable).collect();
2075        store
2076            .persist_input_states_atomically(&rid, &initial)
2077            .await
2078            .unwrap();
2079
2080        let winner = replacement_records(&expected, 1);
2081        let stale_candidate = replacement_records(&expected, 2);
2082        assert_eq!(
2083            store
2084                .compare_and_swap_input_states_atomically(&rid, &expected, &winner)
2085                .await
2086                .unwrap(),
2087            InputStateBatchCasOutcome::Swapped
2088        );
2089        assert_eq!(
2090            store
2091                .compare_and_swap_input_states_atomically(&rid, &expected, &winner)
2092                .await
2093                .unwrap(),
2094            InputStateBatchCasOutcome::Swapped,
2095            "retry after a lost CAS acknowledgement must observe the exact replacement as success"
2096        );
2097        assert_eq!(
2098            store
2099                .compare_and_swap_input_states_atomically(&rid, &expected, &stale_candidate)
2100                .await
2101                .unwrap(),
2102            InputStateBatchCasOutcome::Stale
2103        );
2104        let rows = store.load_input_states(&rid).await.unwrap();
2105        assert_eq!(rows.len(), 3);
2106        assert!(rows.iter().all(|row| row.state.recovery_count == 1));
2107    }
2108
2109    #[tokio::test]
2110    async fn input_state_batch_cas_memory_rejects_missing_extra_and_key_mismatch() {
2111        let store = InMemoryRuntimeStore::new();
2112        let rid = LogicalRuntimeId::new("input-cas-shape");
2113        let expected: Vec<_> = (0..2)
2114            .map(|_| StoredInputState::new_accepted(InputId::new()))
2115            .collect();
2116        store
2117            .persist_input_state(&rid, &persistable(expected[0].clone()))
2118            .await
2119            .unwrap();
2120        let replacements = replacement_records(&expected, 1);
2121
2122        assert_eq!(
2123            store
2124                .compare_and_swap_input_states_atomically(&rid, &expected, &replacements)
2125                .await
2126                .unwrap(),
2127            InputStateBatchCasOutcome::Stale,
2128            "one missing durable row must stale the entire batch"
2129        );
2130        assert_eq!(
2131            store
2132                .load_input_state(&rid, &expected[0].state.input_id)
2133                .await
2134                .unwrap()
2135                .unwrap()
2136                .state
2137                .recovery_count,
2138            0,
2139            "stale comparison must not update the matching prefix row"
2140        );
2141
2142        let extra = vec![
2143            replacements[0].clone(),
2144            replacements[1].clone(),
2145            persistable(StoredInputState::new_accepted(InputId::new())),
2146        ];
2147        assert!(matches!(
2148            store
2149                .compare_and_swap_input_states_atomically(&rid, &expected, &extra)
2150                .await,
2151            Err(RuntimeStoreError::InvalidInputStateBatchCas { .. })
2152        ));
2153
2154        let wrong_key = vec![
2155            replacements[0].clone(),
2156            persistable(StoredInputState::new_accepted(InputId::new())),
2157        ];
2158        assert!(matches!(
2159            store
2160                .compare_and_swap_input_states_atomically(&rid, &expected, &wrong_key)
2161                .await,
2162            Err(RuntimeStoreError::InvalidInputStateBatchCas { .. })
2163        ));
2164    }
2165
2166    #[tokio::test]
2167    async fn load_nonexistent_returns_none() {
2168        let store = InMemoryRuntimeStore::new();
2169        let rid = LogicalRuntimeId::new("test");
2170
2171        let states = store.load_input_states(&rid).await.unwrap();
2172        assert!(states.is_empty());
2173
2174        let state = store.load_input_state(&rid, &InputId::new()).await.unwrap();
2175        assert!(state.is_none());
2176
2177        let receipt = store
2178            .load_boundary_receipt(&rid, &RunId::new(), 0)
2179            .await
2180            .unwrap();
2181        assert!(receipt.is_none());
2182    }
2183
2184    #[tokio::test]
2185    async fn atomic_apply_updates_existing() {
2186        let store = InMemoryRuntimeStore::new();
2187        let rid = LogicalRuntimeId::new("test");
2188        let input_id = InputId::new();
2189
2190        // First write
2191        let bundle1 = StoredInputState::new_accepted(input_id.clone());
2192        store
2193            .atomic_apply(
2194                &rid,
2195                None,
2196                make_receipt(RunId::new(), 0),
2197                vec![persistable(bundle1)],
2198                None,
2199            )
2200            .await
2201            .unwrap();
2202
2203        // Second write with updated seed phase
2204        let mut bundle2 = StoredInputState::new_accepted(input_id.clone());
2205        bundle2.seed.phase = crate::input_state::InputLifecycleState::Queued;
2206        store
2207            .atomic_apply(
2208                &rid,
2209                None,
2210                make_receipt(RunId::new(), 1),
2211                vec![persistable(bundle2)],
2212                None,
2213            )
2214            .await
2215            .unwrap();
2216
2217        let states = store.load_input_states(&rid).await.unwrap();
2218        assert_eq!(states.len(), 1);
2219        assert_eq!(
2220            states[0].seed.phase,
2221            crate::input_state::InputLifecycleState::Queued
2222        );
2223    }
2224
2225    #[tokio::test]
2226    async fn atomic_apply_validates_session_store_key_without_aliasing_snapshot() {
2227        let store = InMemoryRuntimeStore::new();
2228        let rid = LogicalRuntimeId::new("runtime-key");
2229        let session = meerkat_core::Session::new();
2230        let session_id = session.id().clone();
2231        let snapshot = serde_json::to_vec(&session).unwrap();
2232
2233        store
2234            .atomic_apply(
2235                &rid,
2236                Some(SessionDelta {
2237                    session_snapshot: snapshot.clone(),
2238                }),
2239                make_receipt(RunId::new(), 0),
2240                vec![],
2241                Some(session_id.clone()),
2242            )
2243            .await
2244            .unwrap();
2245
2246        assert_eq!(
2247            store.load_session_snapshot(&rid).await.unwrap(),
2248            Some(snapshot)
2249        );
2250        assert!(
2251            store
2252                .load_session_snapshot(&LogicalRuntimeId::legacy_session_uuid_alias(&session_id))
2253                .await
2254                .unwrap()
2255                .is_none(),
2256            "session_store_key must validate the snapshot identity, not create a raw UUID runtime alias"
2257        );
2258    }
2259
2260    #[tokio::test]
2261    async fn atomic_apply_rejects_mismatched_session_store_key() {
2262        let store = InMemoryRuntimeStore::new();
2263        let rid = LogicalRuntimeId::new("runtime-key");
2264        let session = meerkat_core::Session::new();
2265        let wrong_session_id = meerkat_core::Session::new().id().clone();
2266        let snapshot = serde_json::to_vec(&session).unwrap();
2267
2268        let err = store
2269            .atomic_apply(
2270                &rid,
2271                Some(SessionDelta {
2272                    session_snapshot: snapshot,
2273                }),
2274                make_receipt(RunId::new(), 0),
2275                vec![],
2276                Some(wrong_session_id),
2277            )
2278            .await
2279            .expect_err("mismatched session_store_key should fail");
2280
2281        assert!(matches!(err, RuntimeStoreError::SessionKeyMismatch { .. }));
2282        assert!(store.load_session_snapshot(&rid).await.unwrap().is_none());
2283    }
2284
2285    #[tokio::test]
2286    async fn atomic_apply_persists_machine_owned_receipt() {
2287        let store = InMemoryRuntimeStore::new();
2288        let rid = LogicalRuntimeId::new("test");
2289        let run_id = RunId::new();
2290        let input_id = InputId::new();
2291        let session = meerkat_core::Session::new();
2292        let snapshot = serde_json::to_vec(&session).unwrap();
2293        let receipt = RunBoundaryReceipt {
2294            run_id: run_id.clone(),
2295            boundary: RunApplyBoundary::Immediate,
2296            contributing_input_ids: vec![input_id.clone()],
2297            conversation_digest: Some("machine-owned-digest".to_string()),
2298            message_count: 42,
2299            sequence: 7,
2300        };
2301
2302        store
2303            .atomic_apply(
2304                &rid,
2305                Some(SessionDelta {
2306                    session_snapshot: snapshot,
2307                }),
2308                receipt.clone(),
2309                vec![persistable(StoredInputState::new_accepted(input_id))],
2310                None,
2311            )
2312            .await
2313            .unwrap();
2314
2315        assert_eq!(receipt.run_id, run_id);
2316        assert!(receipt.conversation_digest.is_some());
2317        let loaded = store
2318            .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
2319            .await
2320            .unwrap();
2321        assert!(loaded.is_some(), "receipt should be persisted");
2322        let Some(loaded) = loaded else {
2323            unreachable!("asserted above");
2324        };
2325        assert_eq!(loaded, receipt);
2326    }
2327
2328    #[tokio::test]
2329    async fn multiple_runtimes_isolated() {
2330        let store = InMemoryRuntimeStore::new();
2331        let rid1 = LogicalRuntimeId::new("runtime-1");
2332        let rid2 = LogicalRuntimeId::new("runtime-2");
2333
2334        store
2335            .persist_input_state(
2336                &rid1,
2337                &persistable(StoredInputState::new_accepted(InputId::new())),
2338            )
2339            .await
2340            .unwrap();
2341        store
2342            .persist_input_state(
2343                &rid2,
2344                &persistable(StoredInputState::new_accepted(InputId::new())),
2345            )
2346            .await
2347            .unwrap();
2348        store
2349            .persist_input_state(
2350                &rid2,
2351                &persistable(StoredInputState::new_accepted(InputId::new())),
2352            )
2353            .await
2354            .unwrap();
2355
2356        let s1 = store.load_input_states(&rid1).await.unwrap();
2357        let s2 = store.load_input_states(&rid2).await.unwrap();
2358        assert_eq!(s1.len(), 1);
2359        assert_eq!(s2.len(), 2);
2360    }
2361
2362    #[tokio::test]
2363    async fn load_session_snapshot_roundtrip() {
2364        let store = InMemoryRuntimeStore::new();
2365        let rid = LogicalRuntimeId::new("runtime");
2366        let snapshot = serde_json::to_vec(&meerkat_core::Session::new()).unwrap();
2367
2368        store
2369            .atomic_apply(
2370                &rid,
2371                Some(SessionDelta {
2372                    session_snapshot: snapshot.clone(),
2373                }),
2374                make_receipt(RunId::new(), 0),
2375                vec![],
2376                None,
2377            )
2378            .await
2379            .unwrap();
2380
2381        let loaded = store.load_session_snapshot(&rid).await.unwrap();
2382        assert_eq!(loaded, Some(snapshot));
2383    }
2384
2385    #[tokio::test]
2386    async fn commit_session_snapshot_rejects_stale_runtime_parent() {
2387        let store = InMemoryRuntimeStore::new();
2388        let rid = LogicalRuntimeId::new("runtime-stale-parent");
2389        let accepted = session_with_user("accepted runtime turn");
2390        let mut stale = meerkat_core::Session::with_id(accepted.id().clone());
2391        stale.push(meerkat_core::types::Message::User(
2392            meerkat_core::types::UserMessage::text("stale runtime turn".to_string()),
2393        ));
2394        let accepted_snapshot = serde_json::to_vec(&accepted).unwrap();
2395
2396        store
2397            .commit_session_snapshot(
2398                &rid,
2399                SessionDelta {
2400                    session_snapshot: accepted_snapshot.clone(),
2401                },
2402            )
2403            .await
2404            .unwrap();
2405
2406        let err = store
2407            .commit_session_snapshot(
2408                &rid,
2409                SessionDelta {
2410                    session_snapshot: serde_json::to_vec(&stale).unwrap(),
2411                },
2412            )
2413            .await
2414            .expect_err("stale non-continuation must not overwrite runtime snapshot");
2415
2416        assert!(matches!(err, RuntimeStoreError::WriteFailed(_)));
2417        assert_eq!(
2418            store.load_session_snapshot(&rid).await.unwrap(),
2419            Some(accepted_snapshot)
2420        );
2421    }
2422
2423    #[tokio::test]
2424    async fn atomic_apply_keeps_current_snapshot_when_incoming_is_superseded() {
2425        let store = InMemoryRuntimeStore::new();
2426        let rid = LogicalRuntimeId::new("runtime-superseded-terminal");
2427        let incoming = session_with_user("turn input");
2428        let mut current = incoming.clone();
2429        current.push(meerkat_core::types::Message::BlockAssistant(
2430            meerkat_core::types::BlockAssistantMessage {
2431                blocks: vec![meerkat_core::types::AssistantBlock::Text {
2432                    text: "peer response already applied".to_string(),
2433                    meta: None,
2434                }],
2435                stop_reason: meerkat_core::types::StopReason::EndTurn,
2436                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
2437                created_at: meerkat_core::types::message_timestamp_now(),
2438            },
2439        ));
2440        let current_snapshot = serde_json::to_vec(&current).unwrap();
2441        let receipt = make_receipt(RunId::new(), 11);
2442
2443        store
2444            .commit_session_snapshot(
2445                &rid,
2446                SessionDelta {
2447                    session_snapshot: current_snapshot.clone(),
2448                },
2449            )
2450            .await
2451            .unwrap();
2452
2453        let error = store
2454            .atomic_apply(
2455                &rid,
2456                Some(SessionDelta {
2457                    session_snapshot: serde_json::to_vec(&incoming).unwrap(),
2458                }),
2459                receipt.clone(),
2460                vec![],
2461                Some(incoming.id().clone()),
2462            )
2463            .await
2464            .expect_err("superseded atomic commit must be explicitly rejected");
2465        assert!(matches!(
2466            error,
2467            RuntimeStoreError::SessionSnapshotSuperseded { .. }
2468        ));
2469
2470        assert_eq!(
2471            store.load_session_snapshot(&rid).await.unwrap(),
2472            Some(current_snapshot)
2473        );
2474        // The session snapshot was classified superseded and skipped, so the
2475        // boundary receipt for that boundary must NOT advance against the
2476        // retained (more-advanced) session snapshot.
2477        assert_eq!(
2478            store
2479                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
2480                .await
2481                .unwrap(),
2482            None
2483        );
2484    }
2485
2486    #[tokio::test]
2487    async fn atomic_apply_skips_inputs_when_session_snapshot_superseded() {
2488        let store = InMemoryRuntimeStore::new();
2489        let rid = LogicalRuntimeId::new("runtime-superseded-inputs");
2490        let incoming = session_with_user("turn input");
2491        let mut current = incoming.clone();
2492        current.push(meerkat_core::types::Message::BlockAssistant(
2493            meerkat_core::types::BlockAssistantMessage {
2494                blocks: vec![meerkat_core::types::AssistantBlock::Text {
2495                    text: "peer response already applied".to_string(),
2496                    meta: None,
2497                }],
2498                stop_reason: meerkat_core::types::StopReason::EndTurn,
2499                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
2500                created_at: meerkat_core::types::message_timestamp_now(),
2501            },
2502        ));
2503        let current_snapshot = serde_json::to_vec(&current).unwrap();
2504        let receipt = make_receipt(RunId::new(), 21);
2505        let input_id = InputId::new();
2506        let bundle = StoredInputState::new_accepted(input_id.clone());
2507
2508        store
2509            .commit_session_snapshot(
2510                &rid,
2511                SessionDelta {
2512                    session_snapshot: current_snapshot.clone(),
2513                },
2514            )
2515            .await
2516            .unwrap();
2517
2518        let error = store
2519            .atomic_apply(
2520                &rid,
2521                Some(SessionDelta {
2522                    session_snapshot: serde_json::to_vec(&incoming).unwrap(),
2523                }),
2524                receipt.clone(),
2525                vec![persistable(bundle)],
2526                Some(incoming.id().clone()),
2527            )
2528            .await
2529            .expect_err("superseded atomic commit must be explicitly rejected");
2530        assert!(matches!(
2531            error,
2532            RuntimeStoreError::SessionSnapshotSuperseded { .. }
2533        ));
2534
2535        // Snapshot retained, receipt + input-state writes skipped as a unit.
2536        assert_eq!(
2537            store.load_session_snapshot(&rid).await.unwrap(),
2538            Some(current_snapshot)
2539        );
2540        assert_eq!(
2541            store
2542                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
2543                .await
2544                .unwrap(),
2545            None
2546        );
2547        assert!(store.load_input_states(&rid).await.unwrap().is_empty());
2548    }
2549
2550    #[tokio::test]
2551    async fn atomic_apply_allows_first_generated_snapshot_after_placeholder() {
2552        let store = InMemoryRuntimeStore::new();
2553        let rid = LogicalRuntimeId::new("runtime-placeholder");
2554        let mut placeholder = meerkat_core::Session::new();
2555        placeholder.set_system_prompt("base system".to_string());
2556        let mut incoming = meerkat_core::Session::with_id(placeholder.id().clone());
2557        incoming.set_system_prompt("base system".to_string());
2558        incoming.push(meerkat_core::types::Message::User(
2559            meerkat_core::types::UserMessage::text("verbose first turn".to_string()),
2560        ));
2561        let parent_revision = incoming.transcript_revision().unwrap();
2562        incoming
2563            .commit_transcript_rewrite(
2564                meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2565                vec![meerkat_core::types::Message::User(
2566                    meerkat_core::types::UserMessage::compaction_summary(
2567                        "[Context compacted] first turn",
2568                    ),
2569                )],
2570                meerkat_core::TranscriptRewriteReason::new("compaction"),
2571                Some("meerkat-core".to_string()),
2572                Some(parent_revision),
2573            )
2574            .unwrap();
2575        let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
2576        let receipt = make_receipt(RunId::new(), 12);
2577
2578        store
2579            .commit_session_snapshot(
2580                &rid,
2581                SessionDelta {
2582                    session_snapshot: serde_json::to_vec(&placeholder).unwrap(),
2583                },
2584            )
2585            .await
2586            .unwrap();
2587
2588        store
2589            .atomic_apply(
2590                &rid,
2591                Some(SessionDelta {
2592                    session_snapshot: incoming_snapshot.clone(),
2593                }),
2594                receipt.clone(),
2595                vec![],
2596                Some(incoming.id().clone()),
2597            )
2598            .await
2599            .unwrap();
2600
2601        assert_eq!(
2602            store.load_session_snapshot(&rid).await.unwrap(),
2603            Some(incoming_snapshot)
2604        );
2605        assert_eq!(
2606            store
2607                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
2608                .await
2609                .unwrap(),
2610            Some(receipt)
2611        );
2612    }
2613
2614    #[tokio::test]
2615    async fn atomic_apply_allows_generated_compaction_before_retained_tail() {
2616        let store = InMemoryRuntimeStore::new();
2617        let rid = LogicalRuntimeId::new("runtime-compaction-tail");
2618        let mut previous = meerkat_core::Session::new();
2619        previous.set_system_prompt("runtime system before context refresh".to_string());
2620        previous.push(meerkat_core::types::Message::User(
2621            meerkat_core::types::UserMessage::text("Turn 1 request".to_string()),
2622        ));
2623        previous.push(meerkat_core::types::Message::BlockAssistant(
2624            meerkat_core::types::BlockAssistantMessage {
2625                blocks: vec![meerkat_core::types::AssistantBlock::Text {
2626                    text: "Turn 1 answer".to_string(),
2627                    meta: None,
2628                }],
2629                stop_reason: meerkat_core::types::StopReason::EndTurn,
2630                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
2631                created_at: meerkat_core::types::message_timestamp_now(),
2632            },
2633        ));
2634
2635        let mut incoming = meerkat_core::Session::with_id(previous.id().clone());
2636        incoming.set_system_prompt("runtime system after context refresh".to_string());
2637        incoming.push(meerkat_core::types::Message::User(
2638            meerkat_core::types::UserMessage::text(
2639                "Verbose context that will be compacted".to_string(),
2640            ),
2641        ));
2642        for message in previous.messages()[1..].iter().cloned() {
2643            incoming.push(message);
2644        }
2645        incoming.push(meerkat_core::types::Message::BlockAssistant(
2646            meerkat_core::types::BlockAssistantMessage {
2647                blocks: vec![meerkat_core::types::AssistantBlock::Text {
2648                    text: "Turn 2 generated answer".to_string(),
2649                    meta: None,
2650                }],
2651                stop_reason: meerkat_core::types::StopReason::EndTurn,
2652                identity: meerkat_core::types::TranscriptMessageIdentity::default(),
2653                created_at: meerkat_core::types::message_timestamp_now(),
2654            },
2655        ));
2656        let parent_revision = incoming.transcript_revision().unwrap();
2657        incoming
2658            .commit_transcript_rewrite(
2659                meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2660                vec![meerkat_core::types::Message::User(
2661                    meerkat_core::types::UserMessage::compaction_summary(
2662                        "[Context compacted] Earlier runtime context".to_string(),
2663                    ),
2664                )],
2665                meerkat_core::TranscriptRewriteReason::new("compaction"),
2666                Some("meerkat-core".to_string()),
2667                Some(parent_revision),
2668            )
2669            .unwrap();
2670        let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
2671        let receipt = make_receipt(RunId::new(), 13);
2672
2673        store
2674            .commit_session_snapshot(
2675                &rid,
2676                SessionDelta {
2677                    session_snapshot: serde_json::to_vec(&previous).unwrap(),
2678                },
2679            )
2680            .await
2681            .unwrap();
2682
2683        store
2684            .atomic_apply(
2685                &rid,
2686                Some(SessionDelta {
2687                    session_snapshot: incoming_snapshot.clone(),
2688                }),
2689                receipt.clone(),
2690                vec![],
2691                Some(incoming.id().clone()),
2692            )
2693            .await
2694            .unwrap();
2695
2696        assert_eq!(
2697            store.load_session_snapshot(&rid).await.unwrap(),
2698            Some(incoming_snapshot)
2699        );
2700        assert_eq!(
2701            store
2702                .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
2703                .await
2704                .unwrap(),
2705            Some(receipt)
2706        );
2707    }
2708
2709    #[tokio::test]
2710    async fn commit_machine_lifecycle_persists_binding_facts() {
2711        use crate::runtime_state::RuntimeState;
2712
2713        let store = InMemoryRuntimeStore::new();
2714        let rid = LogicalRuntimeId::new("runtime-binding");
2715        let binding = MachineLifecycleBindingFacts::new(
2716            Some("rt:session:abc".to_string()),
2717            Some(7),
2718            Some(3),
2719            Some("epoch-1".to_string()),
2720        );
2721
2722        store
2723            .commit_machine_lifecycle(
2724                &rid,
2725                MachineLifecycleCommit::new_with_binding(
2726                    RuntimeState::Retired,
2727                    binding.clone(),
2728                    crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
2729                ),
2730                &[],
2731            )
2732            .await
2733            .unwrap();
2734
2735        let lifecycle = crate::store::load_machine_lifecycle(&store, &rid)
2736            .await
2737            .unwrap()
2738            .expect("machine lifecycle snapshot");
2739        assert_eq!(lifecycle.runtime_state(), RuntimeState::Retired);
2740        assert_eq!(lifecycle.binding(), &binding);
2741        assert_eq!(
2742            crate::store::load_runtime_state(&store, &rid)
2743                .await
2744                .unwrap(),
2745            Some(RuntimeState::Retired)
2746        );
2747    }
2748
2749    #[tokio::test]
2750    async fn concurrent_ops_initializers_return_one_canonical_snapshot() {
2751        let store = InMemoryRuntimeStore::new();
2752        let runtime_id = LogicalRuntimeId::new("runtime-concurrent-ops-initialize");
2753        let registry = crate::ops_lifecycle::RuntimeOpsLifecycleRegistry::new();
2754        let first_candidate = registry
2755            .capture_persistence_snapshot(
2756                meerkat_core::RuntimeEpochId::new(),
2757                &meerkat_core::EpochCursorState::new(),
2758            )
2759            .unwrap();
2760        let second_candidate = registry
2761            .capture_persistence_snapshot(
2762                meerkat_core::RuntimeEpochId::new(),
2763                &meerkat_core::EpochCursorState::new(),
2764            )
2765            .unwrap();
2766        assert_ne!(first_candidate.epoch_id, second_candidate.epoch_id);
2767
2768        let (first, second) = tokio::join!(
2769            store.initialize_ops_lifecycle_if_absent(&runtime_id, &first_candidate),
2770            store.initialize_ops_lifecycle_if_absent(&runtime_id, &second_candidate),
2771        );
2772        let first = first.unwrap();
2773        let second = second.unwrap();
2774
2775        assert_eq!(first.epoch_id, second.epoch_id);
2776        assert_eq!(
2777            store
2778                .load_ops_lifecycle(&runtime_id)
2779                .await
2780                .unwrap()
2781                .expect("canonical snapshot")
2782                .epoch_id,
2783            first.epoch_id
2784        );
2785    }
2786
2787    #[tokio::test]
2788    async fn unregister_finalization_atomically_retires_ops_epoch_and_is_idempotent() {
2789        let store = InMemoryRuntimeStore::new();
2790        let reopened = store.clone();
2791        let runtime_id = LogicalRuntimeId::new("runtime-unregister-finalization");
2792        let stale_ops = crate::ops_lifecycle::RuntimeOpsLifecycleRegistry::new()
2793            .capture_persistence_snapshot(
2794                meerkat_core::RuntimeEpochId::new(),
2795                &meerkat_core::EpochCursorState::new(),
2796            )
2797            .unwrap();
2798        store
2799            .persist_ops_lifecycle(&runtime_id, &stale_ops)
2800            .await
2801            .unwrap();
2802        let retired_ops_epoch = stale_ops.epoch_id.clone();
2803
2804        for _ in 0..2 {
2805            store
2806                .commit_unregister_finalization(
2807                    &runtime_id,
2808                    crate::store::UnregisterFinalizationCommit::new(
2809                        MachineLifecycleCommit::new_with_binding(
2810                            RuntimeState::Stopped,
2811                            MachineLifecycleBindingFacts::new(None, None, None, None),
2812                            crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
2813                        ),
2814                        vec![],
2815                        retired_ops_epoch.clone(),
2816                        crate::meerkat_machine::DeleteOpsFinalizationAuthority::for_store_test(),
2817                    ),
2818                )
2819                .await
2820                .unwrap();
2821        }
2822
2823        assert_eq!(
2824            crate::store::load_runtime_state(&reopened, &runtime_id)
2825                .await
2826                .unwrap(),
2827            Some(RuntimeState::Stopped)
2828        );
2829        assert!(
2830            reopened
2831                .load_ops_lifecycle(&runtime_id)
2832                .await
2833                .unwrap()
2834                .is_none(),
2835            "the same critical section that publishes terminal lifecycle must remove the ops epoch"
2836        );
2837        let late_error = reopened
2838            .persist_ops_lifecycle(&runtime_id, &stale_ops)
2839            .await
2840            .expect_err("a detached callback must not resurrect its retired ops epoch");
2841        assert!(matches!(
2842            late_error,
2843            RuntimeStoreError::OpsLifecycleEpochRetired { epoch_id, .. }
2844                if epoch_id == retired_ops_epoch
2845        ));
2846        assert!(matches!(
2847            reopened
2848                .initialize_ops_lifecycle_if_absent(&runtime_id, &stale_ops)
2849                .await
2850                .expect_err("initialization must honor the same retired-epoch fence"),
2851            RuntimeStoreError::OpsLifecycleEpochRetired { epoch_id, .. }
2852                if epoch_id == retired_ops_epoch
2853        ));
2854        assert!(
2855            reopened
2856                .load_ops_lifecycle(&runtime_id)
2857                .await
2858                .unwrap()
2859                .is_none()
2860        );
2861    }
2862
2863    #[tokio::test]
2864    async fn delayed_old_epoch_finalizer_cannot_delete_or_overwrite_new_ops_epoch() {
2865        let store = InMemoryRuntimeStore::new();
2866        let runtime_id = LogicalRuntimeId::new("runtime-old-finalizer-new-epoch");
2867        let registry = crate::ops_lifecycle::RuntimeOpsLifecycleRegistry::new();
2868        let old_ops = registry
2869            .capture_persistence_snapshot(
2870                meerkat_core::RuntimeEpochId::new(),
2871                &meerkat_core::EpochCursorState::new(),
2872            )
2873            .unwrap();
2874        let new_ops = registry
2875            .capture_persistence_snapshot(
2876                meerkat_core::RuntimeEpochId::new(),
2877                &meerkat_core::EpochCursorState::new(),
2878            )
2879            .unwrap();
2880        store
2881            .persist_ops_lifecycle(&runtime_id, &old_ops)
2882            .await
2883            .unwrap();
2884        store
2885            .persist_ops_lifecycle(&runtime_id, &new_ops)
2886            .await
2887            .unwrap();
2888
2889        store
2890            .commit_unregister_finalization(
2891                &runtime_id,
2892                crate::store::UnregisterFinalizationCommit::new(
2893                    MachineLifecycleCommit::new_with_binding(
2894                        RuntimeState::Stopped,
2895                        MachineLifecycleBindingFacts::new(None, None, None, None),
2896                        crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
2897                    ),
2898                    vec![],
2899                    old_ops.epoch_id.clone(),
2900                    crate::meerkat_machine::DeleteOpsFinalizationAuthority::for_store_test(),
2901                ),
2902            )
2903            .await
2904            .unwrap();
2905
2906        assert_eq!(
2907            store
2908                .load_ops_lifecycle(&runtime_id)
2909                .await
2910                .unwrap()
2911                .expect("new epoch row must survive delayed old finalization")
2912                .epoch_id,
2913            new_ops.epoch_id
2914        );
2915        assert!(matches!(
2916            store
2917                .persist_ops_lifecycle(&runtime_id, &old_ops)
2918                .await
2919                .expect_err("retired old epoch stays fenced"),
2920            RuntimeStoreError::OpsLifecycleEpochRetired { .. }
2921        ));
2922        store
2923            .persist_ops_lifecycle(&runtime_id, &new_ops)
2924            .await
2925            .unwrap();
2926    }
2927
2928    #[tokio::test]
2929    async fn clear_session_snapshot_if_current_sets_quarantine_marker_cleared_on_write() {
2930        let store = InMemoryRuntimeStore::new();
2931        let rid = LogicalRuntimeId::new("runtime-quarantine");
2932        let rejected = serde_json::to_vec(&session_with_user("rejected")).unwrap();
2933
2934        assert!(!store.is_runtime_projection_quarantined(&rid).await.unwrap());
2935        store
2936            .commit_session_snapshot(
2937                &rid,
2938                SessionDelta {
2939                    session_snapshot: rejected.clone(),
2940                },
2941            )
2942            .await
2943            .unwrap();
2944        assert!(
2945            store
2946                .clear_session_snapshot_if_current(&rid, &rejected)
2947                .await
2948                .unwrap()
2949        );
2950        assert!(
2951            store.is_runtime_projection_quarantined(&rid).await.unwrap(),
2952            "clearing the rejected snapshot must record the in-memory quarantine marker"
2953        );
2954
2955        // A live snapshot write reclaims runtime authority and clears the marker.
2956        store
2957            .commit_session_snapshot(
2958                &rid,
2959                SessionDelta {
2960                    session_snapshot: serde_json::to_vec(&session_with_user("revived")).unwrap(),
2961                },
2962            )
2963            .await
2964            .unwrap();
2965        assert!(
2966            !store.is_runtime_projection_quarantined(&rid).await.unwrap(),
2967            "a live snapshot write must clear the in-memory quarantine marker"
2968        );
2969    }
2970
2971    #[tokio::test]
2972    async fn lifecycle_observation_and_missing_or_version_cas_are_target_local() {
2973        let store = InMemoryRuntimeStore::new();
2974        let runtime_id = LogicalRuntimeId::new("runtime-lifecycle-cas");
2975        let other_runtime_id = LogicalRuntimeId::new("runtime-lifecycle-other");
2976        assert_eq!(
2977            store.observe_machine_lifecycle(&runtime_id).await.unwrap(),
2978            MachineLifecycleObservation::Missing
2979        );
2980
2981        let MachineLifecycleCasOutcome::Applied { version } = store
2982            .compare_and_swap_machine_lifecycle(
2983                &runtime_id,
2984                MachineLifecycleExpectedVersion::Missing,
2985                lifecycle_commit(&runtime_id, RuntimeState::Idle, 7, 3),
2986            )
2987            .await
2988            .unwrap()
2989        else {
2990            panic!("missing row must be inserted");
2991        };
2992        let observed = store.observe_machine_lifecycle(&runtime_id).await.unwrap();
2993        let MachineLifecycleObservation::Decoded {
2994            record,
2995            version: observed_version,
2996        } = &observed
2997        else {
2998            panic!("committed lifecycle row must decode");
2999        };
3000        assert_eq!(observed_version, &version);
3001        assert_eq!(record.runtime_state(), Some(RuntimeState::Idle));
3002        assert_eq!(record.binding().fence_token(), Some(7));
3003
3004        let conflict = store
3005            .compare_and_swap_machine_lifecycle(
3006                &runtime_id,
3007                MachineLifecycleExpectedVersion::Missing,
3008                lifecycle_commit(&runtime_id, RuntimeState::Stopped, 8, 4),
3009            )
3010            .await
3011            .unwrap();
3012        assert_eq!(
3013            conflict,
3014            MachineLifecycleCasOutcome::Conflict {
3015                current: observed.clone()
3016            }
3017        );
3018        assert_eq!(
3019            store
3020                .observe_machine_lifecycle(&other_runtime_id)
3021                .await
3022                .unwrap(),
3023            MachineLifecycleObservation::Missing
3024        );
3025
3026        assert!(matches!(
3027            store
3028                .compare_and_swap_machine_lifecycle(
3029                    &runtime_id,
3030                    MachineLifecycleExpectedVersion::Version(version),
3031                    lifecycle_commit(&runtime_id, RuntimeState::Stopped, 8, 4),
3032                )
3033                .await
3034                .unwrap(),
3035            MachineLifecycleCasOutcome::Applied { .. }
3036        ));
3037    }
3038
3039    #[tokio::test]
3040    async fn malformed_lifecycle_repair_is_blocked_even_with_apparent_highwater() {
3041        let store = InMemoryRuntimeStore::new();
3042        let runtime_id = LogicalRuntimeId::new("runtime-malformed-lifecycle");
3043        let raw = serde_json::to_vec(&serde_json::json!({
3044            "record_version": crate::store::MACHINE_LIFECYCLE_STORE_RECORD_VERSION,
3045            "runtime_state": "idle",
3046            "binding": {
3047                "agent_runtime_id": runtime_id.0.clone(),
3048                "fence_token": 9,
3049                "runtime_generation": 5,
3050                "runtime_epoch_id": "epoch-5"
3051            },
3052            "current_run_id": null,
3053            "pre_run_phase": null,
3054            "unregister_progress": null
3055        }))
3056        .unwrap();
3057        store
3058            .inner
3059            .lock()
3060            .await
3061            .runtime_lifecycle
3062            .insert(runtime_id.0.clone(), raw.clone());
3063
3064        let observed = store.observe_machine_lifecycle(&runtime_id).await.unwrap();
3065        let MachineLifecycleObservation::Malformed { version, .. } = observed else {
3066            panic!("structurally incomplete row must remain malformed evidence");
3067        };
3068        assert!(matches!(
3069            store
3070                .compare_and_swap_machine_lifecycle(
3071                    &runtime_id,
3072                    MachineLifecycleExpectedVersion::Version(version.clone()),
3073                    lifecycle_commit(&runtime_id, RuntimeState::Idle, 8, 5),
3074                )
3075                .await
3076                .expect_err("repair must not lower an independently readable fence"),
3077            RuntimeStoreError::MachineLifecycleRepairBlocked { .. }
3078        ));
3079        assert_eq!(
3080            store
3081                .load_machine_lifecycle_record(&runtime_id)
3082                .await
3083                .unwrap(),
3084            Some(raw.clone())
3085        );
3086
3087        assert!(matches!(
3088            store
3089                .compare_and_swap_machine_lifecycle(
3090                    &runtime_id,
3091                    MachineLifecycleExpectedVersion::Version(version),
3092                    lifecycle_commit(&runtime_id, RuntimeState::Idle, 10, 6),
3093                )
3094                .await
3095                .expect_err("decodable fragments inside malformed bytes are not repair authority"),
3096            RuntimeStoreError::MachineLifecycleRepairBlocked { .. }
3097        ));
3098        assert_eq!(
3099            store
3100                .load_machine_lifecycle_record(&runtime_id)
3101                .await
3102                .unwrap(),
3103            Some(raw)
3104        );
3105    }
3106
3107    #[tokio::test]
3108    async fn malformed_lifecycle_duplicate_highwater_keys_are_repair_blocked() {
3109        let store = InMemoryRuntimeStore::new();
3110        let runtime_id = LogicalRuntimeId::new("runtime-duplicate-lifecycle-fence");
3111        let raw = format!(
3112            r#"{{"record_version":4,"runtime_state":"idle","binding":{{"agent_runtime_id":"{}","fence_token":99,"fence_token":1,"runtime_generation":3,"runtime_epoch_id":"epoch-3"}},"current_run_id":null,"pre_run_phase":null,"supervisor_authority":{{"kind":"unbound_no_receipt"}},"unregister_progress":null}}"#,
3113            runtime_id.0
3114        )
3115        .into_bytes();
3116        store
3117            .inner
3118            .lock()
3119            .await
3120            .runtime_lifecycle
3121            .insert(runtime_id.0.clone(), raw.clone());
3122        let MachineLifecycleObservation::Malformed { version, .. } =
3123            store.observe_machine_lifecycle(&runtime_id).await.unwrap()
3124        else {
3125            panic!("duplicate high-water keys must classify as malformed");
3126        };
3127
3128        assert!(matches!(
3129            store
3130                .compare_and_swap_machine_lifecycle(
3131                    &runtime_id,
3132                    MachineLifecycleExpectedVersion::Version(version),
3133                    lifecycle_commit(&runtime_id, RuntimeState::Idle, 2, 3),
3134                )
3135                .await
3136                .expect_err("ambiguous duplicate high-water must block repair"),
3137            RuntimeStoreError::MachineLifecycleRepairBlocked { .. }
3138        ));
3139        assert_eq!(
3140            store
3141                .load_machine_lifecycle_record(&runtime_id)
3142                .await
3143                .unwrap(),
3144            Some(raw)
3145        );
3146    }
3147}