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 };
708
709 let prompt_before = req.prompt.clone();
713 let max_tokens_before = req.max_tokens;
714 let event_tx_was_some = req.event_tx.is_some();
715 let initial_turn_before = req.initial_turn;
716 let build_before_is_none = req.build.is_none();
717
718 self.hook
719 .before_create(&mut req)
720 .await
721 .map_err(|e| CustomizerError::BuildFailed(format!("session hook: {e}")))?;
722
723 let mut unsupported_mutations: Vec<&str> = Vec::new();
726
727 if req.prompt != prompt_before {
728 unsupported_mutations.push("prompt");
729 }
730 if req.max_tokens != max_tokens_before {
731 unsupported_mutations.push("max_tokens");
732 }
733 if req.event_tx.is_some() != event_tx_was_some {
734 unsupported_mutations.push("event_tx");
735 }
736 if req.initial_turn != initial_turn_before {
737 unsupported_mutations.push("initial_turn");
738 }
739 if let Some(ref build) = req.build {
744 if build_before_is_none {
745 unsupported_mutations.push("build");
747 if build.resume_session.is_some() {
748 unsupported_mutations.push("build.resume_session");
749 }
750 } else if build.resume_session.is_some() {
751 unsupported_mutations.push("build.resume_session");
753 }
754 }
755
756 if !unsupported_mutations.is_empty() {
757 tracing::warn!(
758 identity = %spec.identity,
759 fields = ?unsupported_mutations,
760 "SessionHook mutated unsupported CreateSessionRequest fields — \
761 these mutations are NOT applied in the identity-first model. \
762 Migrate to AgentCustomizer."
763 );
764 }
765
766 if !req.model.is_empty() {
774 draft.model = Some(req.model);
775 }
776 draft.system_prompt = req.system_prompt.as_set_prompt().map(ToString::to_string);
777 draft.labels = req.labels.unwrap_or_default();
778
779 Ok(())
780 }
781
782 async fn after_create(
783 &self,
784 _identity: &AgentIdentity,
785 session_id: &meerkat_core::types::SessionId,
786 context: &SessionCreatedContext,
787 ) -> Result<(), CustomizerError> {
788 self.hook.after_create(session_id, context).await;
789 Ok(())
790 }
791}
792
793#[cfg(test)]
794#[allow(clippy::expect_used, clippy::panic)]
795mod tests {
796 use std::sync::Arc;
797 use std::sync::atomic::{AtomicBool, Ordering as AtomicOrdering};
798
799 use serde_json::json;
800
801 use super::super::contracts::ContinuityStore;
802 use super::super::local_store::LocalContinuityStore;
803 use super::super::types::{
804 AgentIdentity, AgentRuntimeId, CheckpointVersion, ContinuityGeneration, ContinuityRecord,
805 ContinuityResolveState, ContinuityStoreError, FencingToken, SessionSnapshot,
806 };
807 use super::*;
808
809 struct FailSaveContinuityStore {
810 inner: Arc<LocalContinuityStore>,
811 fail_save: AtomicBool,
812 }
813
814 impl FailSaveContinuityStore {
815 fn new(inner: Arc<LocalContinuityStore>) -> Self {
816 Self {
817 inner,
818 fail_save: AtomicBool::new(false),
819 }
820 }
821
822 fn fail_saves(&self, fail: bool) {
823 self.fail_save.store(fail, AtomicOrdering::SeqCst);
824 }
825 }
826
827 #[async_trait]
828 impl ContinuityStore for FailSaveContinuityStore {
829 async fn resolve_many(
830 &self,
831 identities: &[AgentIdentity],
832 ) -> Result<
833 std::collections::BTreeMap<AgentIdentity, ContinuityResolveState>,
834 ContinuityStoreError,
835 > {
836 self.inner.resolve_many(identities).await
837 }
838
839 async fn load_session_snapshot(
840 &self,
841 session_id: &meerkat_core::types::SessionId,
842 ) -> Result<Option<SessionSnapshot>, ContinuityStoreError> {
843 self.inner.load_session_snapshot(session_id).await
844 }
845
846 async fn delete_session_snapshot_if_current_revision(
847 &self,
848 session_id: &meerkat_core::types::SessionId,
849 expected_current_revision: &str,
850 ) -> Result<bool, ContinuityStoreError> {
851 self.inner
852 .delete_session_snapshot_if_current_revision(session_id, expected_current_revision)
853 .await
854 }
855
856 async fn save_session_snapshot(
857 &self,
858 identity: &AgentIdentity,
859 session_id: &meerkat_core::types::SessionId,
860 generation: ContinuityGeneration,
861 version: CheckpointVersion,
862 fencing_token: FencingToken,
863 snapshot: &SessionSnapshot,
864 ) -> Result<(), ContinuityStoreError> {
865 if self.fail_save.load(AtomicOrdering::SeqCst) {
866 return Err(ContinuityStoreError::Io("forced save failure".to_string()));
867 }
868 self.inner
869 .save_session_snapshot(
870 identity,
871 session_id,
872 generation,
873 version,
874 fencing_token,
875 snapshot,
876 )
877 .await
878 }
879
880 async fn upsert_continuity_record(
881 &self,
882 record: &ContinuityRecord,
883 fencing_token: FencingToken,
884 ) -> Result<(), ContinuityStoreError> {
885 self.inner
886 .upsert_continuity_record(record, fencing_token)
887 .await
888 }
889
890 async fn delete_continuity_record(
891 &self,
892 identity: &AgentIdentity,
893 fencing_token: FencingToken,
894 ) -> Result<(), ContinuityStoreError> {
895 self.inner
896 .delete_continuity_record(identity, fencing_token)
897 .await
898 }
899 }
900
901 #[tokio::test]
902 async fn continuity_session_store_adapter_seeds_registered_checkpoint_version() {
903 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
904 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
905 let session = meerkat_core::Session::new();
906 let identity = AgentIdentity::parse("agent:restored").expect("identity");
907 let record = ContinuityRecord {
908 identity: identity.clone(),
909 agent_runtime_id: AgentRuntimeId::parse("rt:agent:restored:0").expect("runtime id"),
910 session_id: session.id().clone(),
911 generation: ContinuityGeneration::new(2),
912 checkpoint_version: CheckpointVersion::new(5),
913 };
914 let fencing_token = FencingToken::new(9);
915 store
916 .upsert_continuity_record(&record, fencing_token)
917 .await
918 .expect("seed record");
919
920 adapter
921 .register_session(
922 session.id(),
923 SessionRuntimeState {
924 identity: identity.clone(),
925 generation: record.generation,
926 fencing_token,
927 checkpoint_version: record.checkpoint_version,
928 },
929 )
930 .await
931 .expect("register");
932
933 meerkat::SessionStore::save(&adapter, &session)
934 .await
935 .expect("save should advance from restored checkpoint");
936 let effective_version = adapter
937 .register_session(
938 session.id(),
939 SessionRuntimeState {
940 identity: identity.clone(),
941 generation: record.generation,
942 fencing_token,
943 checkpoint_version: record.checkpoint_version,
944 },
945 )
946 .await
947 .expect("post-save register should report advanced version");
948 assert_eq!(effective_version, CheckpointVersion::new(6));
949 let resolved = store
950 .resolve_many(std::slice::from_ref(&identity))
951 .await
952 .expect("resolve");
953 let ContinuityResolveState::Ready { record } = resolved.get(&identity).expect("record")
954 else {
955 panic!("expected ready record");
956 };
957 assert_eq!(record.checkpoint_version, CheckpointVersion::new(6));
958 }
959
960 #[tokio::test]
961 async fn continuity_session_store_adapter_flushes_pending_save_under_registered_identity() {
962 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
963 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
964 let session = meerkat_core::Session::new();
965 let identity = AgentIdentity::parse("agent:fresh").expect("identity");
966 let record = ContinuityRecord {
967 identity: identity.clone(),
968 agent_runtime_id: AgentRuntimeId::parse("rt:agent:fresh:0").expect("runtime id"),
969 session_id: session.id().clone(),
970 generation: ContinuityGeneration::new(0),
971 checkpoint_version: CheckpointVersion::new(0),
972 };
973 let fencing_token = FencingToken::new(3);
974 store
975 .upsert_continuity_record(&record, fencing_token)
976 .await
977 .expect("seed record");
978
979 meerkat::SessionStore::save(&adapter, &session)
980 .await
981 .expect("unregistered save should be delayed, not written under fallback identity");
982 assert!(
983 store
984 .load_session_snapshot(session.id())
985 .await
986 .expect("load before register")
987 .is_none(),
988 "unregistered save must not be visible in continuity store"
989 );
990
991 adapter
992 .register_session(
993 session.id(),
994 SessionRuntimeState {
995 identity: identity.clone(),
996 generation: record.generation,
997 fencing_token,
998 checkpoint_version: record.checkpoint_version,
999 },
1000 )
1001 .await
1002 .expect("register flushes pending");
1003
1004 assert!(
1005 store
1006 .load_session_snapshot(session.id())
1007 .await
1008 .expect("load after register")
1009 .is_some(),
1010 "pending save should flush under the registered identity"
1011 );
1012 let resolved = store
1013 .resolve_many(std::slice::from_ref(&identity))
1014 .await
1015 .expect("resolve");
1016 let ContinuityResolveState::Ready { record } = resolved.get(&identity).expect("record")
1017 else {
1018 panic!("expected ready record");
1019 };
1020 assert_eq!(record.checkpoint_version, CheckpointVersion::new(1));
1021 }
1022
1023 #[tokio::test]
1024 async fn continuity_session_store_adapter_rejects_saves_after_unregister() {
1025 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1026 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1027 let session = meerkat_core::Session::new();
1028 let identity = AgentIdentity::parse("agent:retired").expect("identity");
1029 let record = ContinuityRecord {
1030 identity: identity.clone(),
1031 agent_runtime_id: AgentRuntimeId::parse("rt:agent:retired:0").expect("runtime id"),
1032 session_id: session.id().clone(),
1033 generation: ContinuityGeneration::new(0),
1034 checkpoint_version: CheckpointVersion::new(0),
1035 };
1036 let fencing_token = FencingToken::new(9);
1037 store
1038 .upsert_continuity_record(&record, fencing_token)
1039 .await
1040 .expect("seed record");
1041 adapter
1042 .register_session(
1043 session.id(),
1044 SessionRuntimeState {
1045 identity: identity.clone(),
1046 generation: record.generation,
1047 fencing_token,
1048 checkpoint_version: record.checkpoint_version,
1049 },
1050 )
1051 .await
1052 .expect("register");
1053
1054 adapter
1055 .unregister_session(session.id())
1056 .await
1057 .expect("unregister");
1058 let err = meerkat::SessionStore::save(&adapter, &session)
1059 .await
1060 .expect_err("post-unregister save must fail closed");
1061 assert!(
1062 err.to_string().contains("was unregistered"),
1063 "unexpected error: {err}"
1064 );
1065 assert!(
1066 store
1067 .load_session_snapshot(session.id())
1068 .await
1069 .expect("load")
1070 .is_none(),
1071 "post-unregister save must not be queued as pending"
1072 );
1073
1074 adapter
1075 .register_session(
1076 session.id(),
1077 SessionRuntimeState {
1078 identity,
1079 generation: record.generation,
1080 fencing_token,
1081 checkpoint_version: record.checkpoint_version,
1082 },
1083 )
1084 .await
1085 .expect("registering the same id later should not flush stale pending data");
1086 assert!(
1087 store
1088 .load_session_snapshot(session.id())
1089 .await
1090 .expect("load after re-register")
1091 .is_none(),
1092 "stale post-unregister save must not flush on a later registration"
1093 );
1094 }
1095
1096 #[tokio::test]
1097 async fn continuity_session_store_adapter_register_keeps_pending_snapshot_on_flush_failure() {
1098 let inner = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1099 let fail_store = Arc::new(FailSaveContinuityStore::new(inner.clone()));
1100 let adapter = ContinuitySessionStoreAdapter::new(fail_store.clone());
1101 let mut session = meerkat_core::Session::new();
1102 session.set_metadata("pending", json!(true));
1103 let identity = AgentIdentity::parse("agent:pending-fail").expect("identity");
1104 let record = ContinuityRecord {
1105 identity: identity.clone(),
1106 agent_runtime_id: AgentRuntimeId::parse("rt:agent:pending-fail:0").expect("runtime id"),
1107 session_id: session.id().clone(),
1108 generation: ContinuityGeneration::new(0),
1109 checkpoint_version: CheckpointVersion::new(0),
1110 };
1111 let fencing_token = FencingToken::new(14);
1112 inner
1113 .upsert_continuity_record(&record, fencing_token)
1114 .await
1115 .expect("seed record");
1116
1117 meerkat::SessionStore::save(&adapter, &session)
1118 .await
1119 .expect("pending save");
1120 fail_store.fail_saves(true);
1121 adapter
1122 .register_session(
1123 session.id(),
1124 SessionRuntimeState {
1125 identity: identity.clone(),
1126 generation: record.generation,
1127 fencing_token,
1128 checkpoint_version: record.checkpoint_version,
1129 },
1130 )
1131 .await
1132 .expect_err("forced pending flush failure");
1133 assert!(
1134 meerkat::SessionStore::load(&adapter, session.id())
1135 .await
1136 .expect("load after failed register")
1137 .is_none(),
1138 "failed register must not leave a synthetic registered session"
1139 );
1140
1141 fail_store.fail_saves(false);
1142 adapter
1143 .register_session(
1144 session.id(),
1145 SessionRuntimeState {
1146 identity: identity.clone(),
1147 generation: record.generation,
1148 fencing_token,
1149 checkpoint_version: record.checkpoint_version,
1150 },
1151 )
1152 .await
1153 .expect("retry register should flush preserved pending snapshot");
1154 let loaded = meerkat::SessionStore::load(&adapter, session.id())
1155 .await
1156 .expect("load after retry")
1157 .expect("snapshot");
1158 assert_eq!(loaded.metadata().get("pending"), Some(&json!(true)));
1159 }
1160
1161 #[tokio::test]
1162 async fn continuity_session_store_adapter_delete_if_current_revision_removes_matching_snapshot()
1163 {
1164 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1165 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1166 let session = meerkat_core::Session::new();
1167 let identity = AgentIdentity::parse("agent:quarantine").expect("identity");
1168 let record = ContinuityRecord {
1169 identity: identity.clone(),
1170 agent_runtime_id: AgentRuntimeId::parse("rt:agent:quarantine:0").expect("runtime id"),
1171 session_id: session.id().clone(),
1172 generation: ContinuityGeneration::new(0),
1173 checkpoint_version: CheckpointVersion::new(0),
1174 };
1175 let fencing_token = FencingToken::new(4);
1176 store
1177 .upsert_continuity_record(&record, fencing_token)
1178 .await
1179 .expect("seed record");
1180 adapter
1181 .register_session(
1182 session.id(),
1183 SessionRuntimeState {
1184 identity,
1185 generation: record.generation,
1186 fencing_token,
1187 checkpoint_version: record.checkpoint_version,
1188 },
1189 )
1190 .await
1191 .expect("register");
1192 meerkat::SessionStore::save(&adapter, &session)
1193 .await
1194 .expect("save snapshot");
1195
1196 let stale_revision = "row-sha256:not-current".to_string();
1197 assert!(
1198 !meerkat::SessionStore::delete_if_current_revision(
1199 &adapter,
1200 session.id(),
1201 &stale_revision
1202 )
1203 .await
1204 .expect("stale delete should be clean"),
1205 "stale revision must not delete"
1206 );
1207 assert!(
1208 store
1209 .load_session_snapshot(session.id())
1210 .await
1211 .expect("load after stale")
1212 .is_some(),
1213 "stale CAS delete must leave snapshot in place"
1214 );
1215
1216 let current_revision =
1217 meerkat_core::session_store::session_projection_cas_token(&session).expect("revision");
1218 assert!(
1219 meerkat::SessionStore::delete_if_current_revision(
1220 &adapter,
1221 session.id(),
1222 ¤t_revision
1223 )
1224 .await
1225 .expect("matching delete should succeed"),
1226 "matching revision should delete"
1227 );
1228 assert!(
1229 store
1230 .load_session_snapshot(session.id())
1231 .await
1232 .expect("load after delete")
1233 .is_none(),
1234 "matching CAS delete must remove the continuity snapshot"
1235 );
1236 assert!(
1237 meerkat::SessionStore::load(&adapter, session.id())
1238 .await
1239 .expect("adapter load after delete")
1240 .is_none(),
1241 "adapter must not synthesize a session after successful CAS delete"
1242 );
1243 }
1244
1245 #[tokio::test]
1246 async fn continuity_session_store_adapter_save_rejects_transcript_shrink() {
1247 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1248 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1249 let mut session = meerkat_core::Session::new();
1250 session.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
1251 session
1252 .append_external_user_content(meerkat_core::ContentInput::Text("second".to_string()));
1253 let identity = AgentIdentity::parse("agent:append-only").expect("identity");
1254 let record = ContinuityRecord {
1255 identity: identity.clone(),
1256 agent_runtime_id: AgentRuntimeId::parse("rt:agent:append-only:0").expect("runtime id"),
1257 session_id: session.id().clone(),
1258 generation: ContinuityGeneration::new(0),
1259 checkpoint_version: CheckpointVersion::new(0),
1260 };
1261 let fencing_token = FencingToken::new(12);
1262 store
1263 .upsert_continuity_record(&record, fencing_token)
1264 .await
1265 .expect("seed record");
1266 adapter
1267 .register_session(
1268 session.id(),
1269 SessionRuntimeState {
1270 identity,
1271 generation: record.generation,
1272 fencing_token,
1273 checkpoint_version: record.checkpoint_version,
1274 },
1275 )
1276 .await
1277 .expect("register");
1278 meerkat::SessionStore::save(&adapter, &session)
1279 .await
1280 .expect("initial save");
1281
1282 let mut stale = meerkat_core::Session::with_id(session.id().clone());
1283 stale.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
1284 let err = meerkat::SessionStore::save(&adapter, &stale)
1285 .await
1286 .expect_err("plain save must reject transcript shrink");
1287 assert!(
1288 err.to_string().contains("transcript")
1289 || err.to_string().contains("monotonicity")
1290 || err.to_string().contains("continuity"),
1291 "unexpected shrink error: {err}"
1292 );
1293 }
1294
1295 #[tokio::test]
1296 async fn continuity_session_store_adapter_saves_transcript_rewrite() {
1297 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1298 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1299 let mut session = meerkat_core::Session::new();
1300 session.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
1301 session
1302 .append_external_user_content(meerkat_core::ContentInput::Text("second".to_string()));
1303 let identity = AgentIdentity::parse("agent:rewrite").expect("identity");
1304 let record = ContinuityRecord {
1305 identity: identity.clone(),
1306 agent_runtime_id: AgentRuntimeId::parse("rt:agent:rewrite:0").expect("runtime id"),
1307 session_id: session.id().clone(),
1308 generation: ContinuityGeneration::new(0),
1309 checkpoint_version: CheckpointVersion::new(0),
1310 };
1311 let fencing_token = FencingToken::new(13);
1312 store
1313 .upsert_continuity_record(&record, fencing_token)
1314 .await
1315 .expect("seed record");
1316 adapter
1317 .register_session(
1318 session.id(),
1319 SessionRuntimeState {
1320 identity,
1321 generation: record.generation,
1322 fencing_token,
1323 checkpoint_version: record.checkpoint_version,
1324 },
1325 )
1326 .await
1327 .expect("register");
1328 meerkat::SessionStore::save(&adapter, &session)
1329 .await
1330 .expect("initial save");
1331
1332 let parent_revision = session.transcript_revision().expect("parent revision");
1333 let mut rewritten = session.clone();
1334 let commit = rewritten
1335 .commit_transcript_rewrite(
1336 meerkat_core::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1337 vec![meerkat_core::Message::User(
1338 meerkat_core::UserMessage::text("compacted first".to_string()),
1339 )],
1340 meerkat_core::TranscriptRewriteReason::new("test"),
1341 Some("mobkit-test".to_string()),
1342 Some(parent_revision),
1343 )
1344 .expect("rewrite commit");
1345
1346 meerkat::SessionStore::save_transcript_rewrite(&adapter, &rewritten, &commit)
1347 .await
1348 .expect("rewrite save should be supported");
1349 let loaded = meerkat::SessionStore::load(&adapter, session.id())
1350 .await
1351 .expect("load rewritten")
1352 .expect("rewritten session");
1353 assert_eq!(loaded.messages().len(), rewritten.messages().len());
1354 assert_eq!(
1355 loaded.transcript_revision().expect("loaded revision"),
1356 commit.revision
1357 );
1358 }
1359
1360 #[tokio::test]
1361 async fn continuity_session_store_adapter_delete_removes_current_snapshot() {
1362 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1363 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1364 let session = meerkat_core::Session::new();
1365 let identity = AgentIdentity::parse("agent:delete").expect("identity");
1366 let record = ContinuityRecord {
1367 identity: identity.clone(),
1368 agent_runtime_id: AgentRuntimeId::parse("rt:agent:delete:0").expect("runtime id"),
1369 session_id: session.id().clone(),
1370 generation: ContinuityGeneration::new(0),
1371 checkpoint_version: CheckpointVersion::new(0),
1372 };
1373 let fencing_token = FencingToken::new(7);
1374 store
1375 .upsert_continuity_record(&record, fencing_token)
1376 .await
1377 .expect("seed record");
1378 adapter
1379 .register_session(
1380 session.id(),
1381 SessionRuntimeState {
1382 identity,
1383 generation: record.generation,
1384 fencing_token,
1385 checkpoint_version: record.checkpoint_version,
1386 },
1387 )
1388 .await
1389 .expect("register");
1390 meerkat::SessionStore::save(&adapter, &session)
1391 .await
1392 .expect("save snapshot");
1393
1394 meerkat::SessionStore::delete(&adapter, session.id())
1395 .await
1396 .expect("delete should remove current snapshot");
1397 assert!(
1398 store
1399 .load_session_snapshot(session.id())
1400 .await
1401 .expect("load after delete")
1402 .is_none(),
1403 "delete must not be a successful no-op"
1404 );
1405 assert!(
1406 meerkat::SessionStore::load(&adapter, session.id())
1407 .await
1408 .expect("adapter load after delete")
1409 .is_none(),
1410 "adapter must forget registry state after delete"
1411 );
1412 }
1413
1414 #[tokio::test]
1415 async fn continuity_session_store_adapter_queues_unregistered_authoritative_projection() {
1416 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1417 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1418 let session = meerkat_core::Session::new();
1419
1420 meerkat::SessionStore::save_authoritative_projection(&adapter, &session)
1421 .await
1422 .expect("create-time authoritative projection should queue before registration");
1423 assert!(
1424 meerkat::SessionStore::load(&adapter, session.id())
1425 .await
1426 .expect("load")
1427 .is_none(),
1428 "pending authoritative projection must stay invisible until registration"
1429 );
1430
1431 let identity = AgentIdentity::parse("agent:queued").expect("identity");
1432 let record = ContinuityRecord {
1433 identity: identity.clone(),
1434 agent_runtime_id: AgentRuntimeId::parse("rt:agent:queued:0").expect("runtime id"),
1435 session_id: session.id().clone(),
1436 generation: ContinuityGeneration::new(0),
1437 checkpoint_version: CheckpointVersion::new(0),
1438 };
1439 let fencing_token = FencingToken::new(7);
1440 store
1441 .upsert_continuity_record(&record, fencing_token)
1442 .await
1443 .expect("seed record");
1444 adapter
1445 .register_session(
1446 session.id(),
1447 SessionRuntimeState {
1448 identity,
1449 generation: record.generation,
1450 fencing_token,
1451 checkpoint_version: record.checkpoint_version,
1452 },
1453 )
1454 .await
1455 .expect("register flushes pending authoritative projection");
1456 assert!(
1457 meerkat::SessionStore::load(&adapter, session.id())
1458 .await
1459 .expect("load after register")
1460 .is_some(),
1461 "registration must flush the pending authoritative projection"
1462 );
1463 }
1464
1465 #[tokio::test]
1466 async fn continuity_session_store_adapter_delete_forgets_registered_session_without_snapshot() {
1467 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1468 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1469 let session = meerkat_core::Session::new();
1470 let identity = AgentIdentity::parse("agent:delete-empty").expect("identity");
1471 let record = ContinuityRecord {
1472 identity: identity.clone(),
1473 agent_runtime_id: AgentRuntimeId::parse("rt:agent:delete-empty:0").expect("runtime id"),
1474 session_id: session.id().clone(),
1475 generation: ContinuityGeneration::new(0),
1476 checkpoint_version: CheckpointVersion::new(0),
1477 };
1478 let fencing_token = FencingToken::new(11);
1479 store
1480 .upsert_continuity_record(&record, fencing_token)
1481 .await
1482 .expect("seed record");
1483 adapter
1484 .register_session(
1485 session.id(),
1486 SessionRuntimeState {
1487 identity,
1488 generation: record.generation,
1489 fencing_token,
1490 checkpoint_version: record.checkpoint_version,
1491 },
1492 )
1493 .await
1494 .expect("register");
1495 assert!(
1496 meerkat::SessionStore::load(&adapter, session.id())
1497 .await
1498 .expect("synthetic load before delete")
1499 .is_some()
1500 );
1501
1502 meerkat::SessionStore::delete(&adapter, session.id())
1503 .await
1504 .expect("delete with no persisted snapshot should be idempotent");
1505 assert!(
1506 meerkat::SessionStore::load(&adapter, session.id())
1507 .await
1508 .expect("load after delete")
1509 .is_none(),
1510 "delete must forget registry state when no persisted row exists"
1511 );
1512 }
1513
1514 #[tokio::test]
1515 async fn continuity_session_store_adapter_authoritative_projection_cas_guards_rewrites() {
1516 let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
1517 let adapter = ContinuitySessionStoreAdapter::new(store.clone());
1518 let mut session = meerkat_core::Session::new();
1519 let identity = AgentIdentity::parse("agent:projection").expect("identity");
1520 let record = ContinuityRecord {
1521 identity: identity.clone(),
1522 agent_runtime_id: AgentRuntimeId::parse("rt:agent:projection:0").expect("runtime id"),
1523 session_id: session.id().clone(),
1524 generation: ContinuityGeneration::new(0),
1525 checkpoint_version: CheckpointVersion::new(0),
1526 };
1527 let fencing_token = FencingToken::new(5);
1528 store
1529 .upsert_continuity_record(&record, fencing_token)
1530 .await
1531 .expect("seed record");
1532 adapter
1533 .register_session(
1534 session.id(),
1535 SessionRuntimeState {
1536 identity: identity.clone(),
1537 generation: record.generation,
1538 fencing_token,
1539 checkpoint_version: record.checkpoint_version,
1540 },
1541 )
1542 .await
1543 .expect("register");
1544
1545 meerkat::SessionStore::save_authoritative_projection_if_current_revision(
1546 &adapter, &session, None,
1547 )
1548 .await
1549 .expect("initial projection should accept missing current revision");
1550 let original_revision =
1551 meerkat_core::session_store::session_projection_cas_token(&session).expect("revision");
1552
1553 let mut stale_rewrite = session.clone();
1554 stale_rewrite.set_metadata("projection", json!("stale"));
1555 let stale_error = meerkat::SessionStore::save_authoritative_projection_if_current_revision(
1556 &adapter,
1557 &stale_rewrite,
1558 Some("row-sha256:not-current".to_string()),
1559 )
1560 .await
1561 .expect_err("stale CAS projection must reject");
1562 assert!(
1563 stale_error.to_string().contains("not a continuation"),
1564 "unexpected stale error: {stale_error}"
1565 );
1566
1567 let loaded = meerkat::SessionStore::load(&adapter, session.id())
1568 .await
1569 .expect("load")
1570 .expect("snapshot");
1571 assert_eq!(
1572 meerkat_core::session_store::session_projection_cas_token(&loaded).expect("revision"),
1573 original_revision,
1574 "stale authoritative projection must leave stored row unchanged"
1575 );
1576
1577 session.set_metadata("projection", json!("current"));
1578 meerkat::SessionStore::save_authoritative_projection_if_current_revision(
1579 &adapter,
1580 &session,
1581 Some(original_revision),
1582 )
1583 .await
1584 .expect("matching CAS projection should save");
1585
1586 let loaded = meerkat::SessionStore::load(&adapter, session.id())
1587 .await
1588 .expect("load after save")
1589 .expect("snapshot after save");
1590 assert_eq!(loaded.metadata().get("projection"), Some(&json!("current")));
1591 }
1592}