1use std::collections::{HashMap, HashSet};
9use std::sync::atomic::{AtomicU64, Ordering};
10use std::sync::{Arc, Mutex};
11
12use async_trait::async_trait;
13
14use super::contracts::{AgentCustomizer, RosterProvider, TopologyProvider};
15use super::types::{
16 AgentAddressability, AgentBuildContext, AgentBuildDraft, AgentIdentity, ContinuityStoreError,
17 CustomizerError, DurableAgentSpec, ManagedPeerEdge, RosterContext, RosterError,
18 TopologyContext, TopologyError,
19};
20use crate::mob_handle_runtime::{SessionCreatedContext, SessionHook};
21use crate::types::AgentDiscoverySpec;
22use crate::unified_runtime::edge_types::{Discovery, EdgeDiscovery};
23
24pub struct DiscoveryRosterAdapter {
40 inner: Box<dyn Discovery>,
41}
42
43impl DiscoveryRosterAdapter {
44 pub fn new(discovery: impl Discovery + 'static) -> Self {
45 Self {
46 inner: Box::new(discovery),
47 }
48 }
49}
50
51pub fn agent_discovery_to_durable(
53 spec: &AgentDiscoverySpec,
54) -> Result<DurableAgentSpec, RosterError> {
55 let identity = AgentIdentity::parse(&spec.meerkat_id)
56 .map_err(|e| RosterError::Io(format!("invalid meerkat_id: {e}")))?;
57 Ok(DurableAgentSpec {
58 identity,
59 profile: meerkat_mob::ProfileName::from(spec.profile.as_str()),
60 addressability: AgentAddressability::Addressable,
61 display_name: None,
62 labels: spec.labels.clone().unwrap_or_default(),
63 context: spec.context.clone(),
64 additional_instructions: spec.additional_instructions.clone(),
65 initial_message: None,
66 runtime_mode_override: None,
67 backend: None,
68 binding: None,
69 })
70}
71
72#[async_trait]
73impl RosterProvider for DiscoveryRosterAdapter {
74 async fn roster(&self, _context: &RosterContext) -> Result<Vec<DurableAgentSpec>, RosterError> {
75 let specs = self.inner.discover(serde_json::Value::Null).await;
76 specs.iter().map(agent_discovery_to_durable).collect()
77 }
78}
79
80pub struct EdgeDiscoveryTopologyAdapter {
89 inner: Box<dyn EdgeDiscovery>,
90}
91
92impl EdgeDiscoveryTopologyAdapter {
93 pub fn new(edge_discovery: impl EdgeDiscovery + 'static) -> Self {
94 Self {
95 inner: Box::new(edge_discovery),
96 }
97 }
98}
99
100#[async_trait]
101impl TopologyProvider for EdgeDiscoveryTopologyAdapter {
102 async fn compute_edges(
103 &self,
104 _target_identities: &[AgentIdentity],
105 context: &TopologyContext,
106 ) -> Result<Vec<ManagedPeerEdge>, TopologyError> {
107 let member_views: Vec<crate::unified_runtime::edge_types::EdgeMemberView> = context
110 .roster
111 .iter()
112 .map(|spec| crate::unified_runtime::edge_types::EdgeMemberView {
113 agent_identity: spec.identity.as_str().to_string(),
114 role: spec.profile.as_str().to_string(),
115 wired_to: std::collections::BTreeSet::new(),
116 labels: spec.labels.clone(),
117 })
118 .collect();
119
120 let desired_edges = self.inner.discover_edges(member_views).await;
121 let mut edges = Vec::with_capacity(desired_edges.len());
122 for edge in &desired_edges {
123 let (a_str, b_str) = edge.endpoints();
124 let a = AgentIdentity::parse(a_str)
125 .map_err(|e| TopologyError::InvalidEdge(format!("endpoint {a_str:?}: {e}")))?;
126 let b = AgentIdentity::parse(b_str)
127 .map_err(|e| TopologyError::InvalidEdge(format!("endpoint {b_str:?}: {e}")))?;
128 let managed = ManagedPeerEdge::new(a, b)
129 .map_err(|e| TopologyError::InvalidEdge(format!("{e}")))?;
130 edges.push(managed);
131 }
132 Ok(edges)
133 }
134}
135
136#[derive(Clone)]
142pub(crate) struct SessionRuntimeState {
143 pub identity: AgentIdentity,
144 pub generation: super::types::ContinuityGeneration,
145 pub fencing_token: super::types::FencingToken,
146 pub checkpoint_version: super::types::CheckpointVersion,
147}
148
149pub struct ContinuitySessionStoreAdapter {
159 store: Arc<dyn super::contracts::ContinuityStore>,
160 versions: Mutex<HashMap<String, AtomicU64>>,
162 session_registry: Mutex<HashMap<String, SessionRuntimeState>>,
164 pending_unregistered: Mutex<HashMap<String, Vec<u8>>>,
167 unregistered_sessions: Mutex<HashSet<String>>,
171 save_guard: tokio::sync::Mutex<()>,
173}
174
175impl ContinuitySessionStoreAdapter {
176 pub fn new(store: Arc<dyn super::contracts::ContinuityStore>) -> Self {
177 Self {
178 store,
179 versions: Mutex::new(HashMap::new()),
180 session_registry: Mutex::new(HashMap::new()),
181 pending_unregistered: Mutex::new(HashMap::new()),
182 unregistered_sessions: Mutex::new(HashSet::new()),
183 save_guard: tokio::sync::Mutex::new(()),
184 }
185 }
186
187 #[allow(dead_code)]
192 pub(crate) async fn register_session(
193 &self,
194 session_id: &meerkat_core::types::SessionId,
195 state: SessionRuntimeState,
196 ) -> Result<super::types::CheckpointVersion, meerkat_store::SessionStoreError> {
197 let _guard = self.save_guard.lock().await;
198 let session_key = session_id.to_string();
199 let checkpoint_version = state.checkpoint_version.get();
200 let previous_registry = {
201 let mut registry = self
202 .session_registry
203 .lock()
204 .unwrap_or_else(std::sync::PoisonError::into_inner);
205 registry.insert(session_key.clone(), state.clone())
206 };
207 self.unregistered_sessions
208 .lock()
209 .unwrap_or_else(std::sync::PoisonError::into_inner)
210 .remove(&session_key);
211
212 let previous_version = {
213 let mut versions = self
214 .versions
215 .lock()
216 .unwrap_or_else(std::sync::PoisonError::into_inner);
217 let counter = versions
218 .entry(session_key)
219 .or_insert_with(|| AtomicU64::new(checkpoint_version));
220 let previous_version = counter.load(Ordering::Relaxed);
221 counter.fetch_max(checkpoint_version, Ordering::Relaxed);
222 previous_version
223 };
224
225 let pending = {
226 let pending = self
227 .pending_unregistered
228 .lock()
229 .unwrap_or_else(std::sync::PoisonError::into_inner);
230 pending.get(&session_id.to_string()).cloned()
231 };
232 let mut effective_checkpoint_version = self.current_version(session_id);
233 if let Some(data) = pending {
234 let flush_result = self.save_registered_snapshot(session_id, data, state).await;
235 match flush_result {
236 Ok(version) => {
237 self.pending_unregistered
238 .lock()
239 .unwrap_or_else(std::sync::PoisonError::into_inner)
240 .remove(&session_id.to_string());
241 effective_checkpoint_version = version;
242 }
243 Err(err) => {
244 self.restore_registration_state(
245 session_id,
246 previous_registry,
247 previous_version,
248 );
249 return Err(err);
250 }
251 }
252 }
253 Ok(effective_checkpoint_version)
254 }
255
256 #[allow(dead_code)]
258 pub(crate) fn update_fencing_token(
259 &self,
260 session_id: &meerkat_core::types::SessionId,
261 token: super::types::FencingToken,
262 ) {
263 let mut registry = self
264 .session_registry
265 .lock()
266 .unwrap_or_else(std::sync::PoisonError::into_inner);
267 if let Some(state) = registry.get_mut(&session_id.to_string()) {
268 state.fencing_token = token;
269 }
270 }
271
272 fn forget_session(&self, session_id: &meerkat_core::types::SessionId) {
273 let key = session_id.to_string();
274 self.session_registry
275 .lock()
276 .unwrap_or_else(std::sync::PoisonError::into_inner)
277 .remove(&key);
278 self.pending_unregistered
279 .lock()
280 .unwrap_or_else(std::sync::PoisonError::into_inner)
281 .remove(&key);
282 self.versions
283 .lock()
284 .unwrap_or_else(std::sync::PoisonError::into_inner)
285 .remove(&key);
286 }
287
288 pub(crate) async fn unregister_session(
289 &self,
290 session_id: &meerkat_core::types::SessionId,
291 ) -> Result<(), ContinuityStoreError> {
292 let _guard = self.save_guard.lock().await;
293 self.forget_session(session_id);
294 self.unregistered_sessions
295 .lock()
296 .unwrap_or_else(std::sync::PoisonError::into_inner)
297 .insert(session_id.to_string());
298 Ok(())
299 }
300
301 fn session_was_unregistered(&self, session_id: &meerkat_core::types::SessionId) -> bool {
302 self.unregistered_sessions
303 .lock()
304 .unwrap_or_else(std::sync::PoisonError::into_inner)
305 .contains(&session_id.to_string())
306 }
307
308 fn next_version(&self, session_id: &str) -> u64 {
310 let mut map = self
311 .versions
312 .lock()
313 .unwrap_or_else(std::sync::PoisonError::into_inner);
314 let counter = map
315 .entry(session_id.to_string())
316 .or_insert_with(|| AtomicU64::new(0));
317 counter.fetch_add(1, Ordering::Relaxed) + 1
318 }
319
320 fn restore_registration_state(
321 &self,
322 session_id: &meerkat_core::types::SessionId,
323 previous_registry: Option<SessionRuntimeState>,
324 previous_version: u64,
325 ) {
326 let key = session_id.to_string();
327 {
328 let mut registry = self
329 .session_registry
330 .lock()
331 .unwrap_or_else(std::sync::PoisonError::into_inner);
332 match previous_registry {
333 Some(state) => {
334 registry.insert(key.clone(), state);
335 }
336 None => {
337 registry.remove(&key);
338 }
339 }
340 }
341 let mut versions = self
342 .versions
343 .lock()
344 .unwrap_or_else(std::sync::PoisonError::into_inner);
345 if previous_version == 0 {
346 versions.remove(&key);
347 } else {
348 versions
349 .entry(key)
350 .or_insert_with(|| AtomicU64::new(previous_version))
351 .store(previous_version, Ordering::Relaxed);
352 }
353 }
354
355 fn current_version(
356 &self,
357 session_id: &meerkat_core::types::SessionId,
358 ) -> super::types::CheckpointVersion {
359 let map = self
360 .versions
361 .lock()
362 .unwrap_or_else(std::sync::PoisonError::into_inner);
363 let version = map
364 .get(&session_id.to_string())
365 .map(|counter| counter.load(Ordering::Relaxed))
366 .unwrap_or(0);
367 super::types::CheckpointVersion::new(version)
368 }
369
370 fn lookup_session(&self, session_id: &str) -> Option<SessionRuntimeState> {
372 let registry = self
373 .session_registry
374 .lock()
375 .unwrap_or_else(std::sync::PoisonError::into_inner);
376 registry.get(session_id).cloned()
377 }
378
379 async fn load_persisted_session(
380 &self,
381 id: &meerkat_core::types::SessionId,
382 ) -> Result<Option<meerkat_core::Session>, meerkat_store::SessionStoreError> {
383 let snapshot = self.store.load_session_snapshot(id).await.map_err(|e| {
384 meerkat_store::SessionStoreError::Internal(format!("continuity load: {e}"))
385 })?;
386 match snapshot {
387 Some(snap) => {
388 let session: meerkat_core::Session = serde_json::from_slice(&snap.data)
389 .map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
390 Ok(Some(session))
391 }
392 None => Ok(None),
393 }
394 }
395
396 async fn load_previous_session_for_save(
397 &self,
398 id: &meerkat_core::types::SessionId,
399 ) -> Result<Option<meerkat_core::Session>, meerkat_store::SessionStoreError> {
400 if let Some(session) = self.load_persisted_session(id).await? {
401 return Ok(Some(session));
402 }
403 let pending = self
404 .pending_unregistered
405 .lock()
406 .unwrap_or_else(std::sync::PoisonError::into_inner)
407 .get(&id.to_string())
408 .cloned();
409 pending
410 .map(|data| {
411 serde_json::from_slice(&data)
412 .map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))
413 })
414 .transpose()
415 }
416
417 async fn save_registered_snapshot(
418 &self,
419 session_id: &meerkat_core::types::SessionId,
420 data: Vec<u8>,
421 state: SessionRuntimeState,
422 ) -> Result<super::types::CheckpointVersion, meerkat_store::SessionStoreError> {
423 let version = self.next_version(&session_id.to_string());
424 let checkpoint_version = super::types::CheckpointVersion::new(version);
425 let snapshot = super::types::SessionSnapshot { data };
426 self.store
427 .save_session_snapshot(
428 &state.identity,
429 session_id,
430 state.generation,
431 checkpoint_version,
432 state.fencing_token,
433 &snapshot,
434 )
435 .await
436 .map_err(|e| {
437 meerkat_store::SessionStoreError::Internal(format!("continuity save: {e}"))
438 })?;
439 Ok(checkpoint_version)
440 }
441}
442
443#[async_trait]
444impl meerkat::SessionStore for ContinuitySessionStoreAdapter {
445 async fn save(
446 &self,
447 session: &meerkat_core::Session,
448 ) -> Result<(), meerkat_store::SessionStoreError> {
449 let _guard = self.save_guard.lock().await;
450 if self.session_was_unregistered(session.id()) {
451 return Err(meerkat_store::SessionStoreError::Internal(format!(
452 "session {} was unregistered from identity runtime state",
453 session.id()
454 )));
455 }
456 let previous = self.load_previous_session_for_save(session.id()).await?;
457 meerkat_core::session_store::append_only_save_guard(session, previous.as_ref())?;
458 let data = serde_json::to_vec(session)
459 .map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
460 let sid_str = session.id().to_string();
461
462 match self.lookup_session(&sid_str) {
464 Some(state) => {
465 self.save_registered_snapshot(session.id(), data, state)
466 .await?;
467 }
468 None => {
469 tracing::warn!(
474 session_id = %sid_str,
475 "ContinuitySessionStoreAdapter: delaying save until runtime state is registered"
476 );
477 let mut pending = self
478 .pending_unregistered
479 .lock()
480 .unwrap_or_else(std::sync::PoisonError::into_inner);
481 pending.insert(sid_str, data);
482 }
483 }
484 Ok(())
485 }
486
487 async fn save_transcript_rewrite(
488 &self,
489 session: &meerkat_core::Session,
490 commit: &meerkat_core::TranscriptRewriteCommit,
491 ) -> Result<(), meerkat_store::SessionStoreError> {
492 let _guard = self.save_guard.lock().await;
493 if self.session_was_unregistered(session.id()) {
494 return Err(meerkat_store::SessionStoreError::Internal(format!(
495 "session {} was unregistered from identity runtime state",
496 session.id()
497 )));
498 }
499 let previous = self.load_previous_session_for_save(session.id()).await?;
500 meerkat_core::session_store::transcript_rewrite_save_guard(
501 session,
502 previous.as_ref(),
503 commit,
504 )?;
505 let data = serde_json::to_vec(session)
506 .map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
507 let sid_str = session.id().to_string();
508
509 match self.lookup_session(&sid_str) {
510 Some(state) => {
511 self.save_registered_snapshot(session.id(), data, state)
512 .await?;
513 }
514 None => {
515 let mut pending = self
516 .pending_unregistered
517 .lock()
518 .unwrap_or_else(std::sync::PoisonError::into_inner);
519 pending.insert(sid_str, data);
520 }
521 }
522 Ok(())
523 }
524
525 async fn save_authoritative_projection(
526 &self,
527 session: &meerkat_core::Session,
528 ) -> Result<(), meerkat_store::SessionStoreError> {
529 let data = serde_json::to_vec(session)
530 .map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
531 let sid_str = session.id().to_string();
532
533 let _guard = self.save_guard.lock().await;
534 match self.lookup_session(&sid_str) {
535 Some(state) => {
536 self.save_registered_snapshot(session.id(), data, state)
537 .await?;
538 }
539 None => {
540 let mut pending = self
541 .pending_unregistered
542 .lock()
543 .unwrap_or_else(std::sync::PoisonError::into_inner);
544 pending.insert(sid_str, data);
545 }
546 }
547 Ok(())
548 }
549
550 async fn save_authoritative_projection_if_current_revision(
551 &self,
552 session: &meerkat_core::Session,
553 expected_current_revision: Option<String>,
554 ) -> Result<(), meerkat_store::SessionStoreError> {
555 let _guard = self.save_guard.lock().await;
556 let previous = self.load_persisted_session(session.id()).await?;
557 meerkat_core::session_store::authoritative_projection_current_revision_guard(
558 session,
559 previous.as_ref(),
560 expected_current_revision.as_deref(),
561 )?;
562 let data = serde_json::to_vec(session)
563 .map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
564 let sid_str = session.id().to_string();
565 match self.lookup_session(&sid_str) {
566 Some(state) => {
567 self.save_registered_snapshot(session.id(), data, state)
568 .await?;
569 Ok(())
570 }
571 None => {
572 let mut pending = self
573 .pending_unregistered
574 .lock()
575 .unwrap_or_else(std::sync::PoisonError::into_inner);
576 pending.insert(sid_str, data);
577 Ok(())
578 }
579 }
580 }
581
582 async fn load(
583 &self,
584 id: &meerkat_core::types::SessionId,
585 ) -> Result<Option<meerkat_core::Session>, meerkat_store::SessionStoreError> {
586 match self.load_persisted_session(id).await? {
587 Some(session) => Ok(Some(session)),
588 None if self.lookup_session(&id.to_string()).is_some() => {
589 Ok(Some(meerkat_core::Session::with_id(id.clone())))
590 }
591 None => Ok(None),
592 }
593 }
594
595 async fn list(
596 &self,
597 _filter: meerkat_store::SessionFilter,
598 ) -> Result<Vec<meerkat_core::SessionMeta>, meerkat_store::SessionStoreError> {
599 Ok(Vec::new())
602 }
603
604 async fn delete(
605 &self,
606 id: &meerkat_core::types::SessionId,
607 ) -> Result<(), meerkat_store::SessionStoreError> {
608 let _guard = self.save_guard.lock().await;
609 let Some(session) = self.load_persisted_session(id).await? else {
610 self.forget_session(id);
611 return Ok(());
612 };
613 let current_revision = meerkat_core::session_store::session_projection_cas_token(&session)?;
614 let deleted = self
615 .store
616 .delete_session_snapshot_if_current_revision(id, ¤t_revision)
617 .await
618 .map_err(|e| {
619 meerkat_store::SessionStoreError::Internal(format!("continuity delete: {e}"))
620 })?;
621 if !deleted {
622 return Err(meerkat_store::SessionStoreError::Internal(format!(
623 "continuity delete did not remove session snapshot {id}"
624 )));
625 }
626 self.forget_session(id);
627 Ok(())
628 }
629
630 async fn delete_if_current_revision(
631 &self,
632 id: &meerkat_core::types::SessionId,
633 expected_current_revision: &str,
634 ) -> Result<bool, meerkat_store::SessionStoreError> {
635 let _guard = self.save_guard.lock().await;
636 let Some(session) = self.load_persisted_session(id).await? else {
637 self.forget_session(id);
638 return Ok(false);
639 };
640 let current_revision = meerkat_core::session_store::session_projection_cas_token(&session)?;
641 if current_revision != expected_current_revision {
642 return Ok(false);
643 }
644 let deleted = self
645 .store
646 .delete_session_snapshot_if_current_revision(id, expected_current_revision)
647 .await
648 .map_err(|e| {
649 meerkat_store::SessionStoreError::Internal(format!(
650 "continuity delete_if_current_revision: {e}"
651 ))
652 })?;
653 if deleted {
654 self.forget_session(id);
655 }
656 Ok(deleted)
657 }
658}
659
660pub struct SessionHookCustomizerAdapter {
670 hook: Arc<dyn SessionHook>,
671}
672
673impl SessionHookCustomizerAdapter {
674 pub fn new(hook: Arc<dyn SessionHook>) -> Self {
675 Self { hook }
676 }
677}
678
679#[async_trait]
680impl AgentCustomizer for SessionHookCustomizerAdapter {
681 async fn customize_build(
682 &self,
683 _context: &AgentBuildContext,
684 spec: &DurableAgentSpec,
685 draft: &mut AgentBuildDraft,
686 ) -> Result<(), CustomizerError> {
687 let mut req = meerkat_core::service::CreateSessionRequest {
689 model: draft.model.clone().unwrap_or_default(),
690 prompt: meerkat_core::ContentInput::Text(String::new()),
691 system_prompt: match draft.system_prompt.clone() {
694 Some(prompt) => meerkat_core::config::SystemPromptOverride::Set(prompt),
695 None => meerkat_core::config::SystemPromptOverride::Inherit,
696 },
697 max_tokens: None,
698 event_tx: None,
699 initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
700 build: None,
701 labels: if draft.labels.is_empty() {
702 None
703 } else {
704 Some(draft.labels.clone())
705 },
706 deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
707 injected_context: Vec::new(),
711 };
712
713 let prompt_before = req.prompt.clone();
717 let max_tokens_before = req.max_tokens;
718 let event_tx_was_some = req.event_tx.is_some();
719 let initial_turn_before = req.initial_turn;
720 let build_before_is_none = req.build.is_none();
721
722 self.hook
723 .before_create(&mut req)
724 .await
725 .map_err(|e| CustomizerError::BuildFailed(format!("session hook: {e}")))?;
726
727 let mut unsupported_mutations: Vec<&str> = Vec::new();
730
731 if req.prompt != prompt_before {
732 unsupported_mutations.push("prompt");
733 }
734 if req.max_tokens != max_tokens_before {
735 unsupported_mutations.push("max_tokens");
736 }
737 if req.event_tx.is_some() != event_tx_was_some {
738 unsupported_mutations.push("event_tx");
739 }
740 if req.initial_turn != initial_turn_before {
741 unsupported_mutations.push("initial_turn");
742 }
743 if let Some(ref build) = req.build {
748 if build_before_is_none {
749 unsupported_mutations.push("build");
751 if build.resume_session.is_some() {
752 unsupported_mutations.push("build.resume_session");
753 }
754 } else if build.resume_session.is_some() {
755 unsupported_mutations.push("build.resume_session");
757 }
758 }
759
760 if !unsupported_mutations.is_empty() {
761 tracing::warn!(
762 identity = %spec.identity,
763 fields = ?unsupported_mutations,
764 "SessionHook mutated unsupported CreateSessionRequest fields — \
765 these mutations are NOT applied in the identity-first model. \
766 Migrate to AgentCustomizer."
767 );
768 }
769
770 if !req.model.is_empty() {
778 draft.model = Some(req.model);
779 }
780 draft.system_prompt = req.system_prompt.as_set_prompt().map(ToString::to_string);
781 draft.labels = req.labels.unwrap_or_default();
782
783 Ok(())
784 }
785
786 async fn after_create(
787 &self,
788 _identity: &AgentIdentity,
789 session_id: &meerkat_core::types::SessionId,
790 context: &SessionCreatedContext,
791 ) -> Result<(), CustomizerError> {
792 self.hook.after_create(session_id, context).await;
793 Ok(())
794 }
795}
796
797#[cfg(test)]
798#[allow(clippy::expect_used, clippy::panic)]
799mod tests {
800 use std::sync::Arc;
801 use std::sync::atomic::{AtomicBool, Ordering as AtomicOrdering};
802
803 use serde_json::json;
804
805 use super::super::contracts::ContinuityStore;
806 use super::super::local_store::LocalContinuityStore;
807 use super::super::types::{
808 AgentIdentity, AgentRuntimeId, CheckpointVersion, ContinuityGeneration, ContinuityRecord,
809 ContinuityResolveState, ContinuityStoreError, FencingToken, SessionSnapshot,
810 };
811 use super::*;
812
813 struct FailSaveContinuityStore {
814 inner: Arc<LocalContinuityStore>,
815 fail_save: AtomicBool,
816 }
817
818 impl FailSaveContinuityStore {
819 fn new(inner: Arc<LocalContinuityStore>) -> Self {
820 Self {
821 inner,
822 fail_save: AtomicBool::new(false),
823 }
824 }
825
826 fn fail_saves(&self, fail: bool) {
827 self.fail_save.store(fail, AtomicOrdering::SeqCst);
828 }
829 }
830
831 #[async_trait]
832 impl ContinuityStore for FailSaveContinuityStore {
833 async fn resolve_many(
834 &self,
835 identities: &[AgentIdentity],
836 ) -> Result<
837 std::collections::BTreeMap<AgentIdentity, ContinuityResolveState>,
838 ContinuityStoreError,
839 > {
840 self.inner.resolve_many(identities).await
841 }
842
843 async fn load_session_snapshot(
844 &self,
845 session_id: &meerkat_core::types::SessionId,
846 ) -> Result<Option<SessionSnapshot>, ContinuityStoreError> {
847 self.inner.load_session_snapshot(session_id).await
848 }
849
850 async fn delete_session_snapshot_if_current_revision(
851 &self,
852 session_id: &meerkat_core::types::SessionId,
853 expected_current_revision: &str,
854 ) -> Result<bool, ContinuityStoreError> {
855 self.inner
856 .delete_session_snapshot_if_current_revision(session_id, expected_current_revision)
857 .await
858 }
859
860 async fn save_session_snapshot(
861 &self,
862 identity: &AgentIdentity,
863 session_id: &meerkat_core::types::SessionId,
864 generation: ContinuityGeneration,
865 version: CheckpointVersion,
866 fencing_token: FencingToken,
867 snapshot: &SessionSnapshot,
868 ) -> Result<(), ContinuityStoreError> {
869 if self.fail_save.load(AtomicOrdering::SeqCst) {
870 return Err(ContinuityStoreError::Io("forced save failure".to_string()));
871 }
872 self.inner
873 .save_session_snapshot(
874 identity,
875 session_id,
876 generation,
877 version,
878 fencing_token,
879 snapshot,
880 )
881 .await
882 }
883
884 async fn upsert_continuity_record(
885 &self,
886 record: &ContinuityRecord,
887 fencing_token: FencingToken,
888 ) -> Result<(), ContinuityStoreError> {
889 self.inner
890 .upsert_continuity_record(record, fencing_token)
891 .await
892 }
893
894 async fn delete_continuity_record(
895 &self,
896 identity: &AgentIdentity,
897 fencing_token: FencingToken,
898 ) -> Result<(), ContinuityStoreError> {
899 self.inner
900 .delete_continuity_record(identity, fencing_token)
901 .await
902 }
903 }
904
905 #[tokio::test]
906 async fn continuity_session_store_adapter_seeds_registered_checkpoint_version() {
907 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
908 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
909 let session = meerkat_core::Session::new();
910 let identity = AgentIdentity::parse("agent:restored").expect("identity");
911 let record = ContinuityRecord {
912 identity: identity.clone(),
913 agent_runtime_id: AgentRuntimeId::parse("rt:agent:restored:0").expect("runtime id"),
914 session_id: session.id().clone(),
915 generation: ContinuityGeneration::new(2),
916 checkpoint_version: CheckpointVersion::new(5),
917 };
918 let fencing_token = FencingToken::new(9);
919 store
920 .upsert_continuity_record(&record, fencing_token)
921 .await
922 .expect("seed record");
923
924 adapter
925 .register_session(
926 session.id(),
927 SessionRuntimeState {
928 identity: identity.clone(),
929 generation: record.generation,
930 fencing_token,
931 checkpoint_version: record.checkpoint_version,
932 },
933 )
934 .await
935 .expect("register");
936
937 meerkat::SessionStore::save(&adapter, &session)
938 .await
939 .expect("save should advance from restored checkpoint");
940 let effective_version = adapter
941 .register_session(
942 session.id(),
943 SessionRuntimeState {
944 identity: identity.clone(),
945 generation: record.generation,
946 fencing_token,
947 checkpoint_version: record.checkpoint_version,
948 },
949 )
950 .await
951 .expect("post-save register should report advanced version");
952 assert_eq!(effective_version, CheckpointVersion::new(6));
953 let resolved = store
954 .resolve_many(std::slice::from_ref(&identity))
955 .await
956 .expect("resolve");
957 let ContinuityResolveState::Ready { record } = resolved.get(&identity).expect("record")
958 else {
959 panic!("expected ready record");
960 };
961 assert_eq!(record.checkpoint_version, CheckpointVersion::new(6));
962 }
963
964 #[tokio::test]
965 async fn continuity_session_store_adapter_flushes_pending_save_under_registered_identity() {
966 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
967 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
968 let session = meerkat_core::Session::new();
969 let identity = AgentIdentity::parse("agent:fresh").expect("identity");
970 let record = ContinuityRecord {
971 identity: identity.clone(),
972 agent_runtime_id: AgentRuntimeId::parse("rt:agent:fresh:0").expect("runtime id"),
973 session_id: session.id().clone(),
974 generation: ContinuityGeneration::new(0),
975 checkpoint_version: CheckpointVersion::new(0),
976 };
977 let fencing_token = FencingToken::new(3);
978 store
979 .upsert_continuity_record(&record, fencing_token)
980 .await
981 .expect("seed record");
982
983 meerkat::SessionStore::save(&adapter, &session)
984 .await
985 .expect("unregistered save should be delayed, not written under fallback identity");
986 assert!(
987 store
988 .load_session_snapshot(session.id())
989 .await
990 .expect("load before register")
991 .is_none(),
992 "unregistered save must not be visible in continuity store"
993 );
994
995 adapter
996 .register_session(
997 session.id(),
998 SessionRuntimeState {
999 identity: identity.clone(),
1000 generation: record.generation,
1001 fencing_token,
1002 checkpoint_version: record.checkpoint_version,
1003 },
1004 )
1005 .await
1006 .expect("register flushes pending");
1007
1008 assert!(
1009 store
1010 .load_session_snapshot(session.id())
1011 .await
1012 .expect("load after register")
1013 .is_some(),
1014 "pending save should flush under the registered identity"
1015 );
1016 let resolved = store
1017 .resolve_many(std::slice::from_ref(&identity))
1018 .await
1019 .expect("resolve");
1020 let ContinuityResolveState::Ready { record } = resolved.get(&identity).expect("record")
1021 else {
1022 panic!("expected ready record");
1023 };
1024 assert_eq!(record.checkpoint_version, CheckpointVersion::new(1));
1025 }
1026
1027 #[tokio::test]
1028 async fn continuity_session_store_adapter_rejects_saves_after_unregister() {
1029 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1030 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1031 let session = meerkat_core::Session::new();
1032 let identity = AgentIdentity::parse("agent:retired").expect("identity");
1033 let record = ContinuityRecord {
1034 identity: identity.clone(),
1035 agent_runtime_id: AgentRuntimeId::parse("rt:agent:retired:0").expect("runtime id"),
1036 session_id: session.id().clone(),
1037 generation: ContinuityGeneration::new(0),
1038 checkpoint_version: CheckpointVersion::new(0),
1039 };
1040 let fencing_token = FencingToken::new(9);
1041 store
1042 .upsert_continuity_record(&record, fencing_token)
1043 .await
1044 .expect("seed record");
1045 adapter
1046 .register_session(
1047 session.id(),
1048 SessionRuntimeState {
1049 identity: identity.clone(),
1050 generation: record.generation,
1051 fencing_token,
1052 checkpoint_version: record.checkpoint_version,
1053 },
1054 )
1055 .await
1056 .expect("register");
1057
1058 adapter
1059 .unregister_session(session.id())
1060 .await
1061 .expect("unregister");
1062 let err = meerkat::SessionStore::save(&adapter, &session)
1063 .await
1064 .expect_err("post-unregister save must fail closed");
1065 assert!(
1066 err.to_string().contains("was unregistered"),
1067 "unexpected error: {err}"
1068 );
1069 assert!(
1070 store
1071 .load_session_snapshot(session.id())
1072 .await
1073 .expect("load")
1074 .is_none(),
1075 "post-unregister save must not be queued as pending"
1076 );
1077
1078 adapter
1079 .register_session(
1080 session.id(),
1081 SessionRuntimeState {
1082 identity,
1083 generation: record.generation,
1084 fencing_token,
1085 checkpoint_version: record.checkpoint_version,
1086 },
1087 )
1088 .await
1089 .expect("registering the same id later should not flush stale pending data");
1090 assert!(
1091 store
1092 .load_session_snapshot(session.id())
1093 .await
1094 .expect("load after re-register")
1095 .is_none(),
1096 "stale post-unregister save must not flush on a later registration"
1097 );
1098 }
1099
1100 #[tokio::test]
1101 async fn continuity_session_store_adapter_register_keeps_pending_snapshot_on_flush_failure() {
1102 let inner = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1103 let fail_store = Arc::new(FailSaveContinuityStore::new(inner.clone()));
1104 let adapter = ContinuitySessionStoreAdapter::new(fail_store.clone());
1105 let mut session = meerkat_core::Session::new();
1106 session.set_metadata("pending", json!(true));
1107 let identity = AgentIdentity::parse("agent:pending-fail").expect("identity");
1108 let record = ContinuityRecord {
1109 identity: identity.clone(),
1110 agent_runtime_id: AgentRuntimeId::parse("rt:agent:pending-fail:0").expect("runtime id"),
1111 session_id: session.id().clone(),
1112 generation: ContinuityGeneration::new(0),
1113 checkpoint_version: CheckpointVersion::new(0),
1114 };
1115 let fencing_token = FencingToken::new(14);
1116 inner
1117 .upsert_continuity_record(&record, fencing_token)
1118 .await
1119 .expect("seed record");
1120
1121 meerkat::SessionStore::save(&adapter, &session)
1122 .await
1123 .expect("pending save");
1124 fail_store.fail_saves(true);
1125 adapter
1126 .register_session(
1127 session.id(),
1128 SessionRuntimeState {
1129 identity: identity.clone(),
1130 generation: record.generation,
1131 fencing_token,
1132 checkpoint_version: record.checkpoint_version,
1133 },
1134 )
1135 .await
1136 .expect_err("forced pending flush failure");
1137 assert!(
1138 meerkat::SessionStore::load(&adapter, session.id())
1139 .await
1140 .expect("load after failed register")
1141 .is_none(),
1142 "failed register must not leave a synthetic registered session"
1143 );
1144
1145 fail_store.fail_saves(false);
1146 adapter
1147 .register_session(
1148 session.id(),
1149 SessionRuntimeState {
1150 identity: identity.clone(),
1151 generation: record.generation,
1152 fencing_token,
1153 checkpoint_version: record.checkpoint_version,
1154 },
1155 )
1156 .await
1157 .expect("retry register should flush preserved pending snapshot");
1158 let loaded = meerkat::SessionStore::load(&adapter, session.id())
1159 .await
1160 .expect("load after retry")
1161 .expect("snapshot");
1162 assert_eq!(loaded.metadata().get("pending"), Some(&json!(true)));
1163 }
1164
1165 #[tokio::test]
1166 async fn continuity_session_store_adapter_delete_if_current_revision_removes_matching_snapshot()
1167 {
1168 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1169 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1170 let session = meerkat_core::Session::new();
1171 let identity = AgentIdentity::parse("agent:quarantine").expect("identity");
1172 let record = ContinuityRecord {
1173 identity: identity.clone(),
1174 agent_runtime_id: AgentRuntimeId::parse("rt:agent:quarantine:0").expect("runtime id"),
1175 session_id: session.id().clone(),
1176 generation: ContinuityGeneration::new(0),
1177 checkpoint_version: CheckpointVersion::new(0),
1178 };
1179 let fencing_token = FencingToken::new(4);
1180 store
1181 .upsert_continuity_record(&record, fencing_token)
1182 .await
1183 .expect("seed record");
1184 adapter
1185 .register_session(
1186 session.id(),
1187 SessionRuntimeState {
1188 identity,
1189 generation: record.generation,
1190 fencing_token,
1191 checkpoint_version: record.checkpoint_version,
1192 },
1193 )
1194 .await
1195 .expect("register");
1196 meerkat::SessionStore::save(&adapter, &session)
1197 .await
1198 .expect("save snapshot");
1199
1200 let stale_revision = "row-sha256:not-current".to_string();
1201 assert!(
1202 !meerkat::SessionStore::delete_if_current_revision(
1203 &adapter,
1204 session.id(),
1205 &stale_revision
1206 )
1207 .await
1208 .expect("stale delete should be clean"),
1209 "stale revision must not delete"
1210 );
1211 assert!(
1212 store
1213 .load_session_snapshot(session.id())
1214 .await
1215 .expect("load after stale")
1216 .is_some(),
1217 "stale CAS delete must leave snapshot in place"
1218 );
1219
1220 let current_revision =
1221 meerkat_core::session_store::session_projection_cas_token(&session).expect("revision");
1222 assert!(
1223 meerkat::SessionStore::delete_if_current_revision(
1224 &adapter,
1225 session.id(),
1226 ¤t_revision
1227 )
1228 .await
1229 .expect("matching delete should succeed"),
1230 "matching revision should delete"
1231 );
1232 assert!(
1233 store
1234 .load_session_snapshot(session.id())
1235 .await
1236 .expect("load after delete")
1237 .is_none(),
1238 "matching CAS delete must remove the continuity snapshot"
1239 );
1240 assert!(
1241 meerkat::SessionStore::load(&adapter, session.id())
1242 .await
1243 .expect("adapter load after delete")
1244 .is_none(),
1245 "adapter must not synthesize a session after successful CAS delete"
1246 );
1247 }
1248
1249 #[tokio::test]
1250 async fn continuity_session_store_adapter_save_rejects_transcript_shrink() {
1251 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1252 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1253 let mut session = meerkat_core::Session::new();
1254 session.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
1255 session
1256 .append_external_user_content(meerkat_core::ContentInput::Text("second".to_string()));
1257 let identity = AgentIdentity::parse("agent:append-only").expect("identity");
1258 let record = ContinuityRecord {
1259 identity: identity.clone(),
1260 agent_runtime_id: AgentRuntimeId::parse("rt:agent:append-only:0").expect("runtime id"),
1261 session_id: session.id().clone(),
1262 generation: ContinuityGeneration::new(0),
1263 checkpoint_version: CheckpointVersion::new(0),
1264 };
1265 let fencing_token = FencingToken::new(12);
1266 store
1267 .upsert_continuity_record(&record, fencing_token)
1268 .await
1269 .expect("seed record");
1270 adapter
1271 .register_session(
1272 session.id(),
1273 SessionRuntimeState {
1274 identity,
1275 generation: record.generation,
1276 fencing_token,
1277 checkpoint_version: record.checkpoint_version,
1278 },
1279 )
1280 .await
1281 .expect("register");
1282 meerkat::SessionStore::save(&adapter, &session)
1283 .await
1284 .expect("initial save");
1285
1286 let mut stale = meerkat_core::Session::with_id(session.id().clone());
1287 stale.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
1288 let err = meerkat::SessionStore::save(&adapter, &stale)
1289 .await
1290 .expect_err("plain save must reject transcript shrink");
1291 assert!(
1292 err.to_string().contains("transcript")
1293 || err.to_string().contains("monotonicity")
1294 || err.to_string().contains("continuity"),
1295 "unexpected shrink error: {err}"
1296 );
1297 }
1298
1299 #[tokio::test]
1300 async fn continuity_session_store_adapter_saves_transcript_rewrite() {
1301 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1302 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1303 let mut session = meerkat_core::Session::new();
1304 session.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
1305 session
1306 .append_external_user_content(meerkat_core::ContentInput::Text("second".to_string()));
1307 let identity = AgentIdentity::parse("agent:rewrite").expect("identity");
1308 let record = ContinuityRecord {
1309 identity: identity.clone(),
1310 agent_runtime_id: AgentRuntimeId::parse("rt:agent:rewrite:0").expect("runtime id"),
1311 session_id: session.id().clone(),
1312 generation: ContinuityGeneration::new(0),
1313 checkpoint_version: CheckpointVersion::new(0),
1314 };
1315 let fencing_token = FencingToken::new(13);
1316 store
1317 .upsert_continuity_record(&record, fencing_token)
1318 .await
1319 .expect("seed record");
1320 adapter
1321 .register_session(
1322 session.id(),
1323 SessionRuntimeState {
1324 identity,
1325 generation: record.generation,
1326 fencing_token,
1327 checkpoint_version: record.checkpoint_version,
1328 },
1329 )
1330 .await
1331 .expect("register");
1332 meerkat::SessionStore::save(&adapter, &session)
1333 .await
1334 .expect("initial save");
1335
1336 let parent_revision = session.transcript_revision().expect("parent revision");
1337 let mut rewritten = session.clone();
1338 let commit = rewritten
1339 .commit_transcript_rewrite(
1340 meerkat_core::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1341 vec![meerkat_core::Message::User(
1342 meerkat_core::UserMessage::text("compacted first".to_string()),
1343 )],
1344 meerkat_core::TranscriptRewriteReason::new("test"),
1345 Some("mobkit-test".to_string()),
1346 Some(parent_revision),
1347 )
1348 .expect("rewrite commit");
1349
1350 meerkat::SessionStore::save_transcript_rewrite(&adapter, &rewritten, &commit)
1351 .await
1352 .expect("rewrite save should be supported");
1353 let loaded = meerkat::SessionStore::load(&adapter, session.id())
1354 .await
1355 .expect("load rewritten")
1356 .expect("rewritten session");
1357 assert_eq!(loaded.messages().len(), rewritten.messages().len());
1358 assert_eq!(
1359 loaded.transcript_revision().expect("loaded revision"),
1360 commit.revision
1361 );
1362 }
1363
1364 #[tokio::test]
1365 async fn continuity_session_store_adapter_delete_removes_current_snapshot() {
1366 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1367 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1368 let session = meerkat_core::Session::new();
1369 let identity = AgentIdentity::parse("agent:delete").expect("identity");
1370 let record = ContinuityRecord {
1371 identity: identity.clone(),
1372 agent_runtime_id: AgentRuntimeId::parse("rt:agent:delete:0").expect("runtime id"),
1373 session_id: session.id().clone(),
1374 generation: ContinuityGeneration::new(0),
1375 checkpoint_version: CheckpointVersion::new(0),
1376 };
1377 let fencing_token = FencingToken::new(7);
1378 store
1379 .upsert_continuity_record(&record, fencing_token)
1380 .await
1381 .expect("seed record");
1382 adapter
1383 .register_session(
1384 session.id(),
1385 SessionRuntimeState {
1386 identity,
1387 generation: record.generation,
1388 fencing_token,
1389 checkpoint_version: record.checkpoint_version,
1390 },
1391 )
1392 .await
1393 .expect("register");
1394 meerkat::SessionStore::save(&adapter, &session)
1395 .await
1396 .expect("save snapshot");
1397
1398 meerkat::SessionStore::delete(&adapter, session.id())
1399 .await
1400 .expect("delete should remove current snapshot");
1401 assert!(
1402 store
1403 .load_session_snapshot(session.id())
1404 .await
1405 .expect("load after delete")
1406 .is_none(),
1407 "delete must not be a successful no-op"
1408 );
1409 assert!(
1410 meerkat::SessionStore::load(&adapter, session.id())
1411 .await
1412 .expect("adapter load after delete")
1413 .is_none(),
1414 "adapter must forget registry state after delete"
1415 );
1416 }
1417
1418 #[tokio::test]
1419 async fn continuity_session_store_adapter_queues_unregistered_authoritative_projection() {
1420 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1421 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1422 let session = meerkat_core::Session::new();
1423
1424 meerkat::SessionStore::save_authoritative_projection(&adapter, &session)
1425 .await
1426 .expect("create-time authoritative projection should queue before registration");
1427 assert!(
1428 meerkat::SessionStore::load(&adapter, session.id())
1429 .await
1430 .expect("load")
1431 .is_none(),
1432 "pending authoritative projection must stay invisible until registration"
1433 );
1434
1435 let identity = AgentIdentity::parse("agent:queued").expect("identity");
1436 let record = ContinuityRecord {
1437 identity: identity.clone(),
1438 agent_runtime_id: AgentRuntimeId::parse("rt:agent:queued:0").expect("runtime id"),
1439 session_id: session.id().clone(),
1440 generation: ContinuityGeneration::new(0),
1441 checkpoint_version: CheckpointVersion::new(0),
1442 };
1443 let fencing_token = FencingToken::new(7);
1444 store
1445 .upsert_continuity_record(&record, fencing_token)
1446 .await
1447 .expect("seed record");
1448 adapter
1449 .register_session(
1450 session.id(),
1451 SessionRuntimeState {
1452 identity,
1453 generation: record.generation,
1454 fencing_token,
1455 checkpoint_version: record.checkpoint_version,
1456 },
1457 )
1458 .await
1459 .expect("register flushes pending authoritative projection");
1460 assert!(
1461 meerkat::SessionStore::load(&adapter, session.id())
1462 .await
1463 .expect("load after register")
1464 .is_some(),
1465 "registration must flush the pending authoritative projection"
1466 );
1467 }
1468
1469 #[tokio::test]
1470 async fn continuity_session_store_adapter_delete_forgets_registered_session_without_snapshot() {
1471 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1472 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1473 let session = meerkat_core::Session::new();
1474 let identity = AgentIdentity::parse("agent:delete-empty").expect("identity");
1475 let record = ContinuityRecord {
1476 identity: identity.clone(),
1477 agent_runtime_id: AgentRuntimeId::parse("rt:agent:delete-empty:0").expect("runtime id"),
1478 session_id: session.id().clone(),
1479 generation: ContinuityGeneration::new(0),
1480 checkpoint_version: CheckpointVersion::new(0),
1481 };
1482 let fencing_token = FencingToken::new(11);
1483 store
1484 .upsert_continuity_record(&record, fencing_token)
1485 .await
1486 .expect("seed record");
1487 adapter
1488 .register_session(
1489 session.id(),
1490 SessionRuntimeState {
1491 identity,
1492 generation: record.generation,
1493 fencing_token,
1494 checkpoint_version: record.checkpoint_version,
1495 },
1496 )
1497 .await
1498 .expect("register");
1499 assert!(
1500 meerkat::SessionStore::load(&adapter, session.id())
1501 .await
1502 .expect("synthetic load before delete")
1503 .is_some()
1504 );
1505
1506 meerkat::SessionStore::delete(&adapter, session.id())
1507 .await
1508 .expect("delete with no persisted snapshot should be idempotent");
1509 assert!(
1510 meerkat::SessionStore::load(&adapter, session.id())
1511 .await
1512 .expect("load after delete")
1513 .is_none(),
1514 "delete must forget registry state when no persisted row exists"
1515 );
1516 }
1517
1518 #[tokio::test]
1519 async fn continuity_session_store_adapter_authoritative_projection_cas_guards_rewrites() {
1520 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1521 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1522 let mut session = meerkat_core::Session::new();
1523 let identity = AgentIdentity::parse("agent:projection").expect("identity");
1524 let record = ContinuityRecord {
1525 identity: identity.clone(),
1526 agent_runtime_id: AgentRuntimeId::parse("rt:agent:projection:0").expect("runtime id"),
1527 session_id: session.id().clone(),
1528 generation: ContinuityGeneration::new(0),
1529 checkpoint_version: CheckpointVersion::new(0),
1530 };
1531 let fencing_token = FencingToken::new(5);
1532 store
1533 .upsert_continuity_record(&record, fencing_token)
1534 .await
1535 .expect("seed record");
1536 adapter
1537 .register_session(
1538 session.id(),
1539 SessionRuntimeState {
1540 identity: identity.clone(),
1541 generation: record.generation,
1542 fencing_token,
1543 checkpoint_version: record.checkpoint_version,
1544 },
1545 )
1546 .await
1547 .expect("register");
1548
1549 meerkat::SessionStore::save_authoritative_projection_if_current_revision(
1550 &adapter, &session, None,
1551 )
1552 .await
1553 .expect("initial projection should accept missing current revision");
1554 let original_revision =
1555 meerkat_core::session_store::session_projection_cas_token(&session).expect("revision");
1556
1557 let mut stale_rewrite = session.clone();
1558 stale_rewrite.set_metadata("projection", json!("stale"));
1559 let stale_error = meerkat::SessionStore::save_authoritative_projection_if_current_revision(
1560 &adapter,
1561 &stale_rewrite,
1562 Some("row-sha256:not-current".to_string()),
1563 )
1564 .await
1565 .expect_err("stale CAS projection must reject");
1566 assert!(
1567 stale_error.to_string().contains("not a continuation"),
1568 "unexpected stale error: {stale_error}"
1569 );
1570
1571 let loaded = meerkat::SessionStore::load(&adapter, session.id())
1572 .await
1573 .expect("load")
1574 .expect("snapshot");
1575 assert_eq!(
1576 meerkat_core::session_store::session_projection_cas_token(&loaded).expect("revision"),
1577 original_revision,
1578 "stale authoritative projection must leave stored row unchanged"
1579 );
1580
1581 session.set_metadata("projection", json!("current"));
1582 meerkat::SessionStore::save_authoritative_projection_if_current_revision(
1583 &adapter,
1584 &session,
1585 Some(original_revision),
1586 )
1587 .await
1588 .expect("matching CAS projection should save");
1589
1590 let loaded = meerkat::SessionStore::load(&adapter, session.id())
1591 .await
1592 .expect("load after save")
1593 .expect("snapshot after save");
1594 assert_eq!(loaded.metadata().get("projection"), Some(&json!("current")));
1595 }
1596}