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, Clone)]
34struct CompactionOutboxEntry {
35 intent: meerkat_core::CompactionProjectionIntent,
36 finalized: bool,
37}
38
39#[derive(Debug, Default)]
41struct Inner {
42 input_states: HashMap<String, IndexMap<InputId, StoredInputState>>,
44 receipts: HashMap<ReceiptKey, RunBoundaryReceipt>,
46 sessions: HashMap<String, Vec<u8>>,
48 projection_quarantine: HashSet<String>,
55 runtime_lifecycle: HashMap<String, MachineLifecycleSnapshot>,
57 ops_lifecycle_snapshots: HashMap<String, PersistedOpsSnapshot>,
59 retired_ops_epochs: HashSet<(String, meerkat_core::RuntimeEpochId)>,
62 compaction_projection_outbox:
64 HashMap<String, HashMap<meerkat_core::CompactionProjectionId, CompactionOutboxEntry>>,
65}
66
67#[derive(Debug, Clone)]
69pub struct InMemoryRuntimeStore {
70 inner: Arc<Mutex<Inner>>,
71 auth_oauth_flow_snapshot: Arc<StdMutex<Option<Vec<u8>>>>,
72}
73
74impl InMemoryRuntimeStore {
75 pub fn new() -> Self {
76 Self {
77 inner: Arc::new(Mutex::new(Inner::default())),
78 auth_oauth_flow_snapshot: Arc::new(StdMutex::new(None)),
79 }
80 }
81}
82
83impl Default for InMemoryRuntimeStore {
84 fn default() -> Self {
85 Self::new()
86 }
87}
88
89fn is_runtime_placeholder_session(session: &meerkat_core::Session) -> bool {
90 session.transcript_history_state().ok().flatten().is_none()
91 && matches!(
92 session.messages(),
93 [] | [meerkat_core::types::Message::System(_)]
94 )
95}
96
97fn deserialize_persisted_session(bytes: &[u8]) -> Result<meerkat_core::Session, RuntimeStoreError> {
103 serde_json::from_slice(bytes).map_err(|err| RuntimeStoreError::ReadFailed(err.to_string()))
104}
105
106fn ensure_compaction_intents_already_outboxed(
107 inner: &Inner,
108 runtime_id: &LogicalRuntimeId,
109 session: &meerkat_core::Session,
110) -> Result<(), RuntimeStoreError> {
111 let intents = super::validated_compaction_projection_intents(session)?;
112 let existing = inner.compaction_projection_outbox.get(&runtime_id.0);
113 for intent in intents {
114 match existing.and_then(|entries| entries.get(&intent.projection)) {
115 Some(entry) if entry.finalized => {
116 return Err(RuntimeStoreError::WriteFailed(format!(
117 "non-boundary snapshot replays finalized compaction intent {}",
118 intent.projection.revision()
119 )));
120 }
121 Some(entry) if entry.intent == intent => {}
122 Some(_) => {
123 return Err(RuntimeStoreError::WriteFailed(format!(
124 "non-boundary snapshot conflicts with compaction outbox rewrite {}",
125 intent.projection.revision()
126 )));
127 }
128 None => {
129 return Err(RuntimeStoreError::WriteFailed(format!(
130 "non-boundary snapshot introduces compaction intent {} without atomic outbox authority",
131 intent.projection.revision()
132 )));
133 }
134 }
135 }
136 Ok(())
137}
138
139#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
140#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
141impl RuntimeStore for InMemoryRuntimeStore {
142 fn supports_compaction_projection_outbox(&self) -> bool {
143 true
144 }
145
146 fn persist_auth_oauth_flow_snapshot(
147 &self,
148 snapshot_json: &[u8],
149 ) -> Result<(), RuntimeStoreError> {
150 *self
151 .auth_oauth_flow_snapshot
152 .lock()
153 .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))? =
154 Some(snapshot_json.to_vec());
155 Ok(())
156 }
157
158 fn load_auth_oauth_flow_snapshot(&self) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
159 self.auth_oauth_flow_snapshot
160 .lock()
161 .map(|snapshot| snapshot.clone())
162 .map_err(|err| RuntimeStoreError::ReadFailed(err.to_string()))
163 }
164
165 fn update_auth_oauth_flow_snapshot(
166 &self,
167 update: &mut AuthOAuthFlowSnapshotUpdate<'_>,
168 ) -> Result<(), RuntimeStoreError> {
169 let mut snapshot = self
170 .auth_oauth_flow_snapshot
171 .lock()
172 .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
173 let next = update(snapshot.as_deref())?;
174 *snapshot = Some(next);
175 Ok(())
176 }
177
178 async fn commit_session_snapshot(
179 &self,
180 runtime_id: &LogicalRuntimeId,
181 session_delta: SessionDelta,
182 ) -> Result<(), RuntimeStoreError> {
183 let incoming: meerkat_core::Session =
184 serde_json::from_slice(&session_delta.session_snapshot)
185 .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
186 let mut inner = self.inner.lock().await;
187 ensure_compaction_intents_already_outboxed(&inner, runtime_id, &incoming)?;
188 let previous = inner
189 .sessions
190 .get(&runtime_id.0)
191 .map(|snapshot| deserialize_persisted_session(snapshot))
192 .transpose()?;
193 meerkat_core::session_store::run_boundary_snapshot_save_guard(&incoming, previous.as_ref())
194 .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
195 inner
196 .sessions
197 .insert(runtime_id.0.clone(), session_delta.session_snapshot);
198 inner.projection_quarantine.remove(&runtime_id.0);
199 Ok(())
200 }
201
202 async fn commit_session_transcript_rewrite_snapshot(
203 &self,
204 runtime_id: &LogicalRuntimeId,
205 session_delta: SessionDelta,
206 commit: &meerkat_core::TranscriptRewriteCommit,
207 ) -> Result<(), RuntimeStoreError> {
208 let incoming: meerkat_core::Session =
209 serde_json::from_slice(&session_delta.session_snapshot)
210 .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
211 let mut inner = self.inner.lock().await;
212 ensure_compaction_intents_already_outboxed(&inner, runtime_id, &incoming)?;
213 let previous = inner
214 .sessions
215 .get(&runtime_id.0)
216 .map(|snapshot| deserialize_persisted_session(snapshot))
217 .transpose()?;
218 meerkat_core::session_store::transcript_rewrite_save_guard(
219 &incoming,
220 previous.as_ref(),
221 commit,
222 )
223 .map_err(|err| match err {
224 meerkat_core::SessionStoreError::TranscriptRevisionConflict {
225 expected,
226 actual,
227 ..
228 } => RuntimeStoreError::TranscriptRevisionConflict { expected, actual },
229 other => RuntimeStoreError::WriteFailed(other.to_string()),
230 })?;
231 inner
232 .sessions
233 .insert(runtime_id.0.clone(), session_delta.session_snapshot);
234 inner.projection_quarantine.remove(&runtime_id.0);
235 Ok(())
236 }
237
238 async fn atomic_apply(
239 &self,
240 runtime_id: &LogicalRuntimeId,
241 session_delta: Option<SessionDelta>,
242 receipt: RunBoundaryReceipt,
243 input_updates: Vec<InputStatePersistenceRecord>,
244 session_store_key: Option<meerkat_core::types::SessionId>,
245 ) -> Result<(), RuntimeStoreError> {
246 let mut inner = self.inner.lock().await;
247
248 let rid = runtime_id.0.clone();
250
251 let mut session_snapshot_superseded = false;
258 let mut compaction_intents = Vec::new();
259 if let Some(delta) = session_delta {
260 let incoming_session =
261 serde_json::from_slice::<meerkat_core::Session>(&delta.session_snapshot);
262 let mut persist_session_snapshot = true;
263 match (incoming_session, session_store_key) {
264 (Ok(incoming_session), session_store_key) => {
265 compaction_intents =
266 super::validated_compaction_projection_intents(&incoming_session)?;
267 if let Some(existing) = inner.compaction_projection_outbox.get(&rid) {
268 for intent in &compaction_intents {
269 if let Some(entry) = existing.get(&intent.projection) {
270 if entry.finalized {
271 return Err(RuntimeStoreError::WriteFailed(format!(
272 "atomic session snapshot replays finalized compaction intent {}",
273 intent.projection.revision()
274 )));
275 }
276 if entry.intent != *intent {
277 return Err(RuntimeStoreError::WriteFailed(format!(
278 "conflicting compaction outbox intent for rewrite {}",
279 intent.projection.revision()
280 )));
281 }
282 }
283 }
284 }
285 if let Some(session_store_key) = session_store_key
286 && incoming_session.id() != &session_store_key
287 {
288 return Err(RuntimeStoreError::SessionKeyMismatch {
289 expected: session_store_key,
290 actual: incoming_session.id().clone(),
291 });
292 }
293 let previous_session = inner
294 .sessions
295 .get(&rid)
296 .and_then(|snapshot| deserialize_persisted_session(snapshot).ok());
297 if let Err(err) = meerkat_core::session_store::run_boundary_snapshot_save_guard(
298 &incoming_session,
299 previous_session.as_ref(),
300 ) {
301 if previous_session
302 .as_ref()
303 .is_some_and(is_runtime_placeholder_session)
304 {
305 persist_session_snapshot = true;
306 } else if previous_session.as_ref().is_some_and(|previous_session| {
307 meerkat_core::session_store::run_boundary_snapshot_save_guard(
308 previous_session,
309 Some(&incoming_session),
310 )
311 .is_ok()
312 }) {
313 persist_session_snapshot = false;
314 session_snapshot_superseded = true;
315 } else {
316 return Err(RuntimeStoreError::WriteFailed(err.to_string()));
317 }
318 }
319 }
320 (Err(err), Some(session_store_key)) => {
321 return Err(RuntimeStoreError::WriteFailed(format!(
322 "session snapshot for {session_store_key} is not a Session: {err}"
323 )));
324 }
325 (Err(err), None) => {
326 return Err(RuntimeStoreError::WriteFailed(format!(
327 "session snapshot is not a Session: {err}"
328 )));
329 }
330 }
331 if persist_session_snapshot {
332 inner.sessions.insert(rid.clone(), delta.session_snapshot);
333 inner.projection_quarantine.remove(&rid);
334 }
335 }
336
337 if session_snapshot_superseded {
342 return Ok(());
343 }
344
345 let outbox = inner
346 .compaction_projection_outbox
347 .entry(rid.clone())
348 .or_default();
349 for intent in compaction_intents {
350 outbox
351 .entry(intent.projection.clone())
352 .or_insert(CompactionOutboxEntry {
353 intent,
354 finalized: false,
355 });
356 }
357
358 let key = ReceiptKey {
360 runtime_id: rid.clone(),
361 run_id: receipt.run_id.clone(),
362 sequence: receipt.sequence,
363 };
364 inner.receipts.insert(key, receipt);
365
366 let states = inner.input_states.entry(rid).or_default();
368 for record in input_updates {
369 let bundle = record.into_stored();
370 states.insert(bundle.state.input_id.clone(), bundle);
371 }
372
373 Ok(())
374 }
375
376 async fn load_pending_compaction_projections(
377 &self,
378 runtime_id: &LogicalRuntimeId,
379 ) -> Result<Vec<meerkat_core::CompactionProjectionIntent>, RuntimeStoreError> {
380 let inner = self.inner.lock().await;
381 let mut pending = inner
382 .compaction_projection_outbox
383 .get(&runtime_id.0)
384 .into_iter()
385 .flat_map(HashMap::values)
386 .filter(|entry| !entry.finalized)
387 .map(|entry| entry.intent.clone())
388 .collect::<Vec<_>>();
389 pending.sort_by(|left, right| {
390 left.projection
391 .session_id()
392 .to_string()
393 .cmp(&right.projection.session_id().to_string())
394 .then_with(|| {
395 left.projection
396 .parent_revision()
397 .cmp(right.projection.parent_revision())
398 })
399 .then_with(|| left.projection.revision().cmp(right.projection.revision()))
400 .then_with(|| {
401 left.projection
402 .commit_fingerprint()
403 .cmp(right.projection.commit_fingerprint())
404 })
405 });
406 Ok(pending)
407 }
408
409 async fn mark_compaction_projection_finalized(
410 &self,
411 runtime_id: &LogicalRuntimeId,
412 projection: &meerkat_core::CompactionProjectionId,
413 ) -> Result<(), RuntimeStoreError> {
414 let mut inner = self.inner.lock().await;
415 let outbox_exists = inner
416 .compaction_projection_outbox
417 .get(&runtime_id.0)
418 .is_some_and(|entries| entries.contains_key(projection));
419 if !outbox_exists {
420 return Err(RuntimeStoreError::NotFound(format!(
421 "compaction outbox rewrite {}",
422 projection.revision()
423 )));
424 }
425 let cleaned_snapshot = inner
426 .sessions
427 .get(&runtime_id.0)
428 .map(|snapshot| {
429 let mut session = deserialize_persisted_session(snapshot)?;
430 session
431 .complete_compaction_projection_intent(projection)
432 .map_err(|error| RuntimeStoreError::WriteFailed(error.to_string()))?;
433 serde_json::to_vec(&session)
434 .map_err(|error| RuntimeStoreError::WriteFailed(error.to_string()))
435 })
436 .transpose()?;
437 let entry = inner
438 .compaction_projection_outbox
439 .get_mut(&runtime_id.0)
440 .and_then(|entries| entries.get_mut(projection))
441 .ok_or_else(|| {
442 RuntimeStoreError::NotFound(format!(
443 "compaction outbox rewrite {}",
444 projection.revision()
445 ))
446 })?;
447 entry.finalized = true;
448 if let Some(cleaned_snapshot) = cleaned_snapshot {
449 inner
450 .sessions
451 .insert(runtime_id.0.clone(), cleaned_snapshot);
452 }
453 Ok(())
454 }
455
456 async fn load_input_states(
457 &self,
458 runtime_id: &LogicalRuntimeId,
459 ) -> Result<Vec<StoredInputState>, RuntimeStoreError> {
460 let inner = self.inner.lock().await;
461 let states = inner
462 .input_states
463 .get(&runtime_id.0)
464 .map(|m| m.values().cloned().collect())
465 .unwrap_or_default();
466 Ok(states)
467 }
468
469 async fn load_boundary_receipt(
470 &self,
471 runtime_id: &LogicalRuntimeId,
472 run_id: &RunId,
473 sequence: u64,
474 ) -> Result<Option<RunBoundaryReceipt>, RuntimeStoreError> {
475 let inner = self.inner.lock().await;
476 let key = ReceiptKey {
477 runtime_id: runtime_id.0.clone(),
478 run_id: run_id.clone(),
479 sequence,
480 };
481 Ok(inner.receipts.get(&key).cloned())
482 }
483
484 async fn load_session_snapshot(
485 &self,
486 runtime_id: &LogicalRuntimeId,
487 ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
488 let inner = self.inner.lock().await;
489 Ok(inner.sessions.get(&runtime_id.0).cloned())
490 }
491
492 async fn clear_session_snapshot(
493 &self,
494 runtime_id: &LogicalRuntimeId,
495 ) -> Result<(), RuntimeStoreError> {
496 let mut inner = self.inner.lock().await;
497 inner.sessions.remove(&runtime_id.0);
498 Ok(())
499 }
500
501 async fn replace_session_snapshot_if_current(
502 &self,
503 runtime_id: &LogicalRuntimeId,
504 expected_current: &[u8],
505 replacement: Vec<u8>,
506 ) -> Result<bool, RuntimeStoreError> {
507 let replacement_session: meerkat_core::Session = serde_json::from_slice(&replacement)
508 .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
509 let mut inner = self.inner.lock().await;
510 let Some(current) = inner.sessions.get(&runtime_id.0) else {
511 return Ok(false);
512 };
513 if current.as_slice() != expected_current {
514 return Ok(false);
515 }
516 ensure_compaction_intents_already_outboxed(&inner, runtime_id, &replacement_session)?;
517 inner.sessions.insert(runtime_id.0.clone(), replacement);
518 inner.projection_quarantine.remove(&runtime_id.0);
519 Ok(true)
520 }
521
522 async fn clear_session_snapshot_if_current(
523 &self,
524 runtime_id: &LogicalRuntimeId,
525 expected_current: &[u8],
526 ) -> Result<bool, RuntimeStoreError> {
527 let mut inner = self.inner.lock().await;
528 let Some(current) = inner.sessions.get(&runtime_id.0) else {
529 return Ok(false);
530 };
531 if current.as_slice() != expected_current {
532 return Ok(false);
533 }
534 inner.sessions.remove(&runtime_id.0);
535 inner.projection_quarantine.insert(runtime_id.0.clone());
538 Ok(true)
539 }
540
541 async fn is_runtime_projection_quarantined(
542 &self,
543 runtime_id: &LogicalRuntimeId,
544 ) -> Result<bool, RuntimeStoreError> {
545 let inner = self.inner.lock().await;
546 Ok(inner.projection_quarantine.contains(&runtime_id.0))
547 }
548
549 async fn persist_input_state(
550 &self,
551 runtime_id: &LogicalRuntimeId,
552 state: &InputStatePersistenceRecord,
553 ) -> Result<(), RuntimeStoreError> {
554 let mut inner = self.inner.lock().await;
555 let states = inner.input_states.entry(runtime_id.0.clone()).or_default();
556 let bundle = state.as_stored();
557 states.insert(bundle.state.input_id.clone(), bundle.clone());
558 Ok(())
559 }
560
561 async fn load_input_state(
562 &self,
563 runtime_id: &LogicalRuntimeId,
564 input_id: &InputId,
565 ) -> Result<Option<StoredInputState>, RuntimeStoreError> {
566 let inner = self.inner.lock().await;
567 let state = inner
568 .input_states
569 .get(&runtime_id.0)
570 .and_then(|m| m.get(input_id).cloned());
571 Ok(state)
572 }
573
574 async fn load_machine_lifecycle_record(
575 &self,
576 runtime_id: &LogicalRuntimeId,
577 ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
578 let inner = self.inner.lock().await;
579 inner
580 .runtime_lifecycle
581 .get(&runtime_id.0)
582 .map(|snapshot| MachineLifecycleStoreRecord::from_snapshot(snapshot).encode())
583 .transpose()
584 }
585
586 async fn commit_machine_lifecycle(
587 &self,
588 runtime_id: &LogicalRuntimeId,
589 commit: MachineLifecycleCommit,
590 input_states: &[InputStatePersistenceRecord],
591 ) -> Result<(), RuntimeStoreError> {
592 let mut inner = self.inner.lock().await;
593 let rid = runtime_id.0.clone();
594
595 inner
597 .runtime_lifecycle
598 .insert(rid.clone(), commit.into_snapshot());
599 let states = inner.input_states.entry(rid).or_default();
600 for record in input_states {
601 let bundle = record.as_stored();
602 states.insert(bundle.state.input_id.clone(), bundle.clone());
603 }
604
605 Ok(())
606 }
607
608 async fn commit_unregister_finalization(
609 &self,
610 runtime_id: &LogicalRuntimeId,
611 commit: MachineLifecycleCommit,
612 input_states: &[InputStatePersistenceRecord],
613 ) -> Result<(), RuntimeStoreError> {
614 let mut inner = self.inner.lock().await;
615 let rid = runtime_id.0.clone();
616 let retired_ops_epoch = commit.retired_ops_epoch().cloned().ok_or_else(|| {
617 RuntimeStoreError::WriteFailed(
618 "unregister finalization missing exact retired ops epoch".into(),
619 )
620 })?;
621
622 inner
626 .runtime_lifecycle
627 .insert(rid.clone(), commit.into_snapshot());
628 let states = inner.input_states.entry(rid.clone()).or_default();
629 for record in input_states {
630 let bundle = record.as_stored();
631 states.insert(bundle.state.input_id.clone(), bundle.clone());
632 }
633 if inner
634 .ops_lifecycle_snapshots
635 .get(&rid)
636 .is_some_and(|snapshot| snapshot.epoch_id == retired_ops_epoch)
637 {
638 inner.ops_lifecycle_snapshots.remove(&rid);
639 }
640 inner.retired_ops_epochs.insert((rid, retired_ops_epoch));
641 Ok(())
642 }
643
644 async fn persist_ops_lifecycle(
645 &self,
646 runtime_id: &LogicalRuntimeId,
647 snapshot: &PersistedOpsSnapshot,
648 ) -> Result<(), RuntimeStoreError> {
649 let mut inner = self.inner.lock().await;
650 if inner
651 .retired_ops_epochs
652 .contains(&(runtime_id.0.clone(), snapshot.epoch_id.clone()))
653 {
654 return Err(RuntimeStoreError::OpsLifecycleEpochRetired {
655 runtime_id: runtime_id.0.clone(),
656 epoch_id: snapshot.epoch_id.clone(),
657 });
658 }
659 inner
660 .ops_lifecycle_snapshots
661 .insert(runtime_id.0.clone(), snapshot.clone());
662 Ok(())
663 }
664
665 async fn load_ops_lifecycle(
666 &self,
667 runtime_id: &LogicalRuntimeId,
668 ) -> Result<Option<PersistedOpsSnapshot>, RuntimeStoreError> {
669 let inner = self.inner.lock().await;
670 Ok(inner.ops_lifecycle_snapshots.get(&runtime_id.0).cloned())
671 }
672
673 async fn delete_ops_lifecycle(
674 &self,
675 runtime_id: &LogicalRuntimeId,
676 ) -> Result<(), RuntimeStoreError> {
677 let mut inner = self.inner.lock().await;
678 inner.ops_lifecycle_snapshots.remove(&runtime_id.0);
679 Ok(())
680 }
681}
682
683#[cfg(test)]
684#[allow(clippy::unwrap_used)]
685mod tests {
686 use super::*;
687 use crate::RuntimeState;
688 use crate::store::MachineLifecycleBindingFacts;
689 use meerkat_core::lifecycle::run_primitive::RunApplyBoundary;
690
691 fn make_receipt(run_id: RunId, seq: u64) -> RunBoundaryReceipt {
692 RunBoundaryReceipt {
693 run_id,
694 boundary: RunApplyBoundary::RunStart,
695 contributing_input_ids: vec![],
696 conversation_digest: None,
697 message_count: 0,
698 sequence: seq,
699 }
700 }
701
702 fn persistable(bundle: StoredInputState) -> InputStatePersistenceRecord {
703 InputStatePersistenceRecord::from_machine_snapshot(bundle).unwrap()
704 }
705
706 fn session_with_user(content: &str) -> meerkat_core::Session {
707 let mut session = meerkat_core::Session::new();
708 session.push(meerkat_core::types::Message::User(
709 meerkat_core::types::UserMessage::text(content.to_string()),
710 ));
711 session
712 }
713
714 fn session_with_compaction_intent() -> (
715 meerkat_core::Session,
716 meerkat_core::CompactionProjectionIntent,
717 ) {
718 let mut session = session_with_user("verbose context one");
719 session.push(meerkat_core::types::Message::User(
720 meerkat_core::types::UserMessage::text("verbose context two"),
721 ));
722 let parent = session.transcript_revision().unwrap();
723 session
724 .commit_transcript_rewrite(
725 meerkat_core::TranscriptRewriteSelection::MessageRange { start: 0, end: 2 },
726 vec![meerkat_core::types::Message::User(
727 meerkat_core::types::UserMessage::compaction_summary("compacted context"),
728 )],
729 meerkat_core::TranscriptRewriteReason::new("compaction"),
730 Some("runtime-store-test".to_string()),
731 Some(parent),
732 )
733 .unwrap();
734 let mut encoded = serde_json::to_value(&session).unwrap();
735 encoded["metadata"][meerkat_core::SESSION_TRANSCRIPT_HISTORY_STATE_KEY]["commits"][0]["selection"] = serde_json::json!({
736 "type": "compaction_message_range",
737 "range": { "start": 0, "end": 2 }
738 });
739 let mut session: meerkat_core::Session = serde_json::from_value(encoded).unwrap();
740 let commit = session
741 .transcript_history_state()
742 .unwrap()
743 .unwrap()
744 .commits
745 .last()
746 .unwrap()
747 .clone();
748 let intent = meerkat_core::CompactionProjectionIntent {
749 projection: serde_json::from_value(serde_json::json!({
750 "session_id": session.id(),
751 "parent_revision": &commit.parent_revision,
752 "revision": &commit.revision,
753 "commit_fingerprint": "sha256:827d8ee5666e51b2ced4d303640740680d96151d92187fd6e981c29550072c62",
754 }))
755 .unwrap(),
756 summary_tokens: 5,
757 messages_before: 2,
758 messages_after: 1,
759 };
760 session
761 .add_compaction_projection_intent(intent.clone())
762 .unwrap();
763 (session, intent)
764 }
765
766 fn snapshot_with_raw_intents(
767 session: &meerkat_core::Session,
768 intents: &[meerkat_core::CompactionProjectionIntent],
769 ) -> Vec<u8> {
770 let mut value = serde_json::to_value(session).unwrap();
771 value["metadata"][meerkat_core::memory::SESSION_COMPACTION_PROJECTION_INTENTS_KEY] =
772 serde_json::to_value(intents).unwrap();
773 serde_json::to_vec(&value).unwrap()
774 }
775
776 fn unbacked_intent(
777 session_id: &meerkat_core::types::SessionId,
778 ) -> meerkat_core::CompactionProjectionIntent {
779 meerkat_core::CompactionProjectionIntent {
780 projection: serde_json::from_value(serde_json::json!({
781 "session_id": session_id,
782 "parent_revision": "missing-parent",
783 "revision": "missing-revision",
784 "commit_fingerprint": "sha256:unbacked-persisted-fixture",
785 }))
786 .unwrap(),
787 summary_tokens: 1,
788 messages_before: 2,
789 messages_after: 1,
790 }
791 }
792
793 #[tokio::test]
794 async fn atomic_apply_commits_rewrite_and_compaction_outbox_as_one_boundary() {
795 let store = InMemoryRuntimeStore::new();
796 let rid = LogicalRuntimeId::new("runtime-compaction-outbox");
797 let (session, intent) = session_with_compaction_intent();
798 let snapshot = serde_json::to_vec(&session).unwrap();
799 store
800 .atomic_apply(
801 &rid,
802 Some(SessionDelta {
803 session_snapshot: snapshot.clone(),
804 }),
805 make_receipt(RunId::new(), 1),
806 vec![],
807 Some(session.id().clone()),
808 )
809 .await
810 .unwrap();
811 assert_eq!(
812 store.load_session_snapshot(&rid).await.unwrap(),
813 Some(snapshot)
814 );
815 assert_eq!(
816 store
817 .load_pending_compaction_projections(&rid)
818 .await
819 .unwrap(),
820 vec![intent.clone()]
821 );
822 store
823 .mark_compaction_projection_finalized(&rid, &intent.projection)
824 .await
825 .unwrap();
826 store
827 .mark_compaction_projection_finalized(&rid, &intent.projection)
828 .await
829 .unwrap();
830 assert!(
831 store
832 .load_pending_compaction_projections(&rid)
833 .await
834 .unwrap()
835 .is_empty()
836 );
837 let persisted: meerkat_core::Session =
838 serde_json::from_slice(&store.load_session_snapshot(&rid).await.unwrap().unwrap())
839 .unwrap();
840 assert!(
841 persisted
842 .compaction_projection_intents()
843 .unwrap()
844 .is_empty()
845 );
846 }
847
848 #[tokio::test]
849 async fn finalized_outbox_tombstone_rejects_atomic_and_non_boundary_snapshot_replay() {
850 let store = InMemoryRuntimeStore::new();
851 let rid = LogicalRuntimeId::new("runtime-finalized-compaction-replay");
852 let (session, intent) = session_with_compaction_intent();
853 let replay_snapshot = serde_json::to_vec(&session).unwrap();
854 let commit = session
855 .transcript_history_state()
856 .unwrap()
857 .unwrap()
858 .commits
859 .last()
860 .unwrap()
861 .clone();
862 store
863 .atomic_apply(
864 &rid,
865 Some(SessionDelta {
866 session_snapshot: replay_snapshot.clone(),
867 }),
868 make_receipt(RunId::new(), 1),
869 vec![],
870 Some(session.id().clone()),
871 )
872 .await
873 .unwrap();
874 store
875 .mark_compaction_projection_finalized(&rid, &intent.projection)
876 .await
877 .unwrap();
878 let cleaned_snapshot = store.load_session_snapshot(&rid).await.unwrap().unwrap();
879
880 let replay_run_id = RunId::new();
881 let error = store
882 .atomic_apply(
883 &rid,
884 Some(SessionDelta {
885 session_snapshot: replay_snapshot.clone(),
886 }),
887 make_receipt(replay_run_id.clone(), 2),
888 vec![],
889 Some(session.id().clone()),
890 )
891 .await
892 .unwrap_err();
893 assert!(error.to_string().contains("finalized compaction intent"));
894 assert!(
895 store
896 .load_boundary_receipt(&rid, &replay_run_id, 2)
897 .await
898 .unwrap()
899 .is_none(),
900 "finalized replay rejection must roll back the whole atomic boundary"
901 );
902
903 let error = store
904 .commit_session_snapshot(
905 &rid,
906 SessionDelta {
907 session_snapshot: replay_snapshot.clone(),
908 },
909 )
910 .await
911 .unwrap_err();
912 assert!(error.to_string().contains("finalized compaction intent"));
913 let error = store
914 .commit_session_transcript_rewrite_snapshot(
915 &rid,
916 SessionDelta {
917 session_snapshot: replay_snapshot.clone(),
918 },
919 &commit,
920 )
921 .await
922 .unwrap_err();
923 assert!(error.to_string().contains("finalized compaction intent"));
924 let error = store
925 .replace_session_snapshot_if_current(&rid, &cleaned_snapshot, replay_snapshot)
926 .await
927 .unwrap_err();
928 assert!(error.to_string().contains("finalized compaction intent"));
929
930 assert_eq!(
931 store.load_session_snapshot(&rid).await.unwrap(),
932 Some(cleaned_snapshot)
933 );
934 assert!(
935 store
936 .load_pending_compaction_projections(&rid)
937 .await
938 .unwrap()
939 .is_empty(),
940 "a finalized tombstone must never be silently revived or left untracked"
941 );
942 }
943
944 #[tokio::test]
945 async fn invalid_compaction_intent_leaves_snapshot_and_outbox_unmodified() {
946 let store = InMemoryRuntimeStore::new();
947 let rid = LogicalRuntimeId::new("runtime-invalid-compaction-outbox");
948 let (session, mut intent) = session_with_compaction_intent();
949 intent.summary_tokens += 1;
950 let conflicting = vec![
951 session.compaction_projection_intents().unwrap()[0].clone(),
952 intent,
953 ];
954 let error = store
955 .atomic_apply(
956 &rid,
957 Some(SessionDelta {
958 session_snapshot: snapshot_with_raw_intents(&session, &conflicting),
959 }),
960 make_receipt(RunId::new(), 2),
961 vec![],
962 Some(session.id().clone()),
963 )
964 .await
965 .unwrap_err();
966 assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
967 assert_eq!(store.load_session_snapshot(&rid).await.unwrap(), None);
968 assert!(
969 store
970 .load_pending_compaction_projections(&rid)
971 .await
972 .unwrap()
973 .is_empty()
974 );
975
976 let foreign = session_with_compaction_intent().1;
977 for (sequence, invalid) in [foreign, unbacked_intent(session.id())]
978 .into_iter()
979 .enumerate()
980 {
981 let error = store
982 .atomic_apply(
983 &rid,
984 Some(SessionDelta {
985 session_snapshot: snapshot_with_raw_intents(&session, &[invalid]),
986 }),
987 make_receipt(RunId::new(), 10 + sequence as u64),
988 vec![],
989 Some(session.id().clone()),
990 )
991 .await
992 .unwrap_err();
993 assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
994 assert_eq!(store.load_session_snapshot(&rid).await.unwrap(), None);
995 assert!(
996 store
997 .load_pending_compaction_projections(&rid)
998 .await
999 .unwrap()
1000 .is_empty()
1001 );
1002 }
1003 }
1004
1005 #[tokio::test]
1006 async fn superseded_snapshot_does_not_advance_compaction_outbox() {
1007 let store = InMemoryRuntimeStore::new();
1008 let rid = LogicalRuntimeId::new("runtime-superseded-compaction-outbox");
1009 let (incoming, intent) = session_with_compaction_intent();
1010 let mut current = incoming.clone();
1011 current
1012 .complete_compaction_projection_intent(&intent.projection)
1013 .unwrap();
1014 current.push(meerkat_core::types::Message::User(
1015 meerkat_core::types::UserMessage::text("already advanced"),
1016 ));
1017 let current_snapshot = serde_json::to_vec(¤t).unwrap();
1018 store
1019 .commit_session_snapshot(
1020 &rid,
1021 SessionDelta {
1022 session_snapshot: current_snapshot.clone(),
1023 },
1024 )
1025 .await
1026 .unwrap();
1027 store
1028 .atomic_apply(
1029 &rid,
1030 Some(SessionDelta {
1031 session_snapshot: serde_json::to_vec(&incoming).unwrap(),
1032 }),
1033 make_receipt(RunId::new(), 3),
1034 vec![],
1035 Some(incoming.id().clone()),
1036 )
1037 .await
1038 .unwrap();
1039 assert_eq!(
1040 store.load_session_snapshot(&rid).await.unwrap(),
1041 Some(current_snapshot)
1042 );
1043 assert!(
1044 store
1045 .load_pending_compaction_projections(&rid)
1046 .await
1047 .unwrap()
1048 .is_empty()
1049 );
1050 }
1051
1052 #[tokio::test]
1053 async fn existing_outbox_rejects_changed_intent_without_advancing_snapshot() {
1054 let store = InMemoryRuntimeStore::new();
1055 let rid = LogicalRuntimeId::new("runtime-conflicting-compaction-outbox");
1056 let (session, intent) = session_with_compaction_intent();
1057 let original_snapshot = serde_json::to_vec(&session).unwrap();
1058 store
1059 .atomic_apply(
1060 &rid,
1061 Some(SessionDelta {
1062 session_snapshot: original_snapshot.clone(),
1063 }),
1064 make_receipt(RunId::new(), 60),
1065 vec![],
1066 Some(session.id().clone()),
1067 )
1068 .await
1069 .unwrap();
1070
1071 let mut advanced = session.clone();
1072 advanced.push(meerkat_core::types::Message::User(
1073 meerkat_core::types::UserMessage::text("later turn"),
1074 ));
1075 let mut conflicting = intent.clone();
1076 conflicting.summary_tokens += 1;
1077 let error = store
1078 .atomic_apply(
1079 &rid,
1080 Some(SessionDelta {
1081 session_snapshot: snapshot_with_raw_intents(&advanced, &[conflicting]),
1082 }),
1083 make_receipt(RunId::new(), 61),
1084 vec![],
1085 Some(session.id().clone()),
1086 )
1087 .await
1088 .unwrap_err();
1089 assert!(matches!(error, RuntimeStoreError::WriteFailed(_)));
1090 assert_eq!(
1091 store.load_session_snapshot(&rid).await.unwrap(),
1092 Some(original_snapshot)
1093 );
1094 assert_eq!(
1095 store
1096 .load_pending_compaction_projections(&rid)
1097 .await
1098 .unwrap(),
1099 vec![intent]
1100 );
1101 }
1102
1103 #[tokio::test]
1104 async fn non_boundary_snapshot_apis_cannot_bypass_compaction_outbox() {
1105 let store = InMemoryRuntimeStore::new();
1106 let rid = LogicalRuntimeId::new("runtime-compaction-bypass");
1107 let (session, _intent) = session_with_compaction_intent();
1108 let snapshot = serde_json::to_vec(&session).unwrap();
1109 let commit = session
1110 .transcript_history_state()
1111 .unwrap()
1112 .unwrap()
1113 .commits
1114 .last()
1115 .unwrap()
1116 .clone();
1117 assert!(
1118 store
1119 .commit_session_snapshot(
1120 &rid,
1121 SessionDelta {
1122 session_snapshot: snapshot.clone(),
1123 },
1124 )
1125 .await
1126 .is_err()
1127 );
1128 assert!(
1129 store
1130 .commit_session_transcript_rewrite_snapshot(
1131 &rid,
1132 SessionDelta {
1133 session_snapshot: snapshot.clone(),
1134 },
1135 &commit,
1136 )
1137 .await
1138 .is_err()
1139 );
1140 assert_eq!(store.load_session_snapshot(&rid).await.unwrap(), None);
1141 let clean = meerkat_core::Session::with_id(session.id().clone());
1142 let clean_snapshot = serde_json::to_vec(&clean).unwrap();
1143 store
1144 .commit_session_snapshot(
1145 &rid,
1146 SessionDelta {
1147 session_snapshot: clean_snapshot.clone(),
1148 },
1149 )
1150 .await
1151 .unwrap();
1152 assert!(
1153 store
1154 .replace_session_snapshot_if_current(&rid, &clean_snapshot, snapshot)
1155 .await
1156 .is_err()
1157 );
1158 assert_eq!(
1159 store.load_session_snapshot(&rid).await.unwrap(),
1160 Some(clean_snapshot)
1161 );
1162 assert!(
1163 store
1164 .load_pending_compaction_projections(&rid)
1165 .await
1166 .unwrap()
1167 .is_empty()
1168 );
1169 }
1170
1171 #[tokio::test]
1172 async fn atomic_apply_roundtrip() {
1173 let store = InMemoryRuntimeStore::new();
1174 let rid = LogicalRuntimeId::new("test-runtime");
1175 let run_id = RunId::new();
1176 let input_id = InputId::new();
1177
1178 let bundle = StoredInputState::new_accepted(input_id.clone());
1179 let receipt = make_receipt(run_id.clone(), 0);
1180
1181 let session = session_with_user("hello");
1182 let session_snapshot = serde_json::to_vec(&session).unwrap();
1183
1184 store
1185 .atomic_apply(
1186 &rid,
1187 Some(SessionDelta { session_snapshot }),
1188 receipt.clone(),
1189 vec![persistable(bundle)],
1190 None,
1191 )
1192 .await
1193 .unwrap();
1194
1195 let states = store.load_input_states(&rid).await.unwrap();
1197 assert_eq!(states.len(), 1);
1198 assert_eq!(states[0].state.input_id, input_id);
1199
1200 let loaded = store.load_boundary_receipt(&rid, &run_id, 0).await.unwrap();
1202 assert!(loaded.is_some());
1203 }
1204
1205 #[tokio::test]
1206 async fn atomic_apply_rejects_non_session_snapshot_without_owner_context() {
1207 let store = InMemoryRuntimeStore::new();
1208 let rid = LogicalRuntimeId::new("test-runtime");
1209 let run_id = RunId::new();
1210 let input_id = InputId::new();
1211
1212 let bundle = StoredInputState::new_accepted(input_id);
1213 let receipt = make_receipt(run_id, 0);
1214
1215 let err = store
1218 .atomic_apply(
1219 &rid,
1220 Some(SessionDelta {
1221 session_snapshot: b"session-data".to_vec(),
1222 }),
1223 receipt,
1224 vec![persistable(bundle)],
1225 None,
1226 )
1227 .await
1228 .expect_err("non-Session snapshot must be rejected");
1229
1230 match err {
1231 RuntimeStoreError::WriteFailed(message) => {
1232 assert!(
1233 message.contains("not a Session"),
1234 "unexpected WriteFailed message: {message}"
1235 );
1236 }
1237 other => panic!("expected WriteFailed, got {other:?}"),
1238 }
1239 }
1240
1241 #[tokio::test]
1242 async fn persist_and_load_single_state() {
1243 let store = InMemoryRuntimeStore::new();
1244 let rid = LogicalRuntimeId::new("test");
1245 let input_id = InputId::new();
1246 let bundle = StoredInputState::new_accepted(input_id.clone());
1247
1248 store
1249 .persist_input_state(&rid, &persistable(bundle))
1250 .await
1251 .unwrap();
1252
1253 let loaded = store.load_input_state(&rid, &input_id).await.unwrap();
1254 assert!(loaded.is_some());
1255 assert_eq!(loaded.unwrap().state.input_id, input_id);
1256 }
1257
1258 #[tokio::test]
1259 async fn load_nonexistent_returns_none() {
1260 let store = InMemoryRuntimeStore::new();
1261 let rid = LogicalRuntimeId::new("test");
1262
1263 let states = store.load_input_states(&rid).await.unwrap();
1264 assert!(states.is_empty());
1265
1266 let state = store.load_input_state(&rid, &InputId::new()).await.unwrap();
1267 assert!(state.is_none());
1268
1269 let receipt = store
1270 .load_boundary_receipt(&rid, &RunId::new(), 0)
1271 .await
1272 .unwrap();
1273 assert!(receipt.is_none());
1274 }
1275
1276 #[tokio::test]
1277 async fn atomic_apply_updates_existing() {
1278 let store = InMemoryRuntimeStore::new();
1279 let rid = LogicalRuntimeId::new("test");
1280 let input_id = InputId::new();
1281
1282 let bundle1 = StoredInputState::new_accepted(input_id.clone());
1284 store
1285 .atomic_apply(
1286 &rid,
1287 None,
1288 make_receipt(RunId::new(), 0),
1289 vec![persistable(bundle1)],
1290 None,
1291 )
1292 .await
1293 .unwrap();
1294
1295 let mut bundle2 = StoredInputState::new_accepted(input_id.clone());
1297 bundle2.seed.phase = crate::input_state::InputLifecycleState::Queued;
1298 store
1299 .atomic_apply(
1300 &rid,
1301 None,
1302 make_receipt(RunId::new(), 1),
1303 vec![persistable(bundle2)],
1304 None,
1305 )
1306 .await
1307 .unwrap();
1308
1309 let states = store.load_input_states(&rid).await.unwrap();
1310 assert_eq!(states.len(), 1);
1311 assert_eq!(
1312 states[0].seed.phase,
1313 crate::input_state::InputLifecycleState::Queued
1314 );
1315 }
1316
1317 #[tokio::test]
1318 async fn atomic_apply_validates_session_store_key_without_aliasing_snapshot() {
1319 let store = InMemoryRuntimeStore::new();
1320 let rid = LogicalRuntimeId::new("runtime-key");
1321 let session = meerkat_core::Session::new();
1322 let session_id = session.id().clone();
1323 let snapshot = serde_json::to_vec(&session).unwrap();
1324
1325 store
1326 .atomic_apply(
1327 &rid,
1328 Some(SessionDelta {
1329 session_snapshot: snapshot.clone(),
1330 }),
1331 make_receipt(RunId::new(), 0),
1332 vec![],
1333 Some(session_id.clone()),
1334 )
1335 .await
1336 .unwrap();
1337
1338 assert_eq!(
1339 store.load_session_snapshot(&rid).await.unwrap(),
1340 Some(snapshot)
1341 );
1342 assert!(
1343 store
1344 .load_session_snapshot(&LogicalRuntimeId::legacy_session_uuid_alias(&session_id))
1345 .await
1346 .unwrap()
1347 .is_none(),
1348 "session_store_key must validate the snapshot identity, not create a raw UUID runtime alias"
1349 );
1350 }
1351
1352 #[tokio::test]
1353 async fn atomic_apply_rejects_mismatched_session_store_key() {
1354 let store = InMemoryRuntimeStore::new();
1355 let rid = LogicalRuntimeId::new("runtime-key");
1356 let session = meerkat_core::Session::new();
1357 let wrong_session_id = meerkat_core::Session::new().id().clone();
1358 let snapshot = serde_json::to_vec(&session).unwrap();
1359
1360 let err = store
1361 .atomic_apply(
1362 &rid,
1363 Some(SessionDelta {
1364 session_snapshot: snapshot,
1365 }),
1366 make_receipt(RunId::new(), 0),
1367 vec![],
1368 Some(wrong_session_id),
1369 )
1370 .await
1371 .expect_err("mismatched session_store_key should fail");
1372
1373 assert!(matches!(err, RuntimeStoreError::SessionKeyMismatch { .. }));
1374 assert!(store.load_session_snapshot(&rid).await.unwrap().is_none());
1375 }
1376
1377 #[tokio::test]
1378 async fn atomic_apply_persists_machine_owned_receipt() {
1379 let store = InMemoryRuntimeStore::new();
1380 let rid = LogicalRuntimeId::new("test");
1381 let run_id = RunId::new();
1382 let input_id = InputId::new();
1383 let session = meerkat_core::Session::new();
1384 let snapshot = serde_json::to_vec(&session).unwrap();
1385 let receipt = RunBoundaryReceipt {
1386 run_id: run_id.clone(),
1387 boundary: RunApplyBoundary::Immediate,
1388 contributing_input_ids: vec![input_id.clone()],
1389 conversation_digest: Some("machine-owned-digest".to_string()),
1390 message_count: 42,
1391 sequence: 7,
1392 };
1393
1394 store
1395 .atomic_apply(
1396 &rid,
1397 Some(SessionDelta {
1398 session_snapshot: snapshot,
1399 }),
1400 receipt.clone(),
1401 vec![persistable(StoredInputState::new_accepted(input_id))],
1402 None,
1403 )
1404 .await
1405 .unwrap();
1406
1407 assert_eq!(receipt.run_id, run_id);
1408 assert!(receipt.conversation_digest.is_some());
1409 let loaded = store
1410 .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1411 .await
1412 .unwrap();
1413 assert!(loaded.is_some(), "receipt should be persisted");
1414 let Some(loaded) = loaded else {
1415 unreachable!("asserted above");
1416 };
1417 assert_eq!(loaded, receipt);
1418 }
1419
1420 #[tokio::test]
1421 async fn multiple_runtimes_isolated() {
1422 let store = InMemoryRuntimeStore::new();
1423 let rid1 = LogicalRuntimeId::new("runtime-1");
1424 let rid2 = LogicalRuntimeId::new("runtime-2");
1425
1426 store
1427 .persist_input_state(
1428 &rid1,
1429 &persistable(StoredInputState::new_accepted(InputId::new())),
1430 )
1431 .await
1432 .unwrap();
1433 store
1434 .persist_input_state(
1435 &rid2,
1436 &persistable(StoredInputState::new_accepted(InputId::new())),
1437 )
1438 .await
1439 .unwrap();
1440 store
1441 .persist_input_state(
1442 &rid2,
1443 &persistable(StoredInputState::new_accepted(InputId::new())),
1444 )
1445 .await
1446 .unwrap();
1447
1448 let s1 = store.load_input_states(&rid1).await.unwrap();
1449 let s2 = store.load_input_states(&rid2).await.unwrap();
1450 assert_eq!(s1.len(), 1);
1451 assert_eq!(s2.len(), 2);
1452 }
1453
1454 #[tokio::test]
1455 async fn load_session_snapshot_roundtrip() {
1456 let store = InMemoryRuntimeStore::new();
1457 let rid = LogicalRuntimeId::new("runtime");
1458 let snapshot = serde_json::to_vec(&meerkat_core::Session::new()).unwrap();
1459
1460 store
1461 .atomic_apply(
1462 &rid,
1463 Some(SessionDelta {
1464 session_snapshot: snapshot.clone(),
1465 }),
1466 make_receipt(RunId::new(), 0),
1467 vec![],
1468 None,
1469 )
1470 .await
1471 .unwrap();
1472
1473 let loaded = store.load_session_snapshot(&rid).await.unwrap();
1474 assert_eq!(loaded, Some(snapshot));
1475 }
1476
1477 #[tokio::test]
1478 async fn commit_session_snapshot_rejects_stale_runtime_parent() {
1479 let store = InMemoryRuntimeStore::new();
1480 let rid = LogicalRuntimeId::new("runtime-stale-parent");
1481 let accepted = session_with_user("accepted runtime turn");
1482 let mut stale = meerkat_core::Session::with_id(accepted.id().clone());
1483 stale.push(meerkat_core::types::Message::User(
1484 meerkat_core::types::UserMessage::text("stale runtime turn".to_string()),
1485 ));
1486 let accepted_snapshot = serde_json::to_vec(&accepted).unwrap();
1487
1488 store
1489 .commit_session_snapshot(
1490 &rid,
1491 SessionDelta {
1492 session_snapshot: accepted_snapshot.clone(),
1493 },
1494 )
1495 .await
1496 .unwrap();
1497
1498 let err = store
1499 .commit_session_snapshot(
1500 &rid,
1501 SessionDelta {
1502 session_snapshot: serde_json::to_vec(&stale).unwrap(),
1503 },
1504 )
1505 .await
1506 .expect_err("stale non-continuation must not overwrite runtime snapshot");
1507
1508 assert!(matches!(err, RuntimeStoreError::WriteFailed(_)));
1509 assert_eq!(
1510 store.load_session_snapshot(&rid).await.unwrap(),
1511 Some(accepted_snapshot)
1512 );
1513 }
1514
1515 #[tokio::test]
1516 async fn atomic_apply_keeps_current_snapshot_when_incoming_is_superseded() {
1517 let store = InMemoryRuntimeStore::new();
1518 let rid = LogicalRuntimeId::new("runtime-superseded-terminal");
1519 let incoming = session_with_user("turn input");
1520 let mut current = incoming.clone();
1521 current.push(meerkat_core::types::Message::BlockAssistant(
1522 meerkat_core::types::BlockAssistantMessage {
1523 blocks: vec![meerkat_core::types::AssistantBlock::Text {
1524 text: "peer response already applied".to_string(),
1525 meta: None,
1526 }],
1527 stop_reason: meerkat_core::types::StopReason::EndTurn,
1528 identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1529 created_at: meerkat_core::types::message_timestamp_now(),
1530 },
1531 ));
1532 let current_snapshot = serde_json::to_vec(¤t).unwrap();
1533 let receipt = make_receipt(RunId::new(), 11);
1534
1535 store
1536 .commit_session_snapshot(
1537 &rid,
1538 SessionDelta {
1539 session_snapshot: current_snapshot.clone(),
1540 },
1541 )
1542 .await
1543 .unwrap();
1544
1545 store
1546 .atomic_apply(
1547 &rid,
1548 Some(SessionDelta {
1549 session_snapshot: serde_json::to_vec(&incoming).unwrap(),
1550 }),
1551 receipt.clone(),
1552 vec![],
1553 Some(incoming.id().clone()),
1554 )
1555 .await
1556 .unwrap();
1557
1558 assert_eq!(
1559 store.load_session_snapshot(&rid).await.unwrap(),
1560 Some(current_snapshot)
1561 );
1562 assert_eq!(
1566 store
1567 .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1568 .await
1569 .unwrap(),
1570 None
1571 );
1572 }
1573
1574 #[tokio::test]
1575 async fn atomic_apply_skips_inputs_when_session_snapshot_superseded() {
1576 let store = InMemoryRuntimeStore::new();
1577 let rid = LogicalRuntimeId::new("runtime-superseded-inputs");
1578 let incoming = session_with_user("turn input");
1579 let mut current = incoming.clone();
1580 current.push(meerkat_core::types::Message::BlockAssistant(
1581 meerkat_core::types::BlockAssistantMessage {
1582 blocks: vec![meerkat_core::types::AssistantBlock::Text {
1583 text: "peer response already applied".to_string(),
1584 meta: None,
1585 }],
1586 stop_reason: meerkat_core::types::StopReason::EndTurn,
1587 identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1588 created_at: meerkat_core::types::message_timestamp_now(),
1589 },
1590 ));
1591 let current_snapshot = serde_json::to_vec(¤t).unwrap();
1592 let receipt = make_receipt(RunId::new(), 21);
1593 let input_id = InputId::new();
1594 let bundle = StoredInputState::new_accepted(input_id.clone());
1595
1596 store
1597 .commit_session_snapshot(
1598 &rid,
1599 SessionDelta {
1600 session_snapshot: current_snapshot.clone(),
1601 },
1602 )
1603 .await
1604 .unwrap();
1605
1606 store
1607 .atomic_apply(
1608 &rid,
1609 Some(SessionDelta {
1610 session_snapshot: serde_json::to_vec(&incoming).unwrap(),
1611 }),
1612 receipt.clone(),
1613 vec![persistable(bundle)],
1614 Some(incoming.id().clone()),
1615 )
1616 .await
1617 .unwrap();
1618
1619 assert_eq!(
1621 store.load_session_snapshot(&rid).await.unwrap(),
1622 Some(current_snapshot)
1623 );
1624 assert_eq!(
1625 store
1626 .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1627 .await
1628 .unwrap(),
1629 None
1630 );
1631 assert!(store.load_input_states(&rid).await.unwrap().is_empty());
1632 }
1633
1634 #[tokio::test]
1635 async fn atomic_apply_allows_first_generated_snapshot_after_placeholder() {
1636 let store = InMemoryRuntimeStore::new();
1637 let rid = LogicalRuntimeId::new("runtime-placeholder");
1638 let mut placeholder = meerkat_core::Session::new();
1639 placeholder.set_system_prompt("base system".to_string());
1640 let mut incoming = meerkat_core::Session::with_id(placeholder.id().clone());
1641 incoming.set_system_prompt("base system".to_string());
1642 incoming.push(meerkat_core::types::Message::User(
1643 meerkat_core::types::UserMessage::text("verbose first turn".to_string()),
1644 ));
1645 let parent_revision = incoming.transcript_revision().unwrap();
1646 incoming
1647 .commit_transcript_rewrite(
1648 meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
1649 vec![meerkat_core::types::Message::User(
1650 meerkat_core::types::UserMessage::compaction_summary(
1651 "[Context compacted] first turn",
1652 ),
1653 )],
1654 meerkat_core::TranscriptRewriteReason::new("compaction"),
1655 Some("meerkat-core".to_string()),
1656 Some(parent_revision),
1657 )
1658 .unwrap();
1659 let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
1660 let receipt = make_receipt(RunId::new(), 12);
1661
1662 store
1663 .commit_session_snapshot(
1664 &rid,
1665 SessionDelta {
1666 session_snapshot: serde_json::to_vec(&placeholder).unwrap(),
1667 },
1668 )
1669 .await
1670 .unwrap();
1671
1672 store
1673 .atomic_apply(
1674 &rid,
1675 Some(SessionDelta {
1676 session_snapshot: incoming_snapshot.clone(),
1677 }),
1678 receipt.clone(),
1679 vec![],
1680 Some(incoming.id().clone()),
1681 )
1682 .await
1683 .unwrap();
1684
1685 assert_eq!(
1686 store.load_session_snapshot(&rid).await.unwrap(),
1687 Some(incoming_snapshot)
1688 );
1689 assert_eq!(
1690 store
1691 .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1692 .await
1693 .unwrap(),
1694 Some(receipt)
1695 );
1696 }
1697
1698 #[tokio::test]
1699 async fn atomic_apply_allows_generated_compaction_before_retained_tail() {
1700 let store = InMemoryRuntimeStore::new();
1701 let rid = LogicalRuntimeId::new("runtime-compaction-tail");
1702 let mut previous = meerkat_core::Session::new();
1703 previous.set_system_prompt("runtime system before context refresh".to_string());
1704 previous.push(meerkat_core::types::Message::User(
1705 meerkat_core::types::UserMessage::text("Turn 1 request".to_string()),
1706 ));
1707 previous.push(meerkat_core::types::Message::BlockAssistant(
1708 meerkat_core::types::BlockAssistantMessage {
1709 blocks: vec![meerkat_core::types::AssistantBlock::Text {
1710 text: "Turn 1 answer".to_string(),
1711 meta: None,
1712 }],
1713 stop_reason: meerkat_core::types::StopReason::EndTurn,
1714 identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1715 created_at: meerkat_core::types::message_timestamp_now(),
1716 },
1717 ));
1718
1719 let mut incoming = meerkat_core::Session::with_id(previous.id().clone());
1720 incoming.set_system_prompt("runtime system after context refresh".to_string());
1721 incoming.push(meerkat_core::types::Message::User(
1722 meerkat_core::types::UserMessage::text(
1723 "Verbose context that will be compacted".to_string(),
1724 ),
1725 ));
1726 for message in previous.messages()[1..].iter().cloned() {
1727 incoming.push(message);
1728 }
1729 incoming.push(meerkat_core::types::Message::BlockAssistant(
1730 meerkat_core::types::BlockAssistantMessage {
1731 blocks: vec![meerkat_core::types::AssistantBlock::Text {
1732 text: "Turn 2 generated answer".to_string(),
1733 meta: None,
1734 }],
1735 stop_reason: meerkat_core::types::StopReason::EndTurn,
1736 identity: meerkat_core::types::TranscriptMessageIdentity::default(),
1737 created_at: meerkat_core::types::message_timestamp_now(),
1738 },
1739 ));
1740 let parent_revision = incoming.transcript_revision().unwrap();
1741 incoming
1742 .commit_transcript_rewrite(
1743 meerkat_core::TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
1744 vec![meerkat_core::types::Message::User(
1745 meerkat_core::types::UserMessage::compaction_summary(
1746 "[Context compacted] Earlier runtime context".to_string(),
1747 ),
1748 )],
1749 meerkat_core::TranscriptRewriteReason::new("compaction"),
1750 Some("meerkat-core".to_string()),
1751 Some(parent_revision),
1752 )
1753 .unwrap();
1754 let incoming_snapshot = serde_json::to_vec(&incoming).unwrap();
1755 let receipt = make_receipt(RunId::new(), 13);
1756
1757 store
1758 .commit_session_snapshot(
1759 &rid,
1760 SessionDelta {
1761 session_snapshot: serde_json::to_vec(&previous).unwrap(),
1762 },
1763 )
1764 .await
1765 .unwrap();
1766
1767 store
1768 .atomic_apply(
1769 &rid,
1770 Some(SessionDelta {
1771 session_snapshot: incoming_snapshot.clone(),
1772 }),
1773 receipt.clone(),
1774 vec![],
1775 Some(incoming.id().clone()),
1776 )
1777 .await
1778 .unwrap();
1779
1780 assert_eq!(
1781 store.load_session_snapshot(&rid).await.unwrap(),
1782 Some(incoming_snapshot)
1783 );
1784 assert_eq!(
1785 store
1786 .load_boundary_receipt(&rid, &receipt.run_id, receipt.sequence)
1787 .await
1788 .unwrap(),
1789 Some(receipt)
1790 );
1791 }
1792
1793 #[tokio::test]
1794 async fn commit_machine_lifecycle_persists_binding_facts() {
1795 use crate::runtime_state::RuntimeState;
1796
1797 let store = InMemoryRuntimeStore::new();
1798 let rid = LogicalRuntimeId::new("runtime-binding");
1799 let binding = MachineLifecycleBindingFacts::new(
1800 Some("rt:session:abc".to_string()),
1801 Some(7),
1802 Some(3),
1803 Some("epoch-1".to_string()),
1804 );
1805
1806 store
1807 .commit_machine_lifecycle(
1808 &rid,
1809 MachineLifecycleCommit::new_with_binding(
1810 RuntimeState::Retired,
1811 binding.clone(),
1812 crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1813 ),
1814 &[],
1815 )
1816 .await
1817 .unwrap();
1818
1819 let lifecycle = crate::store::load_machine_lifecycle(&store, &rid)
1820 .await
1821 .unwrap()
1822 .expect("machine lifecycle snapshot");
1823 assert_eq!(lifecycle.runtime_state(), RuntimeState::Retired);
1824 assert_eq!(lifecycle.binding(), &binding);
1825 assert_eq!(
1826 crate::store::load_runtime_state(&store, &rid)
1827 .await
1828 .unwrap(),
1829 Some(RuntimeState::Retired)
1830 );
1831 }
1832
1833 #[tokio::test]
1834 async fn unregister_finalization_atomically_retires_ops_epoch_and_is_idempotent() {
1835 let store = InMemoryRuntimeStore::new();
1836 let reopened = store.clone();
1837 let runtime_id = LogicalRuntimeId::new("runtime-unregister-finalization");
1838 let stale_ops = crate::ops_lifecycle::RuntimeOpsLifecycleRegistry::new()
1839 .capture_persistence_snapshot(
1840 meerkat_core::RuntimeEpochId::new(),
1841 &meerkat_core::EpochCursorState::new(),
1842 )
1843 .unwrap();
1844 store
1845 .persist_ops_lifecycle(&runtime_id, &stale_ops)
1846 .await
1847 .unwrap();
1848 let retired_ops_epoch = stale_ops.epoch_id.clone();
1849
1850 for _ in 0..2 {
1851 store
1852 .commit_unregister_finalization(
1853 &runtime_id,
1854 MachineLifecycleCommit::new_with_binding(
1855 RuntimeState::Stopped,
1856 MachineLifecycleBindingFacts::new(None, None, None, None),
1857 crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1858 )
1859 .for_unregister_finalization(retired_ops_epoch.clone()),
1860 &[],
1861 )
1862 .await
1863 .unwrap();
1864 }
1865
1866 assert_eq!(
1867 crate::store::load_runtime_state(&reopened, &runtime_id)
1868 .await
1869 .unwrap(),
1870 Some(RuntimeState::Stopped)
1871 );
1872 assert!(
1873 reopened
1874 .load_ops_lifecycle(&runtime_id)
1875 .await
1876 .unwrap()
1877 .is_none(),
1878 "the same critical section that publishes terminal lifecycle must remove the ops epoch"
1879 );
1880 let late_error = reopened
1881 .persist_ops_lifecycle(&runtime_id, &stale_ops)
1882 .await
1883 .expect_err("a detached callback must not resurrect its retired ops epoch");
1884 assert!(matches!(
1885 late_error,
1886 RuntimeStoreError::OpsLifecycleEpochRetired { epoch_id, .. }
1887 if epoch_id == retired_ops_epoch
1888 ));
1889 assert!(
1890 reopened
1891 .load_ops_lifecycle(&runtime_id)
1892 .await
1893 .unwrap()
1894 .is_none()
1895 );
1896 }
1897
1898 #[tokio::test]
1899 async fn delayed_old_epoch_finalizer_cannot_delete_or_overwrite_new_ops_epoch() {
1900 let store = InMemoryRuntimeStore::new();
1901 let runtime_id = LogicalRuntimeId::new("runtime-old-finalizer-new-epoch");
1902 let registry = crate::ops_lifecycle::RuntimeOpsLifecycleRegistry::new();
1903 let old_ops = registry
1904 .capture_persistence_snapshot(
1905 meerkat_core::RuntimeEpochId::new(),
1906 &meerkat_core::EpochCursorState::new(),
1907 )
1908 .unwrap();
1909 let new_ops = registry
1910 .capture_persistence_snapshot(
1911 meerkat_core::RuntimeEpochId::new(),
1912 &meerkat_core::EpochCursorState::new(),
1913 )
1914 .unwrap();
1915 store
1916 .persist_ops_lifecycle(&runtime_id, &old_ops)
1917 .await
1918 .unwrap();
1919 store
1920 .persist_ops_lifecycle(&runtime_id, &new_ops)
1921 .await
1922 .unwrap();
1923
1924 store
1925 .commit_unregister_finalization(
1926 &runtime_id,
1927 MachineLifecycleCommit::new_with_binding(
1928 RuntimeState::Stopped,
1929 MachineLifecycleBindingFacts::new(None, None, None, None),
1930 crate::store::SupervisorAuthoritySnapshot::UnboundNoReceipt,
1931 )
1932 .for_unregister_finalization(old_ops.epoch_id.clone()),
1933 &[],
1934 )
1935 .await
1936 .unwrap();
1937
1938 assert_eq!(
1939 store
1940 .load_ops_lifecycle(&runtime_id)
1941 .await
1942 .unwrap()
1943 .expect("new epoch row must survive delayed old finalization")
1944 .epoch_id,
1945 new_ops.epoch_id
1946 );
1947 assert!(matches!(
1948 store
1949 .persist_ops_lifecycle(&runtime_id, &old_ops)
1950 .await
1951 .expect_err("retired old epoch stays fenced"),
1952 RuntimeStoreError::OpsLifecycleEpochRetired { .. }
1953 ));
1954 store
1955 .persist_ops_lifecycle(&runtime_id, &new_ops)
1956 .await
1957 .unwrap();
1958 }
1959
1960 #[tokio::test]
1961 async fn clear_session_snapshot_if_current_sets_quarantine_marker_cleared_on_write() {
1962 let store = InMemoryRuntimeStore::new();
1963 let rid = LogicalRuntimeId::new("runtime-quarantine");
1964 let rejected = serde_json::to_vec(&session_with_user("rejected")).unwrap();
1965
1966 assert!(!store.is_runtime_projection_quarantined(&rid).await.unwrap());
1967 store
1968 .commit_session_snapshot(
1969 &rid,
1970 SessionDelta {
1971 session_snapshot: rejected.clone(),
1972 },
1973 )
1974 .await
1975 .unwrap();
1976 assert!(
1977 store
1978 .clear_session_snapshot_if_current(&rid, &rejected)
1979 .await
1980 .unwrap()
1981 );
1982 assert!(
1983 store.is_runtime_projection_quarantined(&rid).await.unwrap(),
1984 "clearing the rejected snapshot must record the in-memory quarantine marker"
1985 );
1986
1987 store
1989 .commit_session_snapshot(
1990 &rid,
1991 SessionDelta {
1992 session_snapshot: serde_json::to_vec(&session_with_user("revived")).unwrap(),
1993 },
1994 )
1995 .await
1996 .unwrap();
1997 assert!(
1998 !store.is_runtime_projection_quarantined(&rid).await.unwrap(),
1999 "a live snapshot write must clear the in-memory quarantine marker"
2000 );
2001 }
2002}