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 identity: meerkat_core::types::TranscriptMessageIdentity::default(),
860 created_at: meerkat_core::types::message_timestamp_now(),
861 },
862 ));
863 let current_snapshot = serde_json::to_vec(¤t).unwrap();
864 let receipt = make_receipt(RunId::new(), 11);
865
866 store
867 .commit_session_snapshot(
868 &rid,
869 SessionDelta {
870 session_snapshot: current_snapshot.clone(),
871 },
872 )
873 .await
874 .unwrap();
875
876 store
877 .atomic_apply(
878 &rid,
879 Some(SessionDelta {
880 session_snapshot: serde_json::to_vec(&incoming).unwrap(),
881 }),
882 receipt.clone(),
883 vec![],
884 Some(incoming.id().clone()),
885 )
886 .await
887 .unwrap();
888
889 assert_eq!(
890 store.load_session_snapshot(&rid).await.unwrap(),
891 Some(current_snapshot)
892 );
893 assert_eq!(
897 store
898 .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
899 .await
900 .unwrap(),
901 None
902 );
903 }
904
905 #[tokio::test]
906 async fn atomic_apply_skips_inputs_when_session_snapshot_superseded() {
907 let store = InMemoryRuntimeStore::new();
908 let rid = LogicalRuntimeId::new("runtime-superseded-inputs");
909 let incoming = session_with_user("turn input");
910 let mut current = incoming.clone();
911 current.push(meerkat_core::types::Message::BlockAssistant(
912 meerkat_core::types::BlockAssistantMessage {
913 blocks: vec![meerkat_core::types::AssistantBlock::Text {
914 text: "peer response already applied".to_string(),
915 meta: None,
916 }],
917 stop_reason: meerkat_core::types::StopReason::EndTurn,
918 identity: meerkat_core::types::TranscriptMessageIdentity::default(),
919 created_at: meerkat_core::types::message_timestamp_now(),
920 },
921 ));
922 let current_snapshot = serde_json::to_vec(¤t).unwrap();
923 let receipt = make_receipt(RunId::new(), 21);
924 let input_id = InputId::new();
925 let bundle = StoredInputState::new_accepted(input_id.clone());
926
927 store
928 .commit_session_snapshot(
929 &rid,
930 SessionDelta {
931 session_snapshot: current_snapshot.clone(),
932 },
933 )
934 .await
935 .unwrap();
936
937 store
938 .atomic_apply(
939 &rid,
940 Some(SessionDelta {
941 session_snapshot: serde_json::to_vec(&incoming).unwrap(),
942 }),
943 receipt.clone(),
944 vec![persistable(bundle)],
945 Some(incoming.id().clone()),
946 )
947 .await
948 .unwrap();
949
950 assert_eq!(
952 store.load_session_snapshot(&rid).await.unwrap(),
953 Some(current_snapshot)
954 );
955 assert_eq!(
956 store
957 .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
958 .await
959 .unwrap(),
960 None
961 );
962 assert!(store.load_input_states(&rid).await.unwrap().is_empty());
963 }
964
965 #[tokio::test]
966 async fn atomic_apply_allows_first_generated_snapshot_after_placeholder() {
967 let store = InMemoryRuntimeStore::new();
968 let rid = LogicalRuntimeId::new("runtime-placeholder");
969 let mut placeholder = meerkat_core::Session::new();
970 placeholder.set_system_prompt("base system".to_string());
971 let mut incoming = meerkat_core::Session::with_id(placeholder.id().clone());
972 incoming.set_system_prompt("base system".to_string());
973 incoming.push(meerkat_core::types::Message::User(
974 meerkat_core::types::UserMessage::text("verbose first turn".to_string()),
975 ));
976 let parent_revision = incoming.transcript_revision().unwrap();
977 incoming
978 .commit_transcript_rewrite(
979 meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
980 vec![meerkat_core::types::Message::User(
981 meerkat_core::types::UserMessage::compaction_summary(
982 "[Context compacted] first turn",
983 ),
984 )],
985 meerkat_core::TranscriptRewriteReason::new("compaction"),
986 Some("meerkat-core".to_string()),
987 Some(parent_revision),
988 )
989 .unwrap();
990 let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
991 let receipt = make_receipt(RunId::new(), 12);
992
993 store
994 .commit_session_snapshot(
995 &rid,
996 SessionDelta {
997 session_snapshot: serde_json::to_vec(&placeholder).unwrap(),
998 },
999 )
1000 .await
1001 .unwrap();
1002
1003 store
1004 .atomic_apply(
1005 &rid,
1006 Some(SessionDelta {
1007 session_snapshot: incoming_snapshot.clone(),
1008 }),
1009 receipt.clone(),
1010 vec![],
1011 Some(incoming.id().clone()),
1012 )
1013 .await
1014 .unwrap();
1015
1016 assert_eq!(
1017 store.load_session_snapshot(&rid).await.unwrap(),
1018 Some(incoming_snapshot)
1019 );
1020 assert_eq!(
1021 store
1022 .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1023 .await
1024 .unwrap(),
1025 Some(receipt)
1026 );
1027 }
1028
1029 #[tokio::test]
1030 async fn atomic_apply_allows_generated_compaction_before_retained_tail() {
1031 let store = InMemoryRuntimeStore::new();
1032 let rid = LogicalRuntimeId::new("runtime-compaction-tail");
1033 let mut previous = meerkat_core::Session::new();
1034 previous.set_system_prompt("runtime system before context refresh".to_string());
1035 previous.push(meerkat_core::types::Message::User(
1036 meerkat_core::types::UserMessage::text("Turn 1 request".to_string()),
1037 ));
1038 previous.push(meerkat_core::types::Message::BlockAssistant(
1039 meerkat_core::types::BlockAssistantMessage {
1040 blocks: vec![meerkat_core::types::AssistantBlock::Text {
1041 text: "Turn 1 answer".to_string(),
1042 meta: None,
1043 }],
1044 stop_reason: meerkat_core::types::StopReason::EndTurn,
1045 identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1046 created_at: meerkat_core::types::message_timestamp_now(),
1047 },
1048 ));
1049
1050 let mut incoming = meerkat_core::Session::with_id(previous.id().clone());
1051 incoming.set_system_prompt("runtime system after context refresh".to_string());
1052 incoming.push(meerkat_core::types::Message::User(
1053 meerkat_core::types::UserMessage::text(
1054 "Verbose context that will be compacted".to_string(),
1055 ),
1056 ));
1057 for message in previous.messages()[1..].iter().cloned() {
1058 incoming.push(message);
1059 }
1060 incoming.push(meerkat_core::types::Message::BlockAssistant(
1061 meerkat_core::types::BlockAssistantMessage {
1062 blocks: vec![meerkat_core::types::AssistantBlock::Text {
1063 text: "Turn 2 generated answer".to_string(),
1064 meta: None,
1065 }],
1066 stop_reason: meerkat_core::types::StopReason::EndTurn,
1067 identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1068 created_at: meerkat_core::types::message_timestamp_now(),
1069 },
1070 ));
1071 let parent_revision = incoming.transcript_revision().unwrap();
1072 incoming
1073 .commit_transcript_rewrite(
1074 meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
1075 vec![meerkat_core::types::Message::User(
1076 meerkat_core::types::UserMessage::compaction_summary(
1077 "[Context compacted] Earlier runtime context".to_string(),
1078 ),
1079 )],
1080 meerkat_core::TranscriptRewriteReason::new("compaction"),
1081 Some("meerkat-core".to_string()),
1082 Some(parent_revision),
1083 )
1084 .unwrap();
1085 let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
1086 let receipt = make_receipt(RunId::new(), 13);
1087
1088 store
1089 .commit_session_snapshot(
1090 &rid,
1091 SessionDelta {
1092 session_snapshot: serde_json::to_vec(&previous).unwrap(),
1093 },
1094 )
1095 .await
1096 .unwrap();
1097
1098 store
1099 .atomic_apply(
1100 &rid,
1101 Some(SessionDelta {
1102 session_snapshot: incoming_snapshot.clone(),
1103 }),
1104 receipt.clone(),
1105 vec![],
1106 Some(incoming.id().clone()),
1107 )
1108 .await
1109 .unwrap();
1110
1111 assert_eq!(
1112 store.load_session_snapshot(&rid).await.unwrap(),
1113 Some(incoming_snapshot)
1114 );
1115 assert_eq!(
1116 store
1117 .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1118 .await
1119 .unwrap(),
1120 Some(receipt)
1121 );
1122 }
1123
1124 #[tokio::test]
1125 async fn commit_machine_lifecycle_persists_binding_facts() {
1126 use crate::runtime_state::RuntimeState;
1127
1128 let store = InMemoryRuntimeStore::new();
1129 let rid = LogicalRuntimeId::new("runtime-binding");
1130 let binding = MachineLifecycleBindingFacts::new(
1131 Some("rt:session:abc".to_string()),
1132 Some(7),
1133 Some(3),
1134 Some("epoch-1".to_string()),
1135 );
1136
1137 store
1138 .commit_machine_lifecycle(
1139 &rid,
1140 MachineLifecycleCommit::new_with_binding(RuntimeState::Retired, binding.clone()),
1141 &[],
1142 )
1143 .await
1144 .unwrap();
1145
1146 let lifecycle = crate::store::load_machine_lifecycle(&store, &rid)
1147 .await
1148 .unwrap()
1149 .expect("machine lifecycle snapshot");
1150 assert_eq!(lifecycle.runtime_state(), RuntimeState::Retired);
1151 assert_eq!(lifecycle.binding(), &binding);
1152 assert_eq!(
1153 crate::store::load_runtime_state(&store, &rid)
1154 .await
1155 .unwrap(),
1156 Some(RuntimeState::Retired)
1157 );
1158 }
1159
1160 #[tokio::test]
1161 async fn clear_session_snapshot_if_current_sets_quarantine_marker_cleared_on_write() {
1162 let store = InMemoryRuntimeStore::new();
1163 let rid = LogicalRuntimeId::new("runtime-quarantine");
1164 let rejected = serde_json::to_vec(&session_with_user("rejected")).unwrap();
1165
1166 assert!(!store.is_runtime_projection_quarantined(&rid).await.unwrap());
1167 store
1168 .commit_session_snapshot(
1169 &rid,
1170 SessionDelta {
1171 session_snapshot: rejected.clone(),
1172 },
1173 )
1174 .await
1175 .unwrap();
1176 assert!(
1177 store
1178 .clear_session_snapshot_if_current(&rid, &rejected)
1179 .await
1180 .unwrap()
1181 );
1182 assert!(
1183 store.is_runtime_projection_quarantined(&rid).await.unwrap(),
1184 "clearing the rejected snapshot must record the in-memory quarantine marker"
1185 );
1186
1187 store
1189 .commit_session_snapshot(
1190 &rid,
1191 SessionDelta {
1192 session_snapshot: serde_json::to_vec(&session_with_user("revived")).unwrap(),
1193 },
1194 )
1195 .await
1196 .unwrap();
1197 assert!(
1198 !store.is_runtime_projection_quarantined(&rid).await.unwrap(),
1199 "a live snapshot write must clear the in-memory quarantine marker"
1200 );
1201 }
1202}