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