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