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