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