1use 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, model_capabilities_for_member,
19 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
48#[derive(Debug, PartialEq, Eq)]
49enum MemberRepairRespawnFailure {
50 DegradedTopologyRestore { failed_peer_ids: Vec<String> },
51 RecoverableCleanup,
52 Fatal(String),
53}
54
55fn classify_member_repair_respawn_failure(
56 error: &meerkat_mob::MobRespawnError,
57) -> MemberRepairRespawnFailure {
58 if let Some(failed_peer_ids) = topology_restore_failed_peer_ids(error) {
59 return MemberRepairRespawnFailure::DegradedTopologyRestore { failed_peer_ids };
60 }
61 if is_recoverable_bridge_respawn_cleanup_error(&error.to_string()) {
62 return MemberRepairRespawnFailure::RecoverableCleanup;
63 }
64 MemberRepairRespawnFailure::Fatal(error.to_string())
65}
66
67#[derive(Debug)]
73pub enum BridgeError {
74 Mob(String),
76 InvalidInput(String),
78}
79
80impl std::fmt::Display for BridgeError {
81 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
82 match self {
83 Self::Mob(msg) => write!(f, "session bridge mob error: {msg}"),
84 Self::InvalidInput(msg) => write!(f, "session bridge invalid input: {msg}"),
85 }
86 }
87}
88
89impl std::error::Error for BridgeError {}
90
91#[derive(Debug, Clone, PartialEq, Eq)]
94pub enum ResumeFallbackReason {
95 RuntimeIdentityIncompatible { detail: String },
98}
99
100#[derive(Debug, Clone, PartialEq, Eq)]
102pub enum ResumeSessionOutcome {
103 Resumed {
105 session_id: meerkat_core::types::SessionId,
106 },
107 FreshSpawned {
110 session_id: meerkat_core::types::SessionId,
111 reason: ResumeFallbackReason,
112 },
113}
114
115impl ResumeSessionOutcome {
116 #[must_use]
117 pub fn session_id(&self) -> &meerkat_core::types::SessionId {
118 match self {
119 Self::Resumed { session_id } | Self::FreshSpawned { session_id, .. } => session_id,
120 }
121 }
122
123 #[must_use]
124 pub fn fallback_reason(&self) -> Option<&ResumeFallbackReason> {
125 match self {
126 Self::Resumed { .. } => None,
127 Self::FreshSpawned { reason, .. } => Some(reason),
128 }
129 }
130}
131
132async fn submit_internal_bridge_work(
133 handle: &MobHandle,
134 member_id: &MobAgentIdentity,
135 content: &meerkat_core::ContentInput,
136 handling_mode: HandlingMode,
137) -> Result<(), BridgeError> {
138 let entry = handle
139 .get_member(member_id)
140 .await
141 .map_err(|err| BridgeError::Mob(err.to_string()))?
142 .ok_or_else(|| BridgeError::Mob(format!("member not found: {member_id}")))?;
143 handle
144 .submit_work_with_mode(
145 entry.agent_runtime_id.clone(),
146 entry.fence_token,
147 WorkRef::new(),
148 WorkSpec::new(content.clone(), WorkOrigin::Internal),
149 handling_mode,
150 )
151 .await
152 .map(|_| ())
153 .map_err(|err| BridgeError::Mob(err.to_string()))
154}
155
156#[async_trait]
164pub trait SessionBridge: Send + Sync {
165 async fn create_session(
167 &self,
168 identity: &AgentIdentity,
169 runtime_id: &AgentRuntimeId,
170 spec: &DurableAgentSpec,
171 draft: &AgentBuildDraft,
172 session_id: &meerkat_core::types::SessionId,
173 ) -> Result<meerkat_core::types::SessionId, BridgeError>;
174
175 async fn resume_session(
183 &self,
184 identity: &AgentIdentity,
185 runtime_id: &AgentRuntimeId,
186 spec: &DurableAgentSpec,
187 draft: &AgentBuildDraft,
188 session_id: &meerkat_core::types::SessionId,
189 snapshot: &SessionSnapshot,
190 ) -> Result<ResumeSessionOutcome, BridgeError>;
191
192 async fn deliver(
194 &self,
195 runtime_id: &AgentRuntimeId,
196 content: &meerkat_core::ContentInput,
197 ) -> Result<meerkat_core::types::SessionId, BridgeError>;
198
199 async fn deliver_with_mode(
203 &self,
204 runtime_id: &AgentRuntimeId,
205 content: &meerkat_core::ContentInput,
206 handling_mode: HandlingMode,
207 ) -> Result<meerkat_core::types::SessionId, BridgeError> {
208 let _ = handling_mode;
209 self.deliver(runtime_id, content).await
210 }
211
212 async fn checkpoint_session(
214 &self,
215 runtime_id: &AgentRuntimeId,
216 session_id: &meerkat_core::types::SessionId,
217 ) -> Result<SessionSnapshot, BridgeError>;
218
219 async fn retire_member(&self, runtime_id: &AgentRuntimeId) -> Result<(), BridgeError>;
221
222 async fn wire_peer(&self, _a: &AgentRuntimeId, _b: &AgentRuntimeId) -> Result<(), BridgeError> {
224 Err(BridgeError::Mob("peer wiring not supported".to_string()))
225 }
226
227 async fn wire_peers_batch(
229 &self,
230 edges: &[(AgentRuntimeId, AgentRuntimeId)],
231 ) -> Result<(), BridgeError> {
232 for (a, b) in edges {
233 self.wire_peer(a, b).await?;
234 }
235 Ok(())
236 }
237
238 async fn current_member_wires(
240 &self,
241 ) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
242 Ok(Vec::new())
243 }
244
245 async fn unwire_peer(
247 &self,
248 _a: &AgentRuntimeId,
249 _b: &AgentRuntimeId,
250 ) -> Result<(), BridgeError> {
251 Err(BridgeError::Mob("peer unwiring not supported".to_string()))
252 }
253
254 async fn inspect_member(
256 &self,
257 _runtime_id: &AgentRuntimeId,
258 ) -> Result<MemberInspection, BridgeError> {
259 Err(BridgeError::Mob("inspect not supported".to_string()))
260 }
261
262 async fn register_session_runtime_state(
268 &self,
269 _session_id: &meerkat_core::types::SessionId,
270 _identity: &AgentIdentity,
271 _generation: ContinuityGeneration,
272 checkpoint_version: CheckpointVersion,
273 _fencing_token: FencingToken,
274 ) -> Result<CheckpointVersion, BridgeError> {
275 Ok(checkpoint_version)
276 }
277
278 async fn unregister_session_runtime_state(
283 &self,
284 _session_id: &meerkat_core::types::SessionId,
285 ) -> Result<(), BridgeError> {
286 Ok(())
287 }
288}
289
290#[derive(Debug, Clone)]
292pub struct MemberInspection {
293 pub output_preview: Option<String>,
294 pub is_final: bool,
295 pub peer_reachable_count: usize,
296}
297
298pub struct MobSessionBridge {
309 handle: MobHandle,
310 session_store: Option<Arc<dyn meerkat::SessionStore>>,
312 session_service: Option<Arc<dyn MobSessionService>>,
314 continuity_session_store: Option<Arc<ContinuitySessionStoreAdapter>>,
316 runtime_members: Arc<tokio::sync::RwLock<HashMap<String, String>>>,
317 runtime_sessions: Arc<tokio::sync::RwLock<HashMap<String, meerkat_core::types::SessionId>>>,
318 generated_external_owner_session: std::sync::OnceLock<meerkat_core::types::SessionId>,
323}
324
325impl MobSessionBridge {
326 pub fn new(handle: MobHandle) -> Self {
328 Self {
329 handle,
330 session_store: None,
331 session_service: None,
332 continuity_session_store: None,
333 runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
334 runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
335 generated_external_owner_session: std::sync::OnceLock::new(),
336 }
337 }
338
339 pub fn with_session_service(
341 handle: MobHandle,
342 session_service: Arc<dyn MobSessionService>,
343 ) -> Self {
344 Self {
345 handle,
346 session_store: None,
347 session_service: Some(session_service),
348 continuity_session_store: None,
349 runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
350 runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
351 generated_external_owner_session: std::sync::OnceLock::new(),
352 }
353 }
354
355 pub fn with_session_store(
357 handle: MobHandle,
358 session_store: Arc<dyn meerkat::SessionStore>,
359 ) -> Self {
360 Self {
361 handle,
362 session_store: Some(session_store),
363 session_service: None,
364 continuity_session_store: None,
365 runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
366 runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
367 generated_external_owner_session: std::sync::OnceLock::new(),
368 }
369 }
370
371 pub fn with_session_store_and_service(
373 handle: MobHandle,
374 session_store: Arc<dyn meerkat::SessionStore>,
375 session_service: Arc<dyn MobSessionService>,
376 ) -> Self {
377 Self {
378 handle,
379 session_store: Some(session_store),
380 session_service: Some(session_service),
381 continuity_session_store: None,
382 runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
383 runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
384 generated_external_owner_session: std::sync::OnceLock::new(),
385 }
386 }
387
388 pub fn with_continuity_session_store(
390 handle: MobHandle,
391 session_store: Arc<ContinuitySessionStoreAdapter>,
392 session_service: Option<Arc<dyn MobSessionService>>,
393 ) -> Self {
394 Self {
395 handle,
396 session_store: Some(session_store.clone()),
397 session_service,
398 continuity_session_store: Some(session_store),
399 runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
400 runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
401 generated_external_owner_session: std::sync::OnceLock::new(),
402 }
403 }
404
405 async fn remember_runtime_member(
406 &self,
407 runtime_id: &AgentRuntimeId,
408 member_id: &MobAgentIdentity,
409 ) {
410 self.runtime_members.write().await.insert(
411 runtime_id.as_str().to_string(),
412 member_id.as_str().to_string(),
413 );
414 }
415
416 async fn remember_runtime_session(
417 &self,
418 runtime_id: &AgentRuntimeId,
419 session_id: &meerkat_core::types::SessionId,
420 ) {
421 self.runtime_sessions
422 .write()
423 .await
424 .insert(runtime_id.as_str().to_string(), session_id.clone());
425 }
426
427 async fn forget_runtime_member(&self, runtime_id: &AgentRuntimeId) {
428 self.runtime_members
429 .write()
430 .await
431 .remove(runtime_id.as_str());
432 self.runtime_sessions
433 .write()
434 .await
435 .remove(runtime_id.as_str());
436 }
437
438 async fn member_id_for_runtime_id(&self, runtime_id: &AgentRuntimeId) -> MobAgentIdentity {
439 let members = self.runtime_members.read().await;
440 members
441 .get(runtime_id.as_str())
442 .map(|member| MobAgentIdentity::from(member.as_str()))
443 .unwrap_or_else(|| crate::member_comms_id::mob_member_id(runtime_id.as_str()))
446 }
447
448 async fn runtime_session_id(
449 &self,
450 runtime_id: &AgentRuntimeId,
451 ) -> Option<meerkat_core::types::SessionId> {
452 self.runtime_sessions
453 .read()
454 .await
455 .get(runtime_id.as_str())
456 .cloned()
457 }
458
459 fn base_profile_for_spec(&self, spec: &DurableAgentSpec) -> Option<meerkat_mob::Profile> {
463 self.handle
464 .definition()
465 .resolve_inline_profile(&spec.profile)
466 .cloned()
467 }
468
469 async fn resolve_runtime_session_id(
470 &self,
471 runtime_id: &AgentRuntimeId,
472 member_id: &MobAgentIdentity,
473 missing_message: &'static str,
474 ) -> Result<meerkat_core::types::SessionId, BridgeError> {
475 if let Some(session_id) = self.handle.resolve_bridge_session_id(member_id).await {
476 self.remember_runtime_session(runtime_id, &session_id).await;
477 return Ok(session_id);
478 }
479
480 if member_id.as_str() != runtime_id.as_str()
483 && self
484 .handle
485 .get_member(member_id)
486 .await
487 .map_err(|err| BridgeError::Mob(err.to_string()))?
488 .is_some()
489 && let Some(session_id) = self.runtime_session_id(runtime_id).await
490 {
491 return Ok(session_id);
492 }
493
494 Err(BridgeError::Mob(missing_message.to_string()))
495 }
496
497 async fn repair_member_for_delivery(
498 &self,
499 runtime_id: &AgentRuntimeId,
500 member_id: &MobAgentIdentity,
501 member_entry_before_delivery: Option<(meerkat_mob::ProfileName, BTreeMap<String, String>)>,
502 ) -> Result<(), BridgeError> {
503 match self.handle.respawn(member_id.clone(), None).await {
504 Ok(_) => Ok(()),
505 Err(respawn_err) => match classify_member_repair_respawn_failure(&respawn_err) {
506 MemberRepairRespawnFailure::DegradedTopologyRestore { failed_peer_ids } => {
507 tracing::warn!(
508 runtime_id = %runtime_id,
509 member_id = %member_id,
510 failed_peer_count = failed_peer_ids.len(),
511 failed_peer_ids = ?failed_peer_ids,
512 "identity bridge respawn restored member with isolated peer edges; continuing delivery"
513 );
514 Ok(())
516 }
517 MemberRepairRespawnFailure::RecoverableCleanup => {
518 if self
522 .handle
523 .get_member(member_id)
524 .await
525 .map_err(|err| BridgeError::Mob(err.to_string()))?
526 .is_none()
527 && let Some((role, labels)) = member_entry_before_delivery
528 {
529 let mut spec = SpawnMemberSpec::new(role, member_id.clone());
530 if !labels.is_empty() {
531 spec = spec.with_labels(labels);
532 }
533 self.handle
534 .ensure_member(spec)
535 .await
536 .map_err(|e| BridgeError::Mob(e.to_string()))?;
537 }
538 Ok(())
539 }
540 MemberRepairRespawnFailure::Fatal(message) => Err(BridgeError::Mob(message)),
541 },
542 }
543 }
544
545 fn external_owner_bridge_session_id(&self) -> meerkat_core::types::SessionId {
555 if let Some(authority) = self.handle.owner_bridge_session_lifecycle_authority() {
556 return authority.bridge_session_id;
557 }
558 self.generated_external_owner_session
559 .get_or_init(meerkat_core::types::SessionId::new)
560 .clone()
561 }
562
563 async fn spawn_member_spec(
567 &self,
568 spawn_spec: SpawnMemberSpec,
569 ) -> Result<(), meerkat_mob::MobError> {
570 if spawn_spec_requires_generated_owner_context(&spawn_spec) {
571 let owner_session_id = self.external_owner_bridge_session_id();
572 Box::pin(
573 self.handle
574 .spawn_spec_with_generated_owner_context(spawn_spec, owner_session_id),
575 )
576 .await
577 .map(|_| ())
578 } else {
579 Box::pin(self.handle.spawn_spec(spawn_spec))
580 .await
581 .map(|_| ())
582 }
583 }
584}
585
586fn peer_reachable_count_from_connectivity(
597 connectivity: Option<&meerkat_contracts::WirePeerConnectivity>,
598) -> Option<usize> {
599 match connectivity {
600 Some(meerkat_contracts::WirePeerConnectivity::Known { snapshot }) => {
601 Some(snapshot.reachable_peer_count)
602 }
603 Some(
604 meerkat_contracts::WirePeerConnectivity::NotApplicable
605 | meerkat_contracts::WirePeerConnectivity::ProbeTimedOut,
606 )
607 | None => None,
608 }
609}
610
611pub(crate) fn spawn_spec_requires_generated_owner_context(spawn_spec: &SpawnMemberSpec) -> bool {
614 matches!(
615 spawn_spec.binding,
616 Some(meerkat_mob::RuntimeBinding::External { .. })
617 )
618}
619
620fn spec_uses_external_binding(spec: &DurableAgentSpec) -> bool {
621 matches!(spec.backend, Some(meerkat_mob::MobBackendKind::External))
622 || matches!(
623 spec.binding.as_ref(),
624 Some(meerkat_contracts::WireRuntimeBinding::External { .. })
625 )
626}
627
628fn member_id_for_spawn_spec(
629 runtime_id: &AgentRuntimeId,
630 spec: &DurableAgentSpec,
631) -> MobAgentIdentity {
632 if spec_uses_external_binding(spec) {
638 crate::member_comms_id::mob_member_id(spec.identity.as_str())
639 } else {
640 crate::member_comms_id::mob_member_id(runtime_id.as_str())
641 }
642}
643
644pub(crate) fn build_spawn_spec(
652 runtime_id: &AgentRuntimeId,
653 spec: &DurableAgentSpec,
654 draft: &AgentBuildDraft,
655 base_profile: Option<&meerkat_mob::Profile>,
656) -> SpawnMemberSpec {
657 let mid = member_id_for_spawn_spec(runtime_id, spec);
658 let mut spawn_spec = SpawnMemberSpec::new(spec.profile.clone(), mid);
659
660 if let Some(message) = spec.initial_message.as_ref() {
661 spawn_spec = spawn_spec.with_initial_message(message.clone());
662 }
663 if let Some(runtime_mode) = spec.runtime_mode_override {
664 spawn_spec = spawn_spec.with_runtime_mode(runtime_mode);
665 }
666 spawn_spec.backend = spec.backend;
667 if let Some(binding) = spec.binding.clone() {
668 spawn_spec.binding = runtime_binding_from_wire(binding);
669 }
670 if let Some(ref ctx) = draft.app_context {
671 spawn_spec = spawn_spec.with_context(ctx.clone());
672 }
673 let mut labels = draft.labels.clone();
674 labels.insert(
675 "agent_identity".to_string(),
676 spec.identity.as_str().to_string(),
677 );
678 if !labels.is_empty() {
679 spawn_spec = spawn_spec.with_labels(labels);
680 }
681 if !draft.additional_instructions.is_empty() {
682 spawn_spec = spawn_spec.with_additional_instructions(draft.additional_instructions.clone());
683 }
684 if let Some(model) = draft.model.as_ref() {
685 match base_profile {
686 Some(base) => {
687 let mut profile = base.clone();
688 profile.model = model.clone();
689 profile.provider = None;
693 profile.self_hosted_server_id = None;
694 spawn_spec.override_profile = Some(profile);
695 }
696 None => {
697 tracing::warn!(
702 identity = %spec.identity,
703 profile = %spec.profile,
704 model = %model,
705 "model override skipped: role profile is not an inline definition profile"
706 );
707 }
708 }
709 }
710 if let Some(system_prompt) = draft.system_prompt.as_ref() {
711 spawn_spec.system_prompt_override =
712 Some(SpawnSystemPromptOverride::Replace(system_prompt.clone()));
713 }
714 if let Some(dispatcher) = draft.local_external_tools.dispatcher() {
715 spawn_spec.external_tools = Some(dispatcher);
716 }
717
718 spawn_spec
719}
720
721fn runtime_binding_from_wire(
722 binding: meerkat_contracts::WireRuntimeBinding,
723) -> Option<meerkat_mob::RuntimeBinding> {
724 match binding {
725 meerkat_contracts::WireRuntimeBinding::Session => {
726 Some(meerkat_mob::RuntimeBinding::Session)
727 }
728 meerkat_contracts::WireRuntimeBinding::External {
729 address,
730 bootstrap_token,
731 identity,
732 } => {
733 let resolved = identity.resolve().ok()?;
734 Some(meerkat_mob::RuntimeBinding::External {
735 peer_id: resolved.peer_id.to_string(),
736 address,
737 bootstrap_token,
738 pubkey: resolved.pubkey,
739 })
740 }
741 }
742}
743
744#[async_trait]
745impl SessionBridge for MobSessionBridge {
746 async fn create_session(
747 &self,
748 _identity: &AgentIdentity,
749 runtime_id: &AgentRuntimeId,
750 spec: &DurableAgentSpec,
751 draft: &AgentBuildDraft,
752 session_id: &meerkat_core::types::SessionId,
753 ) -> Result<meerkat_core::types::SessionId, BridgeError> {
754 let mid = member_id_for_spawn_spec(runtime_id, spec);
755 let spawn_spec = build_spawn_spec(
756 runtime_id,
757 spec,
758 draft,
759 self.base_profile_for_spec(spec).as_ref(),
760 );
761
762 self.spawn_member_spec(spawn_spec)
763 .await
764 .map_err(|e| BridgeError::Mob(e.to_string()))?;
765 self.remember_runtime_member(runtime_id, &mid).await;
766 self.remember_runtime_session(runtime_id, session_id).await;
767
768 self.resolve_runtime_session_id(runtime_id, &mid, "member spawned but has no session ID")
769 .await
770 }
771
772 async fn resume_session(
773 &self,
774 _identity: &AgentIdentity,
775 runtime_id: &AgentRuntimeId,
776 spec: &DurableAgentSpec,
777 draft: &AgentBuildDraft,
778 session_id: &meerkat_core::types::SessionId,
779 _snapshot: &SessionSnapshot,
780 ) -> Result<ResumeSessionOutcome, BridgeError> {
781 if spec_uses_external_binding(spec) {
782 let mut spawn_spec = build_spawn_spec(
783 runtime_id,
784 spec,
785 draft,
786 self.base_profile_for_spec(spec).as_ref(),
787 );
788 spawn_spec.launch_mode = MemberLaunchMode::Resume {
789 bridge_session_id: session_id.clone(),
790 };
791 let mid = member_id_for_spawn_spec(runtime_id, spec);
792 self.spawn_member_spec(spawn_spec)
793 .await
794 .map_err(|e| BridgeError::Mob(e.to_string()))?;
795 self.remember_runtime_member(runtime_id, &mid).await;
796 self.remember_runtime_session(runtime_id, session_id).await;
797 return Ok(ResumeSessionOutcome::Resumed {
798 session_id: session_id.clone(),
799 });
800 }
801
802 let mut spawn_spec = build_spawn_spec(
805 runtime_id,
806 spec,
807 draft,
808 self.base_profile_for_spec(spec).as_ref(),
809 );
810 spawn_spec.launch_mode = MemberLaunchMode::Resume {
811 bridge_session_id: session_id.clone(),
812 };
813
814 let mid = member_id_for_spawn_spec(runtime_id, spec);
815
816 match self.spawn_member_spec(spawn_spec).await {
817 Ok(()) => {
818 self.remember_runtime_member(runtime_id, &mid).await;
819 self.remember_runtime_session(runtime_id, session_id).await;
820 Ok(ResumeSessionOutcome::Resumed {
821 session_id: session_id.clone(),
822 })
823 }
824 Err(e) => {
825 tracing::warn!(
829 identity = %_identity,
830 session_id = %session_id,
831 error = %e,
832 reason = "runtime_identity_incompatible",
833 "resume_session incompatible with current runtime binding, falling back to fresh spawn"
834 );
835 let fresh_spec = build_spawn_spec(
836 runtime_id,
837 spec,
838 draft,
839 self.base_profile_for_spec(spec).as_ref(),
840 );
841 self.spawn_member_spec(fresh_spec)
842 .await
843 .map_err(|e2| BridgeError::Mob(e2.to_string()))?;
844
845 self.remember_runtime_member(runtime_id, &mid).await;
846 let session_id = self
847 .resolve_runtime_session_id(
848 runtime_id,
849 &mid,
850 "member spawned (fresh fallback) but has no session ID",
851 )
852 .await?;
853 Ok(ResumeSessionOutcome::FreshSpawned {
854 session_id,
855 reason: ResumeFallbackReason::RuntimeIdentityIncompatible {
856 detail: e.to_string(),
857 },
858 })
859 }
860 }
861 }
862
863 async fn deliver(
864 &self,
865 runtime_id: &AgentRuntimeId,
866 content: &meerkat_core::ContentInput,
867 ) -> Result<meerkat_core::types::SessionId, BridgeError> {
868 let mid = self.member_id_for_runtime_id(runtime_id).await;
869 let member_entry_before_delivery = self
872 .handle
873 .get_member(&mid)
874 .await
875 .ok()
876 .flatten()
877 .map(|entry| (entry.role, entry.labels));
878 if content_input_has_images(content) {
879 let member_entry = self
880 .handle
881 .get_member(&mid)
882 .await
883 .map_err(|err| BridgeError::Mob(err.to_string()))?
884 .ok_or_else(|| {
885 BridgeError::Mob("member not found while checking image capability".to_string())
886 })?;
887 let caps = model_capabilities_for_member(
888 &self.handle,
889 self.session_service.as_ref(),
890 &member_entry.agent_identity,
891 )
892 .await;
893 if !caps.image_input {
894 return Err(BridgeError::InvalidInput(
895 "target member model cannot accept image input".to_string(),
896 ));
897 }
898 }
899
900 match submit_internal_bridge_work(&self.handle, &mid, content, HandlingMode::Queue).await {
906 Ok(()) => {}
907 Err(err) if is_repairable_bridge_delivery_error(&err.to_string()) => {
908 tracing::warn!(
909 runtime_id = %runtime_id,
910 error = %err,
911 "identity bridge delivery found stale runtime state; repairing member before retry"
912 );
913 Box::pin(self.repair_member_for_delivery(
914 runtime_id,
915 &mid,
916 member_entry_before_delivery,
917 ))
918 .await?;
919 submit_internal_bridge_work(&self.handle, &mid, content, HandlingMode::Queue)
920 .await?;
921 }
922 Err(err) => return Err(BridgeError::Mob(err.to_string())),
923 }
924
925 self.resolve_runtime_session_id(
928 runtime_id,
929 &mid,
930 "member has no bridge session after deliver",
931 )
932 .await
933 }
934
935 async fn deliver_with_mode(
936 &self,
937 runtime_id: &AgentRuntimeId,
938 content: &meerkat_core::ContentInput,
939 handling_mode: HandlingMode,
940 ) -> Result<meerkat_core::types::SessionId, BridgeError> {
941 let mid = self.member_id_for_runtime_id(runtime_id).await;
942 let member_entry_before_delivery = self
945 .handle
946 .get_member(&mid)
947 .await
948 .ok()
949 .flatten()
950 .map(|entry| (entry.role, entry.labels));
951 if content_input_has_images(content) {
952 let member_entry = self
953 .handle
954 .get_member(&mid)
955 .await
956 .map_err(|err| BridgeError::Mob(err.to_string()))?
957 .ok_or_else(|| {
958 BridgeError::Mob("member not found while checking image capability".to_string())
959 })?;
960 let caps = model_capabilities_for_member(
961 &self.handle,
962 self.session_service.as_ref(),
963 &member_entry.agent_identity,
964 )
965 .await;
966 if !caps.image_input {
967 return Err(BridgeError::InvalidInput(
968 "target member model cannot accept image input".to_string(),
969 ));
970 }
971 }
972
973 match submit_internal_bridge_work(&self.handle, &mid, content, handling_mode).await {
974 Ok(()) => {}
975 Err(err) if is_repairable_bridge_delivery_error(&err.to_string()) => {
976 tracing::warn!(
977 runtime_id = %runtime_id,
978 error = %err,
979 "identity bridge delivery found stale runtime state; repairing member before retry"
980 );
981 Box::pin(self.repair_member_for_delivery(
982 runtime_id,
983 &mid,
984 member_entry_before_delivery,
985 ))
986 .await?;
987 submit_internal_bridge_work(&self.handle, &mid, content, handling_mode).await?;
988 }
989 Err(err) => return Err(BridgeError::Mob(err.to_string())),
990 }
991
992 self.resolve_runtime_session_id(
993 runtime_id,
994 &mid,
995 "member has no bridge session after deliver",
996 )
997 .await
998 }
999
1000 async fn checkpoint_session(
1001 &self,
1002 _runtime_id: &AgentRuntimeId,
1003 session_id: &meerkat_core::types::SessionId,
1004 ) -> Result<SessionSnapshot, BridgeError> {
1005 let store = self.session_store.as_ref().ok_or_else(|| {
1006 BridgeError::InvalidInput(
1007 "checkpoint requires a session store but none was configured".to_string(),
1008 )
1009 })?;
1010
1011 let session = store
1012 .load(session_id)
1013 .await
1014 .map_err(|e| BridgeError::Mob(format!("failed to load session for checkpoint: {e}")))?
1015 .ok_or_else(|| {
1016 BridgeError::Mob(format!(
1017 "session {session_id} not found in store for checkpoint"
1018 ))
1019 })?;
1020
1021 let data = serde_json::to_vec(&session)
1022 .map_err(|e| BridgeError::Mob(format!("failed to serialize session: {e}")))?;
1023
1024 Ok(SessionSnapshot { data })
1025 }
1026
1027 async fn retire_member(&self, runtime_id: &AgentRuntimeId) -> Result<(), BridgeError> {
1028 let mid = self.member_id_for_runtime_id(runtime_id).await;
1029 match self.handle.retire(mid).await {
1030 Ok(()) => {
1031 self.forget_runtime_member(runtime_id).await;
1032 Ok(())
1033 }
1034 Err(err) if is_recoverable_lifecycle_cleanup_error(&err.to_string()) => Ok(()),
1035 Err(err) => Err(BridgeError::Mob(err.to_string())),
1036 }
1037 }
1038
1039 async fn wire_peer(&self, a: &AgentRuntimeId, b: &AgentRuntimeId) -> Result<(), BridgeError> {
1040 let member_a = self.member_id_for_runtime_id(a).await;
1041 let member_b = self.member_id_for_runtime_id(b).await;
1042 self.handle
1043 .wire(
1044 meerkat_mob::AgentIdentity::from(member_a.as_str()),
1045 member_b,
1046 )
1047 .await
1048 .map_err(|e| BridgeError::Mob(e.to_string()))
1049 }
1050
1051 async fn wire_peers_batch(
1052 &self,
1053 edges: &[(AgentRuntimeId, AgentRuntimeId)],
1054 ) -> Result<(), BridgeError> {
1055 let mut member_edges = Vec::with_capacity(edges.len());
1056 for (a, b) in edges {
1057 let member_a = self.member_id_for_runtime_id(a).await;
1058 let member_b = self.member_id_for_runtime_id(b).await;
1059 member_edges.push((
1060 meerkat_mob::AgentIdentity::from(member_a.as_str()),
1061 meerkat_mob::AgentIdentity::from(member_b.as_str()),
1062 ));
1063 }
1064 self.handle
1065 .wire_members_batch(member_edges)
1066 .await
1067 .map(|_| ())
1068 .map_err(|e| BridgeError::Mob(e.to_string()))
1069 }
1070
1071 async fn current_member_wires(
1072 &self,
1073 ) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
1074 let members = self.handle.list_members_including_retiring().await;
1075 let runtime_members = self.runtime_members.read().await;
1076 let member_runtimes = runtime_members
1077 .iter()
1078 .map(|(runtime, member)| (member.clone(), runtime.clone()))
1079 .collect::<HashMap<_, _>>();
1080 let active_ids = members
1081 .iter()
1082 .map(|member| member.agent_identity.to_string())
1083 .collect::<std::collections::BTreeSet<_>>();
1084 let mut edges = std::collections::BTreeSet::new();
1085 for member in &members {
1086 let a = member.agent_identity.to_string();
1087 for peer in &member.wired_to {
1088 let b = peer.to_string();
1089 if !active_ids.contains(&b) {
1090 continue;
1091 }
1092 let key = if a <= b {
1093 (a.clone(), b)
1094 } else {
1095 (b, a.clone())
1096 };
1097 edges.insert(key);
1098 }
1099 }
1100 Ok(edges
1101 .into_iter()
1102 .filter_map(|(a, b)| {
1103 let a = member_runtimes
1106 .get(&a)
1107 .cloned()
1108 .unwrap_or_else(|| crate::member_comms_id::runtime_alias_str(&a).into_owned());
1109 let b = member_runtimes
1110 .get(&b)
1111 .cloned()
1112 .unwrap_or_else(|| crate::member_comms_id::runtime_alias_str(&b).into_owned());
1113 Some((
1114 AgentRuntimeId::parse(&a).ok()?,
1115 AgentRuntimeId::parse(&b).ok()?,
1116 ))
1117 })
1118 .collect())
1119 }
1120
1121 async fn unwire_peer(&self, a: &AgentRuntimeId, b: &AgentRuntimeId) -> Result<(), BridgeError> {
1122 let member_a = self.member_id_for_runtime_id(a).await;
1123 let member_b = self.member_id_for_runtime_id(b).await;
1124 match self
1125 .handle
1126 .unwire(
1127 meerkat_mob::AgentIdentity::from(member_a.as_str()),
1128 member_b,
1129 )
1130 .await
1131 {
1132 Ok(()) => Ok(()),
1133 Err(err) => {
1134 let message = err.to_string();
1135 if message.contains("peer not found") || message.contains("not wired") {
1136 Ok(())
1137 } else {
1138 Err(BridgeError::Mob(message))
1139 }
1140 }
1141 }
1142 }
1143
1144 async fn inspect_member(
1145 &self,
1146 runtime_id: &AgentRuntimeId,
1147 ) -> Result<MemberInspection, BridgeError> {
1148 let mid = self.member_id_for_runtime_id(runtime_id).await;
1149 let snap = self
1150 .handle
1151 .member_status(&mid)
1152 .await
1153 .map_err(|e| BridgeError::Mob(e.to_string()))?;
1154 let peer_reachable_count =
1155 match peer_reachable_count_from_connectivity(snap.peer_connectivity.as_ref()) {
1156 Some(count) => count,
1157 None => self
1158 .handle
1159 .get_member(&mid)
1160 .await
1161 .ok()
1162 .flatten()
1163 .map(|entry| entry.wired_to.len())
1164 .unwrap_or(0),
1165 };
1166 Ok(MemberInspection {
1167 output_preview: snap.output_preview.clone(),
1168 is_final: snap.is_final,
1169 peer_reachable_count,
1170 })
1171 }
1172
1173 async fn register_session_runtime_state(
1174 &self,
1175 session_id: &meerkat_core::types::SessionId,
1176 identity: &AgentIdentity,
1177 generation: ContinuityGeneration,
1178 checkpoint_version: CheckpointVersion,
1179 fencing_token: FencingToken,
1180 ) -> Result<CheckpointVersion, BridgeError> {
1181 if let Some(adapter) = self.continuity_session_store.as_ref() {
1182 return adapter
1183 .register_session(
1184 session_id,
1185 SessionRuntimeState {
1186 identity: identity.clone(),
1187 generation,
1188 checkpoint_version,
1189 fencing_token,
1190 },
1191 )
1192 .await
1193 .map_err(|err| BridgeError::Mob(format!("continuity register_session: {err}")));
1194 }
1195 Ok(checkpoint_version)
1196 }
1197
1198 async fn unregister_session_runtime_state(
1199 &self,
1200 session_id: &meerkat_core::types::SessionId,
1201 ) -> Result<(), BridgeError> {
1202 if let Some(adapter) = self.continuity_session_store.as_ref() {
1203 adapter
1204 .unregister_session(session_id)
1205 .await
1206 .map_err(|err| BridgeError::Mob(format!("continuity unregister_session: {err}")))?;
1207 }
1208 Ok(())
1209 }
1210}
1211
1212#[cfg(test)]
1213#[allow(clippy::expect_used)]
1214mod tests {
1215 use std::sync::Arc;
1216
1217 use async_trait::async_trait;
1218 use meerkat_core::agent::AgentToolDispatcher;
1219 use meerkat_core::types::ToolCallView;
1220 use meerkat_core::{ToolDef, error::ToolError, ops::ToolDispatchOutcome};
1221 use meerkat_mob::{MobRespawnError, MobRuntimeMode};
1222
1223 use super::*;
1224 use crate::identity_first::{AgentAddressability, LocalExternalToolOverlay};
1225
1226 struct EmptyDispatcher;
1227
1228 #[async_trait]
1229 impl AgentToolDispatcher for EmptyDispatcher {
1230 fn tools(&self) -> Arc<[Arc<ToolDef>]> {
1231 Arc::from([])
1232 }
1233
1234 async fn dispatch(
1235 &self,
1236 _call: ToolCallView<'_>,
1237 ) -> Result<ToolDispatchOutcome, ToolError> {
1238 Err(ToolError::ExecutionFailed {
1239 message: "not implemented".to_string(),
1240 })
1241 }
1242 }
1243
1244 fn durable_spec() -> DurableAgentSpec {
1245 DurableAgentSpec {
1246 identity: AgentIdentity::parse("agent:alpha").expect("identity"),
1247 profile: meerkat_mob::ProfileName::from("worker"),
1248 addressability: AgentAddressability::Addressable,
1249 display_name: None,
1250 labels: Default::default(),
1251 context: None,
1252 additional_instructions: Vec::new(),
1253 initial_message: Some(meerkat_core::ContentInput::Text("hello".to_string())),
1254 runtime_mode_override: Some(MobRuntimeMode::TurnDriven),
1255 backend: None,
1256 binding: None,
1257 }
1258 }
1259
1260 #[test]
1267 fn peer_reachable_count_tri_state_defers_to_wiring_when_unresolved() {
1268 use meerkat_contracts::{WirePeerConnectivity, WirePeerConnectivitySnapshot};
1269
1270 let known = WirePeerConnectivity::Known {
1271 snapshot: WirePeerConnectivitySnapshot {
1272 reachable_peer_count: 3,
1273 unknown_peer_count: 0,
1274 unreachable_peers: Vec::new(),
1275 },
1276 };
1277 assert_eq!(
1278 peer_reachable_count_from_connectivity(Some(&known)),
1279 Some(3),
1280 "a resolved probe owns the count"
1281 );
1282 assert_eq!(
1283 peer_reachable_count_from_connectivity(Some(&WirePeerConnectivity::NotApplicable)),
1284 None,
1285 "not-applicable must defer to the wiring fallback"
1286 );
1287 assert_eq!(
1288 peer_reachable_count_from_connectivity(Some(&WirePeerConnectivity::ProbeTimedOut)),
1289 None,
1290 "probe timeout must defer to the wiring fallback"
1291 );
1292 assert_eq!(peer_reachable_count_from_connectivity(None), None);
1293 }
1294
1295 #[test]
1296 fn build_spawn_spec_maps_identity_first_overrides() {
1297 let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
1298 let mut labels = std::collections::BTreeMap::new();
1299 labels.insert("team".to_string(), "ops".to_string());
1300 let draft = AgentBuildDraft {
1301 model: Some("gpt-test".to_string()),
1302 system_prompt: Some("system override".to_string()),
1303 additional_instructions: vec!["stay focused".to_string()],
1304 labels,
1305 app_context: Some(serde_json::json!({"ticket": 7})),
1306 external_tools: Vec::new(),
1307 local_external_tools: LocalExternalToolOverlay::new(Arc::new(EmptyDispatcher)),
1308 };
1309
1310 let base_profile: meerkat_mob::Profile =
1311 serde_json::from_value(serde_json::json!({"model": "base-model"}))
1312 .expect("minimal profile");
1313 let spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, Some(&base_profile));
1314
1315 assert_eq!(
1316 spawn
1317 .override_profile
1318 .as_ref()
1319 .map(|profile| profile.model.as_str()),
1320 Some("gpt-test")
1321 );
1322 assert_eq!(
1323 spawn.system_prompt_override,
1324 Some(SpawnSystemPromptOverride::Replace(
1325 "system override".to_string()
1326 ))
1327 );
1328 assert!(spawn.external_tools.is_some());
1329 assert_eq!(spawn.runtime_mode, Some(MobRuntimeMode::TurnDriven));
1330 assert_eq!(
1331 spawn.initial_message,
1332 Some(meerkat_core::ContentInput::Text("hello".to_string()))
1333 );
1334 assert_eq!(
1335 spawn
1336 .labels
1337 .as_ref()
1338 .and_then(|labels| labels.get("team"))
1339 .map(String::as_str),
1340 Some("ops")
1341 );
1342 assert_eq!(
1343 spawn
1344 .labels
1345 .as_ref()
1346 .and_then(|labels| labels.get("agent_identity"))
1347 .map(String::as_str),
1348 Some("agent:alpha")
1349 );
1350 }
1351
1352 #[test]
1353 fn build_spawn_spec_maps_remote_runtime_binding() {
1354 let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
1355 let mut spec = durable_spec();
1356 spec.backend = Some(meerkat_mob::MobBackendKind::External);
1357 spec.binding = Some(
1358 serde_json::from_value(serde_json::json!({
1359 "kind": "external",
1360 "address": "tcp://127.0.0.1:4777",
1361 "identity": {
1362 "kind": "ed25519_public_key",
1363 "public_key": "ed25519:BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc="
1364 }
1365 }))
1366 .expect("wire binding"),
1367 );
1368 let draft = AgentBuildDraft {
1369 model: None,
1370 system_prompt: None,
1371 additional_instructions: Vec::new(),
1372 labels: Default::default(),
1373 app_context: None,
1374 external_tools: Vec::new(),
1375 local_external_tools: Default::default(),
1376 };
1377
1378 let spawn = build_spawn_spec(&runtime_id, &spec, &draft, None);
1379
1380 assert_eq!(
1384 spawn.identity.as_str(),
1385 crate::member_comms_id::mob_member_id_str("agent:alpha").as_ref()
1386 );
1387 assert_eq!(spawn.backend, Some(meerkat_mob::MobBackendKind::External));
1388 assert!(
1389 matches!(
1390 spawn.binding,
1391 Some(meerkat_mob::RuntimeBinding::External { .. })
1392 ),
1393 "expected external runtime binding, got {:?}",
1394 spawn.binding
1395 );
1396 if let Some(meerkat_mob::RuntimeBinding::External {
1397 address, pubkey, ..
1398 }) = spawn.binding
1399 {
1400 assert_eq!(address.as_str(), "tcp://127.0.0.1:4777");
1401 assert_eq!(pubkey, [7; 32]);
1402 }
1403 }
1404
1405 #[test]
1406 fn external_binding_spawn_specs_require_generated_owner_context() {
1407 let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
1408 let draft = AgentBuildDraft {
1409 model: None,
1410 system_prompt: None,
1411 additional_instructions: Vec::new(),
1412 labels: Default::default(),
1413 app_context: None,
1414 external_tools: Vec::new(),
1415 local_external_tools: Default::default(),
1416 };
1417
1418 let session_spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, None);
1420 assert!(
1421 !spawn_spec_requires_generated_owner_context(&session_spawn),
1422 "session-backed spawns must not require a generated owner context"
1423 );
1424
1425 let mut spec = durable_spec();
1430 spec.backend = Some(meerkat_mob::MobBackendKind::External);
1431 spec.binding = Some(
1432 serde_json::from_value(serde_json::json!({
1433 "kind": "external",
1434 "address": "tcp://127.0.0.1:4777",
1435 "identity": {
1436 "kind": "ed25519_public_key",
1437 "public_key": "ed25519:BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc="
1438 }
1439 }))
1440 .expect("wire binding"),
1441 );
1442 let external_spawn = build_spawn_spec(&runtime_id, &spec, &draft, None);
1443 assert!(
1444 spawn_spec_requires_generated_owner_context(&external_spawn),
1445 "external peer-only spawns must carry a generated owner binding on meerkat 0.7.1"
1446 );
1447 }
1448
1449 #[test]
1450 fn bridge_delivery_repair_covers_missing_bridge_session_snapshot() {
1451 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'";
1452
1453 assert!(
1454 is_repairable_bridge_delivery_error(error),
1455 "stale bridge-session bindings should be repaired before retrying delivery"
1456 );
1457 assert!(
1458 is_repairable_bridge_delivery_error("missing event injector capability for member"),
1459 "existing stale event-injector repair path must remain covered"
1460 );
1461 assert!(
1462 is_repairable_bridge_delivery_error(
1463 "mob member rt:us-president:0 missing required capability interaction_event_injector: autonomous member dispatch"
1464 ),
1465 "newer autonomous member dispatch wording should repair and retry instead of dropping the event"
1466 );
1467 assert!(
1468 is_repairable_bridge_delivery_error(
1469 "previous member cleanup ambiguous for member rt:deep-investigator:singleton:0"
1470 ),
1471 "ambiguous Meerkat respawn cleanup should trigger bridge repair instead of failing delivery"
1472 );
1473 assert!(
1474 !is_repairable_bridge_delivery_error("model provider returned rate limit"),
1475 "ordinary turn failures must not trigger member repair"
1476 );
1477 }
1478
1479 #[test]
1480 fn bridge_delivery_repair_classifies_topology_restore_failure_as_degraded() {
1481 let identity = meerkat_mob::AgentIdentity::from("rt:review:singleton:0");
1482 let receipt = meerkat_mob::MemberRespawnReceipt::new(
1483 identity.clone(),
1484 meerkat_mob::AgentRuntimeId::new(identity, meerkat_mob::ids::Generation::INITIAL),
1485 meerkat_mob::FenceToken::new(1),
1486 meerkat_mob::FenceToken::new(2),
1487 );
1488 let err = MobRespawnError::TopologyRestoreFailed {
1489 receipt,
1490 failed_peer_ids: vec![meerkat_mob::RespawnTopologyPeerId::from(
1491 "initiative:broken",
1492 )],
1493 };
1494
1495 assert_eq!(
1496 classify_member_repair_respawn_failure(&err),
1497 MemberRepairRespawnFailure::DegradedTopologyRestore {
1498 failed_peer_ids: vec!["initiative:broken".to_string()]
1499 },
1500 "failed peer edges should degrade bridge repair instead of bricking delivery"
1501 );
1502 assert!(
1503 matches!(
1504 classify_member_repair_respawn_failure(&MobRespawnError::NoRuntimeControl {
1505 identity: meerkat_mob::AgentIdentity::from("rt:review:singleton:0"),
1506 }),
1507 MemberRepairRespawnFailure::Fatal(_)
1508 ),
1509 "ordinary respawn failures must still fail bridge repair"
1510 );
1511 }
1512}