1use 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#[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#[derive(Debug, Default)]
56struct Inner {
57 input_states: HashMap<String, IndexMap<InputId, StoredInputState>>,
59 receipts: HashMap<ReceiptKey, RunBoundaryReceipt>,
61 sessions: HashMap<String, Vec<u8>>,
63 projection_quarantine: HashSet<String>,
70 runtime_lifecycle: HashMap<String, Vec<u8>>,
74 ops_lifecycle_snapshots: HashMap<String, PersistedOpsSnapshot>,
76 retired_ops_epochs: HashSet<(String, meerkat_core::RuntimeEpochId)>,
79 compaction_projection_outbox:
81 HashMap<String, HashMap<meerkat_core::CompactionProjectionId, CompactionOutboxEntry>>,
82 runtime_delivery_authority: HashMap<String, RuntimeDeliveryAuthorityRecord>,
84 runtime_delivery_records: HashMap<String, BTreeMap<u64, RuntimeDeliveryStoreRecord>>,
86}
87
88#[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 #[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 #[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
197fn 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 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 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 let rid = runtime_id.0.clone();
492
493 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 if session_snapshot_superseded {
585 return Err(RuntimeStoreError::SessionSnapshotSuperseded { runtime_id: rid });
586 }
587
588 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 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 } 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 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 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, ¤t) {
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(¤t)?;
1142 validate_machine_lifecycle_replacement(
1143 ¤t,
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, ¤t) {
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(¤t)?;
1182 validate_machine_lifecycle_replacement(
1183 ¤t,
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 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 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(¤t).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 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 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 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 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 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(¤t).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 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(¤t).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 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 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 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}