Skip to main content

meerkat_mobkit/identity_first/
adapters.rs

1//! Compatibility adapters bridging legacy MobKit traits to identity-first contracts.
2//!
3//! - [`DiscoveryRosterAdapter`]: `Discovery` → `RosterProvider` (CONTRACT-08, REQ-27, REQ-28)
4//! - [`EdgeDiscoveryTopologyAdapter`]: `EdgeDiscovery` → `TopologyProvider` (CONTRACT-09, REQ-29)
5//! - [`ContinuitySessionStoreAdapter`]: `ContinuityStore` → `SessionStore` (CONTRACT-10)
6//! - [`SessionHookCustomizerAdapter`]: `SessionHook` → `AgentCustomizer` (CONTRACT-11, REQ-30)
7
8use 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
24// ---------------------------------------------------------------------------
25// CONTRACT-08 / REQ-27 / REQ-28: Discovery → RosterProvider
26// ---------------------------------------------------------------------------
27
28/// Adapts a legacy `Discovery` trait impl into a `RosterProvider`.
29///
30/// Maps `AgentDiscoverySpec` to `DurableAgentSpec` per REQ-27:
31/// - `meerkat_id` → `identity` (parsed as `AgentIdentity`)
32/// - `profile` → `profile`
33/// - `labels` → `labels`
34/// - `context` → `context`
35/// - `additional_instructions` → `additional_instructions`
36/// - `resume_session_id` → ignored
37/// - `addressability` → `Addressable`
38/// - `display_name` → `None`
39pub 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
51/// Convert an `AgentDiscoverySpec` to a `DurableAgentSpec` per REQ-27.
52pub 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
80// ---------------------------------------------------------------------------
81// CONTRACT-09 / REQ-29: EdgeDiscovery → TopologyProvider
82// ---------------------------------------------------------------------------
83
84/// Adapts a legacy `EdgeDiscovery` trait impl into a `TopologyProvider`.
85///
86/// Parses `DesiredPeerEdge` endpoint strings as `AgentIdentity` to produce
87/// `ManagedPeerEdge` instances.
88pub 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        // Project the roster context to EdgeMemberView so legacy EdgeDiscovery
108        // impls see real member identities/labels instead of an empty vec.
109        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// ---------------------------------------------------------------------------
137// CONTRACT-10: ContinuityStore → SessionStore adapter
138// ---------------------------------------------------------------------------
139
140/// Runtime state for a session, used by the adapter to resolve identity/fencing.
141#[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
149/// Adapts a `ContinuityStore` to the Meerkat `SessionStore` interface.
150///
151/// On the external-authoritative path, this is the session persistence layer.
152/// No separate local SQLite is created under scratch_dir — the ContinuityStore
153/// is the single authoritative session truth.
154///
155/// The adapter maintains a session→identity registry populated by the runtime
156/// during restore/activate. This maps SessionId to the owning identity's
157/// continuity parameters (identity, generation, fencing_token).
158pub struct ContinuitySessionStoreAdapter {
159    store: Arc<dyn super::contracts::ContinuityStore>,
160    /// Per-session monotonic version counter to satisfy CAS on repeated saves.
161    versions: Mutex<HashMap<String, AtomicU64>>,
162    /// Session→identity mapping, populated by the runtime.
163    session_registry: Mutex<HashMap<String, SessionRuntimeState>>,
164    /// Session saves that arrive before the bridge can publish the owning
165    /// identity. These are flushed immediately when the session is registered.
166    pending_unregistered: Mutex<HashMap<String, Vec<u8>>>,
167    /// Sessions that were explicitly unregistered. Later writes from those
168    /// actors must fail closed instead of becoming pre-registration pending
169    /// snapshots for a future session with the same id.
170    unregistered_sessions: Mutex<HashSet<String>>,
171    /// Serializes persisted session writes, including projection CAS load/write pairs.
172    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    /// Register a session with its owning identity's runtime state.
188    ///
189    /// Called by the runtime during restore/activate to wire real
190    /// identity/generation/fencing data into the adapter.
191    #[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    /// Update the fencing token for a session (e.g., after lease renewal).
257    #[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    /// Get the next checkpoint version for a session, starting at 1.
309    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    /// Look up the runtime state for a session.
371    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        // Use real identity/generation/fencing from the runtime registry.
463        match self.lookup_session(&sid_str) {
464            Some(state) => {
465                self.save_registered_snapshot(session.id(), data, state)
466                    .await?;
467            }
468            None => {
469                // PersistentSessionService can save during member creation
470                // before the bridge call returns. Hold that first snapshot
471                // until the identity runtime registers the real owner; never
472                // checkpoint under a synthetic `_session:*` identity.
473                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        // Listing is not supported through the continuity store adapter.
600        // The continuity model is identity-keyed, not session-list-keyed.
601        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, &current_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
660// ---------------------------------------------------------------------------
661// CONTRACT-11 / REQ-30: SessionHook → AgentCustomizer adapter
662// ---------------------------------------------------------------------------
663
664/// Adapts a legacy `SessionHook` to the `AgentCustomizer` trait.
665///
666/// Constructs a synthetic `CreateSessionRequest` from the `AgentBuildDraft`,
667/// lets the hook mutate it, and writes supported mutations back. Unsupported
668/// field mutations (e.g., `resume_session`) are detected and logged as warnings.
669pub 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        // Build a synthetic CreateSessionRequest from the draft
688        let mut req = meerkat_core::service::CreateSessionRequest {
689            model: draft.model.clone().unwrap_or_default(),
690            prompt: meerkat_core::ContentInput::Text(String::new()),
691            // Meerkat 0.7: per-request system prompt is the typed tri-state
692            // override; a draft prompt maps to an explicit `Set`.
693            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        // Snapshot "before" state for unsupported-mutation detection (REQ-30).
710        // Supported: model, system_prompt, labels — written back to draft.
711        // Unsupported: everything else — warn if hook mutated them.
712        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        // Detect unsupported mutations by comparing before/after (REQ-30).
724        // Warn for each mutated field — mutations are NOT applied to the draft.
725        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        // Meerkat 0.7 removed the flat `render_metadata` / `skill_references`
740        // request fields; both now live only on the typed
741        // `build.initial_turn_metadata` carrier, which the `build` mutation
742        // detection below already covers.
743        if let Some(ref build) = req.build {
744            if build_before_is_none {
745                // Hook created a build block — any build.* is unsupported
746                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                // build existed before but hook added resume_session
752                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        // Apply supported mutations back to the draft.
767        //
768        // NOTE: `additional_instructions` is part of the AgentCustomizer mutation
769        // surface (REQ-30), but `CreateSessionRequest` does not expose it as a
770        // field. Legacy hooks therefore cannot modify additional_instructions —
771        // the draft's existing value from the DurableAgentSpec passes through
772        // untouched. Native `AgentCustomizer` impls CAN mutate it directly.
773        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                &current_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}