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, 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#[derive(Debug)]
77pub enum BridgeError {
78 Mob(String),
80 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#[derive(Debug, Clone, PartialEq, Eq)]
98pub enum ResumeFallbackReason {
99 RuntimeIdentityIncompatible { detail: String },
102}
103
104#[derive(Debug, Clone, PartialEq, Eq)]
106pub enum ResumeSessionOutcome {
107 Resumed {
109 session_id: meerkat_core::types::SessionId,
110 },
111 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 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#[async_trait]
177pub trait SessionBridge: Send + Sync {
178 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 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 async fn deliver(
207 &self,
208 runtime_id: &AgentRuntimeId,
209 content: &meerkat_core::ContentInput,
210 ) -> Result<meerkat_core::types::SessionId, BridgeError>;
211
212 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 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 async fn checkpoint_session(
243 &self,
244 runtime_id: &AgentRuntimeId,
245 session_id: &meerkat_core::types::SessionId,
246 ) -> Result<SessionSnapshot, BridgeError>;
247
248 async fn retire_member(&self, runtime_id: &AgentRuntimeId) -> Result<(), BridgeError>;
250
251 async fn wire_peer(&self, _a: &AgentRuntimeId, _b: &AgentRuntimeId) -> Result<(), BridgeError> {
253 Err(BridgeError::Mob("peer wiring not supported".to_string()))
254 }
255
256 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 async fn current_member_wires(
269 &self,
270 ) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
271 Ok(Vec::new())
272 }
273
274 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 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 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 async fn unregister_session_runtime_state(
312 &self,
313 _session_id: &meerkat_core::types::SessionId,
314 ) -> Result<(), BridgeError> {
315 Ok(())
316 }
317}
318
319#[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
327pub struct MobSessionBridge {
338 handle: MobHandle,
339 session_store: Option<Arc<dyn meerkat::SessionStore>>,
341 session_service: Option<Arc<dyn MobSessionService>>,
343 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 generated_external_owner_session: std::sync::OnceLock<meerkat_core::types::SessionId>,
352}
353
354impl MobSessionBridge {
355 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 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 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 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 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 .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 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 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 Ok(())
545 }
546 MemberRepairRespawnFailure::RecoverableCleanup => {
547 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 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 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
643fn 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
668pub(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 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
701pub(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 profile.provider = None;
754 profile.self_hosted_server_id = None;
755 spawn_spec.override_profile = Some(profile);
756 }
757 None => {
758 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 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 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 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 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 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 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 Err(err) if is_recoverable_session_owned_retire_cleanup_error(&err.to_string()) => {
1129 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 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 #[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 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 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 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}