Skip to main content

meerkat_mobkit/identity_first/
bridge.rs

1//! Session bridge: connects the identity-first control plane to the Meerkat
2//! session pipeline for real session creation, delivery, and retirement.
3
4use std::collections::{BTreeMap, HashMap};
5use std::sync::Arc;
6
7use async_trait::async_trait;
8use meerkat_core::types::HandlingMode;
9use meerkat_mob::ids::AgentIdentity as MobAgentIdentity;
10use meerkat_mob::launch::MemberLaunchMode;
11use meerkat_mob::{
12    MobHandle, MobSessionService, SpawnMemberSpec, SpawnSystemPromptOverride, WorkOrigin, WorkRef,
13    WorkSpec,
14};
15
16use crate::mob_handle_runtime::{
17    content_input_has_images, is_previous_member_cleanup_ambiguous_error,
18    is_recoverable_lifecycle_cleanup_error, is_recoverable_session_owned_retire_cleanup_error,
19    model_capabilities_for_member, topology_restore_failed_peer_ids,
20};
21
22use super::adapters::{ContinuitySessionStoreAdapter, SessionRuntimeState};
23use super::types::{
24    AgentBuildDraft, AgentIdentity, AgentRuntimeId, CheckpointVersion, ContinuityGeneration,
25    DurableAgentSpec, FencingToken, SessionSnapshot,
26};
27
28fn is_missing_event_injector_error(error: &str) -> bool {
29    error.contains("missing event injector capability")
30        || (error.contains("missing required capability")
31            && error.contains("interaction_event_injector"))
32}
33
34fn is_missing_bridge_session_snapshot_error(error: &str) -> bool {
35    error.contains("missing bridge session snapshot")
36}
37
38fn is_repairable_bridge_delivery_error(error: &str) -> bool {
39    is_missing_event_injector_error(error)
40        || is_missing_bridge_session_snapshot_error(error)
41        || is_previous_member_cleanup_ambiguous_error(error)
42}
43
44fn is_recoverable_bridge_respawn_cleanup_error(error: &str) -> bool {
45    is_recoverable_lifecycle_cleanup_error(error)
46}
47
48fn is_member_already_exists_error(error: &meerkat_mob::MobError) -> bool {
49    matches!(error, meerkat_mob::MobError::MemberAlreadyExists(_))
50}
51
52#[derive(Debug, PartialEq, Eq)]
53enum MemberRepairRespawnFailure {
54    DegradedTopologyRestore { failed_peer_ids: Vec<String> },
55    RecoverableCleanup,
56    Fatal(String),
57}
58
59fn classify_member_repair_respawn_failure(
60    error: &meerkat_mob::MobRespawnError,
61) -> MemberRepairRespawnFailure {
62    if let Some(failed_peer_ids) = topology_restore_failed_peer_ids(error) {
63        return MemberRepairRespawnFailure::DegradedTopologyRestore { failed_peer_ids };
64    }
65    if is_recoverable_bridge_respawn_cleanup_error(&error.to_string()) {
66        return MemberRepairRespawnFailure::RecoverableCleanup;
67    }
68    MemberRepairRespawnFailure::Fatal(error.to_string())
69}
70
71// ---------------------------------------------------------------------------
72// BridgeError
73// ---------------------------------------------------------------------------
74
75/// Errors from session bridge operations.
76#[derive(Debug)]
77pub enum BridgeError {
78    /// The underlying mob operation failed.
79    Mob(String),
80    /// A required field was missing or invalid.
81    InvalidInput(String),
82}
83
84impl std::fmt::Display for BridgeError {
85    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
86        match self {
87            Self::Mob(msg) => write!(f, "session bridge mob error: {msg}"),
88            Self::InvalidInput(msg) => write!(f, "session bridge invalid input: {msg}"),
89        }
90    }
91}
92
93impl std::error::Error for BridgeError {}
94
95/// Typed reason a requested resume could not reuse the persisted runtime
96/// binding and had to fall back to a fresh member spawn.
97#[derive(Debug, Clone, PartialEq, Eq)]
98pub enum ResumeFallbackReason {
99    /// The persisted session/runtime identity is incompatible with the current
100    /// mob runtime binding.
101    RuntimeIdentityIncompatible { detail: String },
102}
103
104/// Result of attempting to materialize a persisted identity through resume.
105#[derive(Debug, Clone, PartialEq, Eq)]
106pub enum ResumeSessionOutcome {
107    /// The persisted session was resumed as-is.
108    Resumed {
109        session_id: meerkat_core::types::SessionId,
110    },
111    /// Resume was rejected for a typed compatibility reason and a fresh member
112    /// was spawned instead.
113    FreshSpawned {
114        session_id: meerkat_core::types::SessionId,
115        reason: ResumeFallbackReason,
116    },
117}
118
119impl ResumeSessionOutcome {
120    #[must_use]
121    pub fn session_id(&self) -> &meerkat_core::types::SessionId {
122        match self {
123            Self::Resumed { session_id } | Self::FreshSpawned { session_id, .. } => session_id,
124        }
125    }
126
127    #[must_use]
128    pub fn fallback_reason(&self) -> Option<&ResumeFallbackReason> {
129        match self {
130            Self::Resumed { .. } => None,
131            Self::FreshSpawned { reason, .. } => Some(reason),
132        }
133    }
134}
135
136async fn submit_internal_bridge_work(
137    handle: &MobHandle,
138    member_id: &MobAgentIdentity,
139    content: &meerkat_core::ContentInput,
140    injected_context: &[meerkat_core::ContentInput],
141    handling_mode: HandlingMode,
142) -> Result<(), BridgeError> {
143    let entry = handle
144        .get_member(member_id)
145        .await
146        .map_err(|err| BridgeError::Mob(err.to_string()))?
147        .ok_or_else(|| BridgeError::Mob(format!("member not found: {member_id}")))?;
148    // Ask 1: attach ambient memory recall as a separate typed injected-context
149    // body rather than fusing it into the user's message text. WorkSpec carries
150    // it to the StartTurnRequest, where meerkat stamps each entry as the
151    // InjectedContext transcript role (excluded from compaction indexing).
152    let mut spec = WorkSpec::new(content.clone(), WorkOrigin::Internal);
153    if !injected_context.is_empty() {
154        spec = spec.with_injected_context(injected_context.to_vec());
155    }
156    handle
157        .submit_work_with_mode(
158            entry.agent_runtime_id.clone(),
159            entry.fence_token,
160            WorkRef::new(),
161            spec,
162            handling_mode,
163        )
164        .await
165        .map(|_| ())
166        .map_err(|err| BridgeError::Mob(err.to_string()))
167}
168
169// ---------------------------------------------------------------------------
170// SessionBridge trait
171// ---------------------------------------------------------------------------
172
173/// Bridge between the identity-first control plane and the Meerkat session
174/// pipeline. Each method maps an identity-layer operation to its concrete
175/// mob-level counterpart.
176#[async_trait]
177pub trait SessionBridge: Send + Sync {
178    /// Spawn a new mob member for a freshly-created identity.
179    async fn create_session(
180        &self,
181        identity: &AgentIdentity,
182        runtime_id: &AgentRuntimeId,
183        spec: &DurableAgentSpec,
184        draft: &AgentBuildDraft,
185        session_id: &meerkat_core::types::SessionId,
186    ) -> Result<meerkat_core::types::SessionId, BridgeError>;
187
188    /// Resume a mob member from a previously checkpointed snapshot.
189    ///
190    /// The `session_id` comes from the ContinuityRecord — it's the session
191    /// that should be loaded from the session store. When the session store
192    /// has the data, this performs a true resume (conversation history intact).
193    /// When the session is missing, implementations should fall back to a
194    /// fresh spawn.
195    async fn resume_session(
196        &self,
197        identity: &AgentIdentity,
198        runtime_id: &AgentRuntimeId,
199        spec: &DurableAgentSpec,
200        draft: &AgentBuildDraft,
201        session_id: &meerkat_core::types::SessionId,
202        snapshot: &SessionSnapshot,
203    ) -> Result<ResumeSessionOutcome, BridgeError>;
204
205    /// Deliver content to an active mob member.
206    async fn deliver(
207        &self,
208        runtime_id: &AgentRuntimeId,
209        content: &meerkat_core::ContentInput,
210    ) -> Result<meerkat_core::types::SessionId, BridgeError>;
211
212    /// Deliver content to an active mob member using a caller-selected turn
213    /// handling mode. Bridge implementations that do not distinguish modes can
214    /// fall back to ordinary delivery.
215    async fn deliver_with_mode(
216        &self,
217        runtime_id: &AgentRuntimeId,
218        content: &meerkat_core::ContentInput,
219        handling_mode: HandlingMode,
220    ) -> Result<meerkat_core::types::SessionId, BridgeError> {
221        let _ = handling_mode;
222        self.deliver(runtime_id, content).await
223    }
224
225    /// Deliver content plus a separate `injected_context` body (meerkat
226    /// 0.7.12 ask 1: typed ambient injection alongside — not fused into —
227    /// the user's message). Bridges that do not carry injected context fall
228    /// back to plain delivery of the user content, dropping the injection.
229    async fn deliver_with_mode_and_context(
230        &self,
231        runtime_id: &AgentRuntimeId,
232        content: &meerkat_core::ContentInput,
233        injected_context: &[meerkat_core::ContentInput],
234        handling_mode: HandlingMode,
235    ) -> Result<meerkat_core::types::SessionId, BridgeError> {
236        let _ = injected_context;
237        self.deliver_with_mode(runtime_id, content, handling_mode)
238            .await
239    }
240
241    /// Checkpoint the current session state for a mob member.
242    async fn checkpoint_session(
243        &self,
244        runtime_id: &AgentRuntimeId,
245        session_id: &meerkat_core::types::SessionId,
246    ) -> Result<SessionSnapshot, BridgeError>;
247
248    /// Retire a mob member.
249    async fn retire_member(&self, runtime_id: &AgentRuntimeId) -> Result<(), BridgeError>;
250
251    /// Wire two active same-mob members by their concrete runtime IDs.
252    async fn wire_peer(&self, _a: &AgentRuntimeId, _b: &AgentRuntimeId) -> Result<(), BridgeError> {
253        Err(BridgeError::Mob("peer wiring not supported".to_string()))
254    }
255
256    /// Wire many active same-mob member pairs by their concrete runtime IDs.
257    async fn wire_peers_batch(
258        &self,
259        edges: &[(AgentRuntimeId, AgentRuntimeId)],
260    ) -> Result<(), BridgeError> {
261        for (a, b) in edges {
262            self.wire_peer(a, b).await?;
263        }
264        Ok(())
265    }
266
267    /// Return currently materialized same-mob member wires by concrete runtime IDs.
268    async fn current_member_wires(
269        &self,
270    ) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
271        Ok(Vec::new())
272    }
273
274    /// Unwire two active same-mob members by their concrete runtime IDs.
275    async fn unwire_peer(
276        &self,
277        _a: &AgentRuntimeId,
278        _b: &AgentRuntimeId,
279    ) -> Result<(), BridgeError> {
280        Err(BridgeError::Mob("peer unwiring not supported".to_string()))
281    }
282
283    /// Inspect the current execution state of a mob member.
284    async fn inspect_member(
285        &self,
286        _runtime_id: &AgentRuntimeId,
287    ) -> Result<MemberInspection, BridgeError> {
288        Err(BridgeError::Mob("inspect not supported".to_string()))
289    }
290
291    /// Register identity ownership for a concrete bridge session.
292    ///
293    /// Bridges that install a continuity-backed session store use this to
294    /// ensure subsequent Meerkat session saves checkpoint under the durable
295    /// identity/generation/fencing tuple.
296    async fn register_session_runtime_state(
297        &self,
298        _session_id: &meerkat_core::types::SessionId,
299        _identity: &AgentIdentity,
300        _generation: ContinuityGeneration,
301        checkpoint_version: CheckpointVersion,
302        _fencing_token: FencingToken,
303    ) -> Result<CheckpointVersion, BridgeError> {
304        Ok(checkpoint_version)
305    }
306
307    /// Remove identity ownership metadata for a concrete bridge session.
308    ///
309    /// This is used when an identity lifecycle operation aborts after session
310    /// runtime state was registered but before continuity commits.
311    async fn unregister_session_runtime_state(
312        &self,
313        _session_id: &meerkat_core::types::SessionId,
314    ) -> Result<(), BridgeError> {
315        Ok(())
316    }
317}
318
319/// Lightweight inspection of a mob member's current execution state.
320#[derive(Debug, Clone)]
321pub struct MemberInspection {
322    pub output_preview: Option<String>,
323    pub is_final: bool,
324    pub peer_reachable_count: usize,
325}
326
327// ---------------------------------------------------------------------------
328// MobSessionBridge — real implementation backed by MobHandle
329// ---------------------------------------------------------------------------
330
331/// Concrete `SessionBridge` backed by a `MobHandle`.
332///
333/// `AgentRuntimeId` is usually used as the `MobAgentIdentity` at the mob layer. Real
334/// external bindings are the exception: Meerkat's external peer names require
335/// identifier-safe `<mob>/<profile>/<member>` segments, so the bridge maps the
336/// runtime ID to the durable identity for those members.
337pub struct MobSessionBridge {
338    handle: MobHandle,
339    /// Session store used for checkpoint (loading session data to serialize).
340    session_store: Option<Arc<dyn meerkat::SessionStore>>,
341    /// Session service used to project the live effective model for capability checks.
342    session_service: Option<Arc<dyn MobSessionService>>,
343    /// Continuity-backed session store, when installed by the identity-first builder.
344    continuity_session_store: Option<Arc<ContinuitySessionStoreAdapter>>,
345    runtime_members: Arc<tokio::sync::RwLock<HashMap<String, String>>>,
346    runtime_sessions: Arc<tokio::sync::RwLock<HashMap<String, meerkat_core::types::SessionId>>>,
347    /// Lazily-minted ops-owner bridge session for external (peer-only)
348    /// members, used when the mob was created without machine-bound owner
349    /// bridge-session authority. Stable for the bridge lifetime so every
350    /// external member shares one generated operation owner.
351    generated_external_owner_session: std::sync::OnceLock<meerkat_core::types::SessionId>,
352}
353
354impl MobSessionBridge {
355    /// Create a new bridge wrapping the given mob handle.
356    pub fn new(handle: MobHandle) -> Self {
357        Self {
358            handle,
359            session_store: None,
360            session_service: None,
361            continuity_session_store: None,
362            runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
363            runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
364            generated_external_owner_session: std::sync::OnceLock::new(),
365        }
366    }
367
368    /// Create a new bridge with session-service access for live model capability checks.
369    pub fn with_session_service(
370        handle: MobHandle,
371        session_service: Arc<dyn MobSessionService>,
372    ) -> Self {
373        Self {
374            handle,
375            session_store: None,
376            session_service: Some(session_service),
377            continuity_session_store: None,
378            runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
379            runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
380            generated_external_owner_session: std::sync::OnceLock::new(),
381        }
382    }
383
384    /// Create a new bridge with an explicit session store for checkpoint support.
385    pub fn with_session_store(
386        handle: MobHandle,
387        session_store: Arc<dyn meerkat::SessionStore>,
388    ) -> Self {
389        Self {
390            handle,
391            session_store: Some(session_store),
392            session_service: None,
393            continuity_session_store: None,
394            runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
395            runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
396            generated_external_owner_session: std::sync::OnceLock::new(),
397        }
398    }
399
400    /// Create a bridge with checkpoint and live capability support.
401    pub fn with_session_store_and_service(
402        handle: MobHandle,
403        session_store: Arc<dyn meerkat::SessionStore>,
404        session_service: Arc<dyn MobSessionService>,
405    ) -> Self {
406        Self {
407            handle,
408            session_store: Some(session_store),
409            session_service: Some(session_service),
410            continuity_session_store: None,
411            runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
412            runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
413            generated_external_owner_session: std::sync::OnceLock::new(),
414        }
415    }
416
417    /// Create a bridge with an identity-owned continuity session store.
418    pub fn with_continuity_session_store(
419        handle: MobHandle,
420        session_store: Arc<ContinuitySessionStoreAdapter>,
421        session_service: Option<Arc<dyn MobSessionService>>,
422    ) -> Self {
423        Self {
424            handle,
425            session_store: Some(session_store.clone()),
426            session_service,
427            continuity_session_store: Some(session_store),
428            runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
429            runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
430            generated_external_owner_session: std::sync::OnceLock::new(),
431        }
432    }
433
434    async fn remember_runtime_member(
435        &self,
436        runtime_id: &AgentRuntimeId,
437        member_id: &MobAgentIdentity,
438    ) {
439        self.runtime_members.write().await.insert(
440            runtime_id.as_str().to_string(),
441            member_id.as_str().to_string(),
442        );
443    }
444
445    async fn remember_runtime_session(
446        &self,
447        runtime_id: &AgentRuntimeId,
448        session_id: &meerkat_core::types::SessionId,
449    ) {
450        self.runtime_sessions
451            .write()
452            .await
453            .insert(runtime_id.as_str().to_string(), session_id.clone());
454    }
455
456    async fn forget_runtime_member(&self, runtime_id: &AgentRuntimeId) {
457        self.runtime_members
458            .write()
459            .await
460            .remove(runtime_id.as_str());
461        self.runtime_sessions
462            .write()
463            .await
464            .remove(runtime_id.as_str());
465    }
466
467    async fn member_id_for_runtime_id(&self, runtime_id: &AgentRuntimeId) -> MobAgentIdentity {
468        let members = self.runtime_members.read().await;
469        members
470            .get(runtime_id.as_str())
471            .map(|member| MobAgentIdentity::from(member.as_str()))
472            // The recompute fallback must mint the same comms-safe roster id
473            // as the spawn path (meerkat 0.7 MemberCommsName rejects `:`).
474            .unwrap_or_else(|| crate::member_comms_id::mob_member_id(runtime_id.as_str()))
475    }
476
477    async fn runtime_session_id(
478        &self,
479        runtime_id: &AgentRuntimeId,
480    ) -> Option<meerkat_core::types::SessionId> {
481        self.runtime_sessions
482            .read()
483            .await
484            .get(runtime_id.as_str())
485            .cloned()
486    }
487
488    /// Resolve the inline definition profile backing `spec.profile`, used as
489    /// the base when a draft model override must be projected into the typed
490    /// `override_profile` spawn owner. Realm-ref bindings resolve to `None`.
491    fn base_profile_for_spec(&self, spec: &DurableAgentSpec) -> Option<meerkat_mob::Profile> {
492        self.handle
493            .definition()
494            .resolve_inline_profile(&spec.profile)
495            .cloned()
496    }
497
498    async fn resolve_runtime_session_id(
499        &self,
500        runtime_id: &AgentRuntimeId,
501        member_id: &MobAgentIdentity,
502        missing_message: &'static str,
503    ) -> Result<meerkat_core::types::SessionId, BridgeError> {
504        if let Some(session_id) = self.handle.resolve_bridge_session_id(member_id).await {
505            self.remember_runtime_session(runtime_id, &session_id).await;
506            return Ok(session_id);
507        }
508
509        // `get_member` faults (machine command/transport errors) must not be
510        // laundered into "member absent"; surface them to the caller.
511        if member_id.as_str() != runtime_id.as_str()
512            && self
513                .handle
514                .get_member(member_id)
515                .await
516                .map_err(|err| BridgeError::Mob(err.to_string()))?
517                .is_some()
518            && let Some(session_id) = self.runtime_session_id(runtime_id).await
519        {
520            return Ok(session_id);
521        }
522
523        Err(BridgeError::Mob(missing_message.to_string()))
524    }
525
526    async fn repair_member_for_delivery(
527        &self,
528        runtime_id: &AgentRuntimeId,
529        member_id: &MobAgentIdentity,
530        member_entry_before_delivery: Option<(meerkat_mob::ProfileName, BTreeMap<String, String>)>,
531    ) -> Result<(), BridgeError> {
532        match self.handle.respawn(member_id.clone(), None).await {
533            Ok(_) => Ok(()),
534            Err(respawn_err) => match classify_member_repair_respawn_failure(&respawn_err) {
535                MemberRepairRespawnFailure::DegradedTopologyRestore { failed_peer_ids } => {
536                    tracing::warn!(
537                        runtime_id = %runtime_id,
538                        member_id = %member_id,
539                        failed_peer_count = failed_peer_ids.len(),
540                        failed_peer_ids = ?failed_peer_ids,
541                        "identity bridge respawn restored member with isolated peer edges; continuing delivery"
542                    );
543                    // meerkat-mob raises this only after the member/session is live; only peer edges are incomplete.
544                    Ok(())
545                }
546                MemberRepairRespawnFailure::RecoverableCleanup => {
547                    // A `get_member` fault must not be read as "member absent"
548                    // (that would trigger a spurious re-spawn); fail the
549                    // delivery repair instead.
550                    if self
551                        .handle
552                        .get_member(member_id)
553                        .await
554                        .map_err(|err| BridgeError::Mob(err.to_string()))?
555                        .is_none()
556                        && let Some((role, labels)) = member_entry_before_delivery
557                    {
558                        let mut spec = SpawnMemberSpec::new(role, member_id.clone());
559                        if !labels.is_empty() {
560                            spec = spec.with_labels(labels);
561                        }
562                        self.handle
563                            .ensure_member(spec)
564                            .await
565                            .map_err(|e| BridgeError::Mob(e.to_string()))?;
566                    }
567                    Ok(())
568                }
569                MemberRepairRespawnFailure::Fatal(message) => Err(BridgeError::Mob(message)),
570            },
571        }
572    }
573
574    /// Ops-owner bridge session for external (peer-only) member operations.
575    ///
576    /// meerkat 0.7.1 fails external member provisioning closed unless the
577    /// spawn carries a generated owner binding (owner bridge session + ops
578    /// registry; see `MultiBackendProvisioner::provision_member`). Prefer the
579    /// mob's machine-bound owner bridge-session authority when it exists;
580    /// otherwise mint one stable session id for this bridge — the runtime
581    /// adapter creates local session resources for it on demand, exactly like
582    /// meerkat's own external smoke supervisors.
583    fn external_owner_bridge_session_id(&self) -> meerkat_core::types::SessionId {
584        if let Some(authority) = self.handle.owner_bridge_session_lifecycle_authority() {
585            return authority.bridge_session_id;
586        }
587        self.generated_external_owner_session
588            .get_or_init(meerkat_core::types::SessionId::new)
589            .clone()
590    }
591
592    /// Spawn a mob member from a fully-built spec, attaching the generated
593    /// owner context that meerkat 0.7.1 requires for external (peer-only)
594    /// bindings. Session-backed specs spawn through the plain path.
595    async fn spawn_member_spec(
596        &self,
597        spawn_spec: SpawnMemberSpec,
598    ) -> Result<(), meerkat_mob::MobError> {
599        if spawn_spec_requires_generated_owner_context(&spawn_spec) {
600            let owner_session_id = self.external_owner_bridge_session_id();
601            Box::pin(
602                self.handle
603                    .spawn_spec_with_generated_owner_context(spawn_spec, owner_session_id),
604            )
605            .await
606            .map(|_| ())
607        } else {
608            Box::pin(self.handle.spawn_spec(spawn_spec))
609                .await
610                .map(|_| ())
611        }
612    }
613
614    async fn spawn_member_spec_replacing_collision(
615        &self,
616        runtime_id: &AgentRuntimeId,
617        member_id: &MobAgentIdentity,
618        spawn_spec: SpawnMemberSpec,
619    ) -> Result<(), meerkat_mob::MobError> {
620        match self.spawn_member_spec(spawn_spec.clone()).await {
621            Ok(()) => Ok(()),
622            Err(error) if is_member_already_exists_error(&error) => {
623                tracing::warn!(
624                    runtime_id = %runtime_id,
625                    member_id = %member_id,
626                    error = %error,
627                    "fresh-spawn fallback collided with an existing member; retiring and retrying with adopted spec"
628                );
629                match self.handle.retire(member_id.clone()).await {
630                    Ok(()) => {}
631                    Err(err)
632                        if is_recoverable_session_owned_retire_cleanup_error(&err.to_string()) => {}
633                    Err(err) => return Err(err),
634                }
635                self.forget_runtime_member(runtime_id).await;
636                self.spawn_member_spec(spawn_spec).await
637            }
638            Err(error) => Err(error),
639        }
640    }
641}
642
643/// Project meerkat 0.7's tri-state peer connectivity into an inspect-level
644/// reachable count, when the tri-state resolves one.
645///
646/// Only a resolved probe ([`WirePeerConnectivity::Known`]) contributes a
647/// live count. The not-applicable / probe-timed-out arms (and an uncomputed
648/// projection) return `None` so the caller falls back to the machine-owned
649/// wiring degree (`wired_to.len()`) instead of projecting 0 — a freshly
650/// wired member has peers regardless of whether a live probe resolved, and
651/// the sibling console alias surface computes the same wire field from
652/// `wired_to`; the two surfaces must agree.
653fn peer_reachable_count_from_connectivity(
654    connectivity: Option<&meerkat_contracts::WirePeerConnectivity>,
655) -> Option<usize> {
656    match connectivity {
657        Some(meerkat_contracts::WirePeerConnectivity::Known { snapshot }) => {
658            Some(snapshot.reachable_peer_count)
659        }
660        Some(
661            meerkat_contracts::WirePeerConnectivity::NotApplicable
662            | meerkat_contracts::WirePeerConnectivity::ProbeTimedOut,
663        )
664        | None => None,
665    }
666}
667
668/// External (peer-only) member provisioning on meerkat 0.7.1 requires a
669/// machine-minted owner context; plain `spawn_spec` fails closed by design.
670pub(crate) fn spawn_spec_requires_generated_owner_context(spawn_spec: &SpawnMemberSpec) -> bool {
671    matches!(
672        spawn_spec.binding,
673        Some(meerkat_mob::RuntimeBinding::External { .. })
674    )
675}
676
677fn spec_uses_external_binding(spec: &DurableAgentSpec) -> bool {
678    matches!(spec.backend, Some(meerkat_mob::MobBackendKind::External))
679        || matches!(
680            spec.binding.as_ref(),
681            Some(meerkat_contracts::WireRuntimeBinding::External { .. })
682        )
683}
684
685fn member_id_for_spawn_spec(
686    runtime_id: &AgentRuntimeId,
687    spec: &DurableAgentSpec,
688) -> MobAgentIdentity {
689    // meerkat 0.7's `MemberCommsName` is fail-closed: roster member ids must
690    // be identifier-safe (no `:`). MobKit's public alias space — durable
691    // identities like `review:singleton` and runtime ids like
692    // `rt:review:singleton:0` — is unchanged; the roster id is the comms-safe
693    // encoding (identity for already-safe names).
694    if spec_uses_external_binding(spec) {
695        crate::member_comms_id::mob_member_id(spec.identity.as_str())
696    } else {
697        crate::member_comms_id::mob_member_id(runtime_id.as_str())
698    }
699}
700
701/// Build a `SpawnMemberSpec` from identity-first types, wiring draft fields.
702///
703/// `base_profile` is the resolved definition profile for `spec.profile`; it is
704/// only needed when the draft carries a model override (meerkat 0.7 removed
705/// `SpawnMemberSpec::model_override` in favor of the typed `override_profile`
706/// owner, so a model override is expressed as the role profile with the model
707/// swapped).
708pub(crate) fn build_spawn_spec(
709    runtime_id: &AgentRuntimeId,
710    spec: &DurableAgentSpec,
711    draft: &AgentBuildDraft,
712    base_profile: Option<&meerkat_mob::Profile>,
713) -> SpawnMemberSpec {
714    let mid = member_id_for_spawn_spec(runtime_id, spec);
715    let mut spawn_spec = SpawnMemberSpec::new(spec.profile.clone(), mid);
716
717    if let Some(message) = spec.initial_message.as_ref() {
718        spawn_spec = spawn_spec.with_initial_message(message.clone());
719    }
720    if let Some(runtime_mode) = spec.runtime_mode_override {
721        spawn_spec = spawn_spec.with_runtime_mode(runtime_mode);
722    }
723    spawn_spec.backend = spec.backend;
724    if let Some(binding) = spec.binding.clone() {
725        spawn_spec.binding = runtime_binding_from_wire(binding);
726    }
727    if let Some(ref ctx) = draft.app_context {
728        spawn_spec = spawn_spec.with_context(ctx.clone());
729    }
730    let mut labels = draft.labels.clone();
731    labels.insert(
732        "agent_identity".to_string(),
733        spec.identity.as_str().to_string(),
734    );
735    labels.insert(
736        "profile_name".to_string(),
737        spec.profile.as_str().to_string(),
738    );
739    if !labels.is_empty() {
740        spawn_spec = spawn_spec.with_labels(labels);
741    }
742    if !draft.additional_instructions.is_empty() {
743        spawn_spec = spawn_spec.with_additional_instructions(draft.additional_instructions.clone());
744    }
745    if let Some(model) = draft.model.as_ref() {
746        match base_profile {
747            Some(base) => {
748                let mut profile = base.clone();
749                profile.model = model.clone();
750                // The base profile's pinned provider (and self-hosted server
751                // binding) belongs to its original model id; clear it so the
752                // catalog re-infers the provider for the overridden model.
753                profile.provider = None;
754                profile.self_hosted_server_id = None;
755                spawn_spec.override_profile = Some(profile);
756            }
757            None => {
758                // No resolvable inline base profile (realm-ref binding or
759                // unknown role). Leave `override_profile` unset so the spawn
760                // resolves through the definition's canonical path; the model
761                // override cannot be applied without a base profile.
762                tracing::warn!(
763                    identity = %spec.identity,
764                    profile = %spec.profile,
765                    model = %model,
766                    "model override skipped: role profile is not an inline definition profile"
767                );
768            }
769        }
770    }
771    if let Some(system_prompt) = draft.system_prompt.as_ref() {
772        spawn_spec.system_prompt_override =
773            Some(SpawnSystemPromptOverride::Replace(system_prompt.clone()));
774    }
775    if let Some(dispatcher) = draft.local_external_tools.dispatcher() {
776        spawn_spec.external_tools = Some(dispatcher);
777    }
778
779    spawn_spec
780}
781
782fn runtime_binding_from_wire(
783    binding: meerkat_contracts::WireRuntimeBinding,
784) -> Option<meerkat_mob::RuntimeBinding> {
785    match binding {
786        meerkat_contracts::WireRuntimeBinding::Session => {
787            Some(meerkat_mob::RuntimeBinding::Session)
788        }
789        meerkat_contracts::WireRuntimeBinding::External {
790            address,
791            bootstrap_token,
792            identity,
793        } => {
794            let resolved = identity.resolve().ok()?;
795            Some(meerkat_mob::RuntimeBinding::External {
796                peer_id: resolved.peer_id.to_string(),
797                address,
798                bootstrap_token,
799                pubkey: resolved.pubkey,
800            })
801        }
802    }
803}
804
805#[async_trait]
806impl SessionBridge for MobSessionBridge {
807    async fn create_session(
808        &self,
809        _identity: &AgentIdentity,
810        runtime_id: &AgentRuntimeId,
811        spec: &DurableAgentSpec,
812        draft: &AgentBuildDraft,
813        session_id: &meerkat_core::types::SessionId,
814    ) -> Result<meerkat_core::types::SessionId, BridgeError> {
815        let mid = member_id_for_spawn_spec(runtime_id, spec);
816        let spawn_spec = build_spawn_spec(
817            runtime_id,
818            spec,
819            draft,
820            self.base_profile_for_spec(spec).as_ref(),
821        );
822
823        self.spawn_member_spec(spawn_spec)
824            .await
825            .map_err(|e| BridgeError::Mob(e.to_string()))?;
826        self.remember_runtime_member(runtime_id, &mid).await;
827        self.remember_runtime_session(runtime_id, session_id).await;
828
829        self.resolve_runtime_session_id(runtime_id, &mid, "member spawned but has no session ID")
830            .await
831    }
832
833    async fn resume_session(
834        &self,
835        _identity: &AgentIdentity,
836        runtime_id: &AgentRuntimeId,
837        spec: &DurableAgentSpec,
838        draft: &AgentBuildDraft,
839        session_id: &meerkat_core::types::SessionId,
840        _snapshot: &SessionSnapshot,
841    ) -> Result<ResumeSessionOutcome, BridgeError> {
842        if spec_uses_external_binding(spec) {
843            let mut spawn_spec = build_spawn_spec(
844                runtime_id,
845                spec,
846                draft,
847                self.base_profile_for_spec(spec).as_ref(),
848            );
849            spawn_spec.launch_mode = MemberLaunchMode::Resume {
850                bridge_session_id: session_id.clone(),
851            };
852            let mid = member_id_for_spawn_spec(runtime_id, spec);
853            self.spawn_member_spec(spawn_spec)
854                .await
855                .map_err(|e| BridgeError::Mob(e.to_string()))?;
856            self.remember_runtime_member(runtime_id, &mid).await;
857            self.remember_runtime_session(runtime_id, session_id).await;
858            return Ok(ResumeSessionOutcome::Resumed {
859                session_id: session_id.clone(),
860            });
861        }
862
863        // Try MemberLaunchMode::Resume first — this loads the existing session
864        // from the session store (conversation history intact).
865        let mut spawn_spec = build_spawn_spec(
866            runtime_id,
867            spec,
868            draft,
869            self.base_profile_for_spec(spec).as_ref(),
870        );
871        spawn_spec.launch_mode = MemberLaunchMode::Resume {
872            bridge_session_id: session_id.clone(),
873        };
874
875        let mid = member_id_for_spawn_spec(runtime_id, spec);
876
877        match self.spawn_member_spec(spawn_spec).await {
878            Ok(()) => {
879                self.remember_runtime_member(runtime_id, &mid).await;
880                self.remember_runtime_session(runtime_id, session_id).await;
881                Ok(ResumeSessionOutcome::Resumed {
882                    session_id: session_id.clone(),
883                })
884            }
885            Err(e) => {
886                // Resume can fail if the old session's comms identity is still
887                // claimed (e.g., in-process restart where the previous mob actor
888                // hasn't fully terminated). Fall back to a fresh spawn.
889                tracing::warn!(
890                    identity = %_identity,
891                    session_id = %session_id,
892                    error = %e,
893                    reason = "runtime_identity_incompatible",
894                    "resume_session incompatible with current runtime binding, falling back to fresh spawn"
895                );
896                let fresh_spec = build_spawn_spec(
897                    runtime_id,
898                    spec,
899                    draft,
900                    self.base_profile_for_spec(spec).as_ref(),
901                );
902                self.spawn_member_spec_replacing_collision(runtime_id, &mid, fresh_spec)
903                    .await
904                    .map_err(|e2| BridgeError::Mob(e2.to_string()))?;
905
906                self.remember_runtime_member(runtime_id, &mid).await;
907                let session_id = self
908                    .resolve_runtime_session_id(
909                        runtime_id,
910                        &mid,
911                        "member spawned (fresh fallback) but has no session ID",
912                    )
913                    .await?;
914                Ok(ResumeSessionOutcome::FreshSpawned {
915                    session_id,
916                    reason: ResumeFallbackReason::RuntimeIdentityIncompatible {
917                        detail: e.to_string(),
918                    },
919                })
920            }
921        }
922    }
923
924    async fn deliver(
925        &self,
926        runtime_id: &AgentRuntimeId,
927        content: &meerkat_core::ContentInput,
928    ) -> Result<meerkat_core::types::SessionId, BridgeError> {
929        let mid = self.member_id_for_runtime_id(runtime_id).await;
930        // Best-effort repair material: a faulted lookup degrades to "no
931        // pre-delivery entry" (the delivery itself will surface the fault).
932        let member_entry_before_delivery = self
933            .handle
934            .get_member(&mid)
935            .await
936            .ok()
937            .flatten()
938            .map(|entry| (entry.role, entry.labels));
939        if content_input_has_images(content) {
940            let member_entry = self
941                .handle
942                .get_member(&mid)
943                .await
944                .map_err(|err| BridgeError::Mob(err.to_string()))?
945                .ok_or_else(|| {
946                    BridgeError::Mob("member not found while checking image capability".to_string())
947                })?;
948            let caps = model_capabilities_for_member(
949                &self.handle,
950                self.session_service.as_ref(),
951                &member_entry.agent_identity,
952            )
953            .await;
954            if !caps.image_input {
955                return Err(BridgeError::InvalidInput(
956                    "target member model cannot accept image input".to_string(),
957                ));
958            }
959        }
960
961        // Submit internal work directly through the mob work lane so delivery
962        // acks at runtime ingress rather than waiting for the full turn to
963        // complete. The identity layer owns addressability enforcement — the
964        // bridge is an internal delivery mechanism regardless of whether the
965        // identity is Addressable or InternalOnly.
966        match submit_internal_bridge_work(&self.handle, &mid, content, &[], HandlingMode::Queue)
967            .await
968        {
969            Ok(()) => {}
970            Err(err) if is_repairable_bridge_delivery_error(&err.to_string()) => {
971                tracing::warn!(
972                    runtime_id = %runtime_id,
973                    error = %err,
974                    "identity bridge delivery found stale runtime state; repairing member before retry"
975                );
976                Box::pin(self.repair_member_for_delivery(
977                    runtime_id,
978                    &mid,
979                    member_entry_before_delivery,
980                ))
981                .await?;
982                submit_internal_bridge_work(&self.handle, &mid, content, &[], HandlingMode::Queue)
983                    .await?;
984            }
985            Err(err) => return Err(BridgeError::Mob(err.to_string())),
986        }
987
988        // Meerkat 0.6: MemberDeliveryReceipt no longer carries session_id.
989        // Query the bridge session id directly from the mob handle.
990        self.resolve_runtime_session_id(
991            runtime_id,
992            &mid,
993            "member has no bridge session after deliver",
994        )
995        .await
996    }
997
998    async fn deliver_with_mode(
999        &self,
1000        runtime_id: &AgentRuntimeId,
1001        content: &meerkat_core::ContentInput,
1002        handling_mode: HandlingMode,
1003    ) -> Result<meerkat_core::types::SessionId, BridgeError> {
1004        self.deliver_with_mode_and_context(runtime_id, content, &[], handling_mode)
1005            .await
1006    }
1007
1008    async fn deliver_with_mode_and_context(
1009        &self,
1010        runtime_id: &AgentRuntimeId,
1011        content: &meerkat_core::ContentInput,
1012        injected_context: &[meerkat_core::ContentInput],
1013        handling_mode: HandlingMode,
1014    ) -> Result<meerkat_core::types::SessionId, BridgeError> {
1015        let mid = self.member_id_for_runtime_id(runtime_id).await;
1016        // Best-effort repair material: a faulted lookup degrades to "no
1017        // pre-delivery entry" (the delivery itself will surface the fault).
1018        let member_entry_before_delivery = self
1019            .handle
1020            .get_member(&mid)
1021            .await
1022            .ok()
1023            .flatten()
1024            .map(|entry| (entry.role, entry.labels));
1025        if content_input_has_images(content) {
1026            let member_entry = self
1027                .handle
1028                .get_member(&mid)
1029                .await
1030                .map_err(|err| BridgeError::Mob(err.to_string()))?
1031                .ok_or_else(|| {
1032                    BridgeError::Mob("member not found while checking image capability".to_string())
1033                })?;
1034            let caps = model_capabilities_for_member(
1035                &self.handle,
1036                self.session_service.as_ref(),
1037                &member_entry.agent_identity,
1038            )
1039            .await;
1040            if !caps.image_input {
1041                return Err(BridgeError::InvalidInput(
1042                    "target member model cannot accept image input".to_string(),
1043                ));
1044            }
1045        }
1046
1047        match submit_internal_bridge_work(
1048            &self.handle,
1049            &mid,
1050            content,
1051            injected_context,
1052            handling_mode,
1053        )
1054        .await
1055        {
1056            Ok(()) => {}
1057            Err(err) if is_repairable_bridge_delivery_error(&err.to_string()) => {
1058                tracing::warn!(
1059                    runtime_id = %runtime_id,
1060                    error = %err,
1061                    "identity bridge delivery found stale runtime state; repairing member before retry"
1062                );
1063                Box::pin(self.repair_member_for_delivery(
1064                    runtime_id,
1065                    &mid,
1066                    member_entry_before_delivery,
1067                ))
1068                .await?;
1069                submit_internal_bridge_work(
1070                    &self.handle,
1071                    &mid,
1072                    content,
1073                    injected_context,
1074                    handling_mode,
1075                )
1076                .await?;
1077            }
1078            Err(err) => return Err(BridgeError::Mob(err.to_string())),
1079        }
1080
1081        self.resolve_runtime_session_id(
1082            runtime_id,
1083            &mid,
1084            "member has no bridge session after deliver",
1085        )
1086        .await
1087    }
1088
1089    async fn checkpoint_session(
1090        &self,
1091        _runtime_id: &AgentRuntimeId,
1092        session_id: &meerkat_core::types::SessionId,
1093    ) -> Result<SessionSnapshot, BridgeError> {
1094        let store = self.session_store.as_ref().ok_or_else(|| {
1095            BridgeError::InvalidInput(
1096                "checkpoint requires a session store but none was configured".to_string(),
1097            )
1098        })?;
1099
1100        let session = store
1101            .load(session_id)
1102            .await
1103            .map_err(|e| BridgeError::Mob(format!("failed to load session for checkpoint: {e}")))?
1104            .ok_or_else(|| {
1105                BridgeError::Mob(format!(
1106                    "session {session_id} not found in store for checkpoint"
1107                ))
1108            })?;
1109
1110        let data = serde_json::to_vec(&session)
1111            .map_err(|e| BridgeError::Mob(format!("failed to serialize session: {e}")))?;
1112
1113        Ok(SessionSnapshot { data })
1114    }
1115
1116    async fn retire_member(&self, runtime_id: &AgentRuntimeId) -> Result<(), BridgeError> {
1117        let mid = self.member_id_for_runtime_id(runtime_id).await;
1118        match self.handle.retire(mid).await {
1119            Ok(()) => {
1120                self.forget_runtime_member(runtime_id).await;
1121                Ok(())
1122            }
1123            // All callers of `retire_member` are identity-first session-owned
1124            // agents, so a mob-archive miss (NotFound for a registered runtime
1125            // session) is the expected outcome of disposing one — not an
1126            // orphan. Tolerate it so reset/delete_identity complete instead of
1127            // bricking the identity until a process restart.
1128            Err(err) if is_recoverable_session_owned_retire_cleanup_error(&err.to_string()) => {
1129                // Disposal completed and the member left the roster, so the
1130                // runtime-member mapping is stale — forget it (matching the
1131                // success path) instead of leaking one entry per tolerated
1132                // retire across the process lifetime.
1133                self.forget_runtime_member(runtime_id).await;
1134                Ok(())
1135            }
1136            Err(err) => Err(BridgeError::Mob(err.to_string())),
1137        }
1138    }
1139
1140    async fn wire_peer(&self, a: &AgentRuntimeId, b: &AgentRuntimeId) -> Result<(), BridgeError> {
1141        let member_a = self.member_id_for_runtime_id(a).await;
1142        let member_b = self.member_id_for_runtime_id(b).await;
1143        self.handle
1144            .wire(
1145                meerkat_mob::AgentIdentity::from(member_a.as_str()),
1146                member_b,
1147            )
1148            .await
1149            .map_err(|e| BridgeError::Mob(e.to_string()))
1150    }
1151
1152    async fn wire_peers_batch(
1153        &self,
1154        edges: &[(AgentRuntimeId, AgentRuntimeId)],
1155    ) -> Result<(), BridgeError> {
1156        let mut member_edges = Vec::with_capacity(edges.len());
1157        for (a, b) in edges {
1158            let member_a = self.member_id_for_runtime_id(a).await;
1159            let member_b = self.member_id_for_runtime_id(b).await;
1160            member_edges.push((
1161                meerkat_mob::AgentIdentity::from(member_a.as_str()),
1162                meerkat_mob::AgentIdentity::from(member_b.as_str()),
1163            ));
1164        }
1165        self.handle
1166            .wire_members_batch(member_edges)
1167            .await
1168            .map(|_| ())
1169            .map_err(|e| BridgeError::Mob(e.to_string()))
1170    }
1171
1172    async fn current_member_wires(
1173        &self,
1174    ) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
1175        let members = self.handle.list_members_including_retiring().await;
1176        let runtime_members = self.runtime_members.read().await;
1177        let member_runtimes = runtime_members
1178            .iter()
1179            .map(|(runtime, member)| (member.clone(), runtime.clone()))
1180            .collect::<HashMap<_, _>>();
1181        let active_ids = members
1182            .iter()
1183            .map(|member| member.agent_identity.to_string())
1184            .collect::<std::collections::BTreeSet<_>>();
1185        let mut edges = std::collections::BTreeSet::new();
1186        for member in &members {
1187            let a = member.agent_identity.to_string();
1188            for peer in &member.wired_to {
1189                let b = peer.to_string();
1190                if !active_ids.contains(&b) {
1191                    continue;
1192                }
1193                let key = if a <= b {
1194                    (a.clone(), b)
1195                } else {
1196                    (b, a.clone())
1197                };
1198                edges.insert(key);
1199            }
1200        }
1201        Ok(edges
1202            .into_iter()
1203            .filter_map(|(a, b)| {
1204                // Fallback for members not in the in-memory map: the roster id
1205                // is the comms-safe encoding of the runtime alias; decode it.
1206                let a = member_runtimes
1207                    .get(&a)
1208                    .cloned()
1209                    .unwrap_or_else(|| crate::member_comms_id::runtime_alias_str(&a).into_owned());
1210                let b = member_runtimes
1211                    .get(&b)
1212                    .cloned()
1213                    .unwrap_or_else(|| crate::member_comms_id::runtime_alias_str(&b).into_owned());
1214                Some((
1215                    AgentRuntimeId::parse(&a).ok()?,
1216                    AgentRuntimeId::parse(&b).ok()?,
1217                ))
1218            })
1219            .collect())
1220    }
1221
1222    async fn unwire_peer(&self, a: &AgentRuntimeId, b: &AgentRuntimeId) -> Result<(), BridgeError> {
1223        let member_a = self.member_id_for_runtime_id(a).await;
1224        let member_b = self.member_id_for_runtime_id(b).await;
1225        match self
1226            .handle
1227            .unwire(
1228                meerkat_mob::AgentIdentity::from(member_a.as_str()),
1229                member_b,
1230            )
1231            .await
1232        {
1233            Ok(()) => Ok(()),
1234            Err(err) => {
1235                let message = err.to_string();
1236                if message.contains("peer not found") || message.contains("not wired") {
1237                    Ok(())
1238                } else {
1239                    Err(BridgeError::Mob(message))
1240                }
1241            }
1242        }
1243    }
1244
1245    async fn inspect_member(
1246        &self,
1247        runtime_id: &AgentRuntimeId,
1248    ) -> Result<MemberInspection, BridgeError> {
1249        let mid = self.member_id_for_runtime_id(runtime_id).await;
1250        let snap = self
1251            .handle
1252            .member_status(&mid)
1253            .await
1254            .map_err(|e| BridgeError::Mob(e.to_string()))?;
1255        let peer_reachable_count =
1256            match peer_reachable_count_from_connectivity(snap.peer_connectivity.as_ref()) {
1257                Some(count) => count,
1258                None => self
1259                    .handle
1260                    .get_member(&mid)
1261                    .await
1262                    .ok()
1263                    .flatten()
1264                    .map(|entry| entry.wired_to.len())
1265                    .unwrap_or(0),
1266            };
1267        Ok(MemberInspection {
1268            output_preview: snap.output_preview.clone(),
1269            is_final: snap.is_final,
1270            peer_reachable_count,
1271        })
1272    }
1273
1274    async fn register_session_runtime_state(
1275        &self,
1276        session_id: &meerkat_core::types::SessionId,
1277        identity: &AgentIdentity,
1278        generation: ContinuityGeneration,
1279        checkpoint_version: CheckpointVersion,
1280        fencing_token: FencingToken,
1281    ) -> Result<CheckpointVersion, BridgeError> {
1282        if let Some(adapter) = self.continuity_session_store.as_ref() {
1283            return adapter
1284                .register_session(
1285                    session_id,
1286                    SessionRuntimeState {
1287                        identity: identity.clone(),
1288                        generation,
1289                        checkpoint_version,
1290                        fencing_token,
1291                    },
1292                )
1293                .await
1294                .map_err(|err| BridgeError::Mob(format!("continuity register_session: {err}")));
1295        }
1296        Ok(checkpoint_version)
1297    }
1298
1299    async fn unregister_session_runtime_state(
1300        &self,
1301        session_id: &meerkat_core::types::SessionId,
1302    ) -> Result<(), BridgeError> {
1303        if let Some(adapter) = self.continuity_session_store.as_ref() {
1304            adapter
1305                .unregister_session(session_id)
1306                .await
1307                .map_err(|err| BridgeError::Mob(format!("continuity unregister_session: {err}")))?;
1308        }
1309        Ok(())
1310    }
1311}
1312
1313#[cfg(test)]
1314#[allow(clippy::expect_used)]
1315mod tests {
1316    use std::sync::Arc;
1317
1318    use async_trait::async_trait;
1319    use meerkat_core::agent::AgentToolDispatcher;
1320    use meerkat_core::types::ToolCallView;
1321    use meerkat_core::{ToolDef, error::ToolError, ops::ToolDispatchOutcome};
1322    use meerkat_mob::{MobRespawnError, MobRuntimeMode};
1323
1324    use super::*;
1325    use crate::identity_first::{AgentAddressability, LocalExternalToolOverlay};
1326
1327    struct EmptyDispatcher;
1328
1329    #[async_trait]
1330    impl AgentToolDispatcher for EmptyDispatcher {
1331        fn tools(&self) -> Arc<[Arc<ToolDef>]> {
1332            Arc::from([])
1333        }
1334
1335        async fn dispatch(
1336            &self,
1337            _call: ToolCallView<'_>,
1338        ) -> Result<ToolDispatchOutcome, ToolError> {
1339            Err(ToolError::ExecutionFailed {
1340                message: "not implemented".to_string(),
1341            })
1342        }
1343    }
1344
1345    fn durable_spec() -> DurableAgentSpec {
1346        DurableAgentSpec {
1347            identity: AgentIdentity::parse("agent:alpha").expect("identity"),
1348            profile: meerkat_mob::ProfileName::from("worker"),
1349            addressability: AgentAddressability::Addressable,
1350            display_name: None,
1351            labels: Default::default(),
1352            context: None,
1353            additional_instructions: Vec::new(),
1354            initial_message: Some(meerkat_core::ContentInput::Text("hello".to_string())),
1355            runtime_mode_override: Some(MobRuntimeMode::TurnDriven),
1356            backend: None,
1357            binding: None,
1358        }
1359    }
1360
1361    /// Regression: meerkat 0.7.1 made `member_status` peer connectivity a
1362    /// tri-state. Only the `Known` arm carries a live count; the
1363    /// not-applicable / probe-timed-out arms must defer to the machine-owned
1364    /// wiring degree (`wired_to.len()`) instead of projecting 0, or
1365    /// `peer_reachable_count` reads 0 for freshly role-wired members and
1366    /// topology verification breaks.
1367    #[test]
1368    fn peer_reachable_count_tri_state_defers_to_wiring_when_unresolved() {
1369        use meerkat_contracts::{WirePeerConnectivity, WirePeerConnectivitySnapshot};
1370
1371        let known = WirePeerConnectivity::Known {
1372            snapshot: WirePeerConnectivitySnapshot {
1373                reachable_peer_count: 3,
1374                unknown_peer_count: 0,
1375                unreachable_peers: Vec::new(),
1376            },
1377        };
1378        assert_eq!(
1379            peer_reachable_count_from_connectivity(Some(&known)),
1380            Some(3),
1381            "a resolved probe owns the count"
1382        );
1383        assert_eq!(
1384            peer_reachable_count_from_connectivity(Some(&WirePeerConnectivity::NotApplicable)),
1385            None,
1386            "not-applicable must defer to the wiring fallback"
1387        );
1388        assert_eq!(
1389            peer_reachable_count_from_connectivity(Some(&WirePeerConnectivity::ProbeTimedOut)),
1390            None,
1391            "probe timeout must defer to the wiring fallback"
1392        );
1393        assert_eq!(peer_reachable_count_from_connectivity(None), None);
1394    }
1395
1396    #[test]
1397    fn build_spawn_spec_maps_identity_first_overrides() {
1398        let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
1399        let mut labels = std::collections::BTreeMap::new();
1400        labels.insert("team".to_string(), "ops".to_string());
1401        let draft = AgentBuildDraft {
1402            model: Some("gpt-test".to_string()),
1403            system_prompt: Some("system override".to_string()),
1404            additional_instructions: vec!["stay focused".to_string()],
1405            labels,
1406            app_context: Some(serde_json::json!({"ticket": 7})),
1407            external_tools: Vec::new(),
1408            local_external_tools: LocalExternalToolOverlay::new(Arc::new(EmptyDispatcher)),
1409        };
1410
1411        let base_profile: meerkat_mob::Profile =
1412            serde_json::from_value(serde_json::json!({"model": "base-model"}))
1413                .expect("minimal profile");
1414        let spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, Some(&base_profile));
1415
1416        assert_eq!(
1417            spawn
1418                .override_profile
1419                .as_ref()
1420                .map(|profile| profile.model.as_str()),
1421            Some("gpt-test")
1422        );
1423        assert_eq!(
1424            spawn.system_prompt_override,
1425            Some(SpawnSystemPromptOverride::Replace(
1426                "system override".to_string()
1427            ))
1428        );
1429        assert!(spawn.external_tools.is_some());
1430        assert_eq!(spawn.runtime_mode, Some(MobRuntimeMode::TurnDriven));
1431        assert_eq!(
1432            spawn.initial_message,
1433            Some(meerkat_core::ContentInput::Text("hello".to_string()))
1434        );
1435        assert_eq!(
1436            spawn
1437                .labels
1438                .as_ref()
1439                .and_then(|labels| labels.get("team"))
1440                .map(String::as_str),
1441            Some("ops")
1442        );
1443        assert_eq!(
1444            spawn
1445                .labels
1446                .as_ref()
1447                .and_then(|labels| labels.get("agent_identity"))
1448                .map(String::as_str),
1449            Some("agent:alpha")
1450        );
1451        assert_eq!(
1452            spawn
1453                .labels
1454                .as_ref()
1455                .and_then(|labels| labels.get("profile_name"))
1456                .map(String::as_str),
1457            Some("worker"),
1458            "identity-first spawn labels must carry the adopted profile so the \
1459             SDK build callback sees the roster profile, not a checkpoint default"
1460        );
1461        assert_eq!(
1462            spawn.role_name.as_str(),
1463            "worker",
1464            "SpawnMemberSpec role remains the authoritative mob profile"
1465        );
1466    }
1467
1468    #[test]
1469    fn fresh_fallback_collision_classifier_matches_member_already_exists() {
1470        let error = meerkat_mob::MobError::MemberAlreadyExists(meerkat_mob::AgentIdentity::from(
1471            "rt-agent-alpha-0",
1472        ));
1473
1474        assert!(
1475            is_member_already_exists_error(&error),
1476            "fresh fallback must retry recreate-over-running-member collisions"
1477        );
1478    }
1479
1480    #[test]
1481    fn build_spawn_spec_maps_remote_runtime_binding() {
1482        let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
1483        let mut spec = durable_spec();
1484        spec.backend = Some(meerkat_mob::MobBackendKind::External);
1485        spec.binding = Some(
1486            serde_json::from_value(serde_json::json!({
1487                "kind": "external",
1488                "address": "tcp://127.0.0.1:4777",
1489                "identity": {
1490                    "kind": "ed25519_public_key",
1491                    "public_key": "ed25519:BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc="
1492                }
1493            }))
1494            .expect("wire binding"),
1495        );
1496        let draft = AgentBuildDraft {
1497            model: None,
1498            system_prompt: None,
1499            additional_instructions: Vec::new(),
1500            labels: Default::default(),
1501            app_context: None,
1502            external_tools: Vec::new(),
1503            local_external_tools: Default::default(),
1504        };
1505
1506        let spawn = build_spawn_spec(&runtime_id, &spec, &draft, None);
1507
1508        // meerkat 0.7: MemberCommsName is fail-closed (no `:` in member-id
1509        // components), so the roster id is the comms-safe encoding of the
1510        // public identity `agent:alpha` (see crate::member_comms_id).
1511        assert_eq!(
1512            spawn.identity.as_str(),
1513            crate::member_comms_id::mob_member_id_str("agent:alpha").as_ref()
1514        );
1515        assert_eq!(spawn.backend, Some(meerkat_mob::MobBackendKind::External));
1516        assert!(
1517            matches!(
1518                spawn.binding,
1519                Some(meerkat_mob::RuntimeBinding::External { .. })
1520            ),
1521            "expected external runtime binding, got {:?}",
1522            spawn.binding
1523        );
1524        if let Some(meerkat_mob::RuntimeBinding::External {
1525            address, pubkey, ..
1526        }) = spawn.binding
1527        {
1528            assert_eq!(address.as_str(), "tcp://127.0.0.1:4777");
1529            assert_eq!(pubkey, [7; 32]);
1530        }
1531    }
1532
1533    #[test]
1534    fn external_binding_spawn_specs_require_generated_owner_context() {
1535        let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
1536        let draft = AgentBuildDraft {
1537            model: None,
1538            system_prompt: None,
1539            additional_instructions: Vec::new(),
1540            labels: Default::default(),
1541            app_context: None,
1542            external_tools: Vec::new(),
1543            local_external_tools: Default::default(),
1544        };
1545
1546        // Session-backed members keep the plain spawn path.
1547        let session_spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, None);
1548        assert!(
1549            !spawn_spec_requires_generated_owner_context(&session_spawn),
1550            "session-backed spawns must not require a generated owner context"
1551        );
1552
1553        // meerkat 0.7.1 MultiBackendProvisioner::provision_member fails
1554        // external (peer-only) members closed without a generated owner
1555        // binding; the bridge must route them through
1556        // spawn_spec_with_generated_owner_context.
1557        let mut spec = durable_spec();
1558        spec.backend = Some(meerkat_mob::MobBackendKind::External);
1559        spec.binding = Some(
1560            serde_json::from_value(serde_json::json!({
1561                "kind": "external",
1562                "address": "tcp://127.0.0.1:4777",
1563                "identity": {
1564                    "kind": "ed25519_public_key",
1565                    "public_key": "ed25519:BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc="
1566                }
1567            }))
1568            .expect("wire binding"),
1569        );
1570        let external_spawn = build_spawn_spec(&runtime_id, &spec, &draft, None);
1571        assert!(
1572            spawn_spec_requires_generated_owner_context(&external_spawn),
1573            "external peer-only spawns must carry a generated owner binding on meerkat 0.7.1"
1574        );
1575    }
1576
1577    #[test]
1578    fn bridge_delivery_repair_covers_missing_bridge_session_snapshot() {
1579        let error = "session bridge mob error: member rt:review:singleton:0 failed to restore session 019e5fc2-dad4-77e2-abbe-a8a66bc15f66: missing bridge session snapshot for '019e5fc2-dad4-77e2-abbe-a8a66bc15f66'";
1580
1581        assert!(
1582            is_repairable_bridge_delivery_error(error),
1583            "stale bridge-session bindings should be repaired before retrying delivery"
1584        );
1585        assert!(
1586            is_repairable_bridge_delivery_error("missing event injector capability for member"),
1587            "existing stale event-injector repair path must remain covered"
1588        );
1589        assert!(
1590            is_repairable_bridge_delivery_error(
1591                "mob member rt:us-president:0 missing required capability interaction_event_injector: autonomous member dispatch"
1592            ),
1593            "newer autonomous member dispatch wording should repair and retry instead of dropping the event"
1594        );
1595        assert!(
1596            is_repairable_bridge_delivery_error(
1597                "previous member cleanup ambiguous for member rt:deep-investigator:singleton:0"
1598            ),
1599            "ambiguous Meerkat respawn cleanup should trigger bridge repair instead of failing delivery"
1600        );
1601        assert!(
1602            !is_repairable_bridge_delivery_error("model provider returned rate limit"),
1603            "ordinary turn failures must not trigger member repair"
1604        );
1605    }
1606
1607    #[test]
1608    fn bridge_delivery_repair_classifies_topology_restore_failure_as_degraded() {
1609        let identity = meerkat_mob::AgentIdentity::from("rt:review:singleton:0");
1610        let receipt = meerkat_mob::MemberRespawnReceipt::new(
1611            identity.clone(),
1612            meerkat_mob::AgentRuntimeId::new(identity, meerkat_mob::ids::Generation::INITIAL),
1613            meerkat_mob::FenceToken::new(1),
1614            meerkat_mob::FenceToken::new(2),
1615        );
1616        let err = MobRespawnError::TopologyRestoreFailed {
1617            receipt,
1618            failed_peer_ids: vec![meerkat_mob::RespawnTopologyPeerId::from(
1619                "initiative:broken",
1620            )],
1621        };
1622
1623        assert_eq!(
1624            classify_member_repair_respawn_failure(&err),
1625            MemberRepairRespawnFailure::DegradedTopologyRestore {
1626                failed_peer_ids: vec!["initiative:broken".to_string()]
1627            },
1628            "failed peer edges should degrade bridge repair instead of bricking delivery"
1629        );
1630        assert!(
1631            matches!(
1632                classify_member_repair_respawn_failure(&MobRespawnError::NoRuntimeControl {
1633                    identity: meerkat_mob::AgentIdentity::from("rt:review:singleton:0"),
1634                }),
1635                MemberRepairRespawnFailure::Fatal(_)
1636            ),
1637            "ordinary respawn failures must still fail bridge repair"
1638        );
1639    }
1640}