1use meerkat_core::comms::TrustedPeerDescriptor;
4use meerkat_core::types::HandlingMode;
5use meerkat_mob::ids::AgentIdentity;
6use meerkat_mob::{MobHandle, PeerTarget};
7
8use crate::auth::peer_keys::GatewayPeerKeys;
9use crate::contact_directory::{ContactDirectory, ContactEntry, MobTransport};
10use crate::runtime::cross_mob_remote::{RemoteMobError, RemoteMobProxy};
11
12use super::UnifiedRuntime;
13
14enum LocalOrRemote {
21 Local(MobHandle),
23 Remote(RemoteMobProxy),
25}
26
27struct MemberPeerInfo {
28 peer_id: String,
29 comms_name: String,
30 pubkey: [u8; 32],
31}
32
33#[derive(Debug)]
35pub enum CrossMobError {
36 NoContactDirectory,
38 UnknownMob(String),
40 NoPeerHandle(String),
42 MissingPeerPubkey { mob_id: Option<String> },
47 MemberNotFound { member_id: String, mob_id: String },
49 NoCommsInfo { member_id: String, mob_id: String },
51 Mob(meerkat_mob::MobError),
53 PeerSpec(String),
55 Remote(RemoteMobError),
60}
61
62impl std::fmt::Display for CrossMobError {
63 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
64 match self {
65 Self::NoContactDirectory => write!(f, "no contact directory configured"),
66 Self::UnknownMob(id) => write!(f, "unknown mob: {id}"),
67 Self::NoPeerHandle(id) => write!(f, "no peer mob handle registered for: {id}"),
68 Self::Remote(e) => write!(f, "cross-process cross-mob: {e}"),
69 Self::MemberNotFound { member_id, mob_id } => {
70 write!(f, "member '{member_id}' not found in mob '{mob_id}'")
71 }
72 Self::NoCommsInfo { member_id, mob_id } => {
73 write!(
74 f,
75 "member '{member_id}' in mob '{mob_id}' has no comms runtime"
76 )
77 }
78 Self::Mob(err) => write!(f, "mob error: {err}"),
79 Self::PeerSpec(reason) => write!(f, "peer spec error: {reason}"),
80 Self::MissingPeerPubkey { mob_id } => match mob_id {
81 Some(id) => write!(
82 f,
83 "non-inproc peer for mob '{id}' has no signing pubkey; \
84 bootstrap via mobkit/peer_pubkey or populate the contact \
85 directory's pubkey field before wiring"
86 ),
87 None => write!(
88 f,
89 "non-inproc peer has no signing pubkey; supply a 32-byte \
90 Ed25519 pubkey or use inproc transport"
91 ),
92 },
93 }
94 }
95}
96
97impl std::error::Error for CrossMobError {}
98
99impl From<meerkat_mob::MobError> for CrossMobError {
100 fn from(err: meerkat_mob::MobError) -> Self {
101 Self::Mob(err)
102 }
103}
104
105impl From<RemoteMobError> for CrossMobError {
106 fn from(err: RemoteMobError) -> Self {
107 Self::Remote(err)
108 }
109}
110
111impl UnifiedRuntime {
112 pub async fn register_peer_mob(&self, mob_id: &str, handle: MobHandle) {
114 self.peer_mob_handles
115 .write()
116 .await
117 .insert(mob_id.to_string(), handle);
118 }
119
120 pub fn set_contact_directory(&mut self, directory: ContactDirectory) {
122 self.contact_directory = Some(directory);
123 }
124
125 pub fn set_gateway_peer_keys(&mut self, keys: GatewayPeerKeys) {
134 self.gateway_peer_keys = Some(keys);
135 }
136
137 pub fn gateway_peer_keys(&self) -> Option<&GatewayPeerKeys> {
139 self.gateway_peer_keys.as_ref()
140 }
141
142 pub async fn wire_cross_mob(
158 &self,
159 local_member_id: &str,
160 remote_member_id: &str,
161 remote_mob_id: &str,
162 ) -> Result<(), CrossMobError> {
163 let entry = self.resolve_contact(remote_mob_id)?;
164 let remote = self.dispatch_for(&entry).await?;
165
166 let local_handle = self.mob_runtime.handle();
167 let local_mob_id = local_handle.mob_id().to_string();
168 let local_mid = crate::member_comms_id::mob_member_id(local_member_id);
172
173 let local_info = self
174 .get_member_peer_info(&local_handle, &local_mid, &local_mob_id)
175 .await?;
176
177 match remote {
178 LocalOrRemote::Local(remote_handle) => {
179 let remote_mid = crate::member_comms_id::mob_member_id(remote_member_id);
180 let remote_info = self
181 .get_member_peer_info(&remote_handle, &remote_mid, remote_mob_id)
182 .await?;
183
184 let remote_spec = build_peer_spec(
185 &remote_info.comms_name,
186 &remote_info.peer_id,
187 &entry.transport,
188 Some(remote_info.pubkey),
189 )?;
190 let local_spec = build_peer_spec(
191 &local_info.comms_name,
192 &local_info.peer_id,
193 &MobTransport::Inproc,
194 Some(local_info.pubkey),
195 )?;
196
197 local_handle
198 .wire(local_mid.clone(), PeerTarget::External(remote_spec))
199 .await
200 .map_err(CrossMobError::Mob)?;
201
202 if let Err(e) = remote_handle
203 .wire(remote_mid.clone(), PeerTarget::External(local_spec))
204 .await
205 {
206 if let Ok(rollback_spec) = build_peer_spec(
207 &remote_info.comms_name,
208 &remote_info.peer_id,
209 &entry.transport,
210 Some(remote_info.pubkey),
211 ) {
212 let _ = local_handle
213 .unwire(local_mid, PeerTarget::External(rollback_spec))
214 .await;
215 }
216 return Err(CrossMobError::Mob(e));
217 }
218
219 Ok(())
220 }
221 LocalOrRemote::Remote(proxy) => {
222 let (remote_peer_id, remote_comms_name) = proxy
230 .lookup_member(remote_member_id)
231 .await
232 .map_err(CrossMobError::Remote)?;
233 let remote_spec = build_peer_spec(
234 &remote_comms_name,
235 &remote_peer_id,
236 &entry.transport,
237 entry.pubkey,
238 )?;
239
240 local_handle
241 .wire(local_mid.clone(), PeerTarget::External(remote_spec))
242 .await
243 .map_err(CrossMobError::Mob)?;
244
245 let pubkey_b64 = self
253 .gateway_peer_keys
254 .as_ref()
255 .map(crate::auth::peer_keys::GatewayPeerKeys::pubkey_b64);
256 let local_spec_address = format!("inproc://{}", local_info.comms_name);
257 if let Err(remote_err) = proxy
258 .wire_remote(
259 remote_member_id,
260 &local_spec_address,
261 &local_info.comms_name,
262 &local_info.peer_id,
263 pubkey_b64,
264 )
265 .await
266 {
267 let rollback_spec = build_peer_spec(
268 &remote_comms_name,
269 &remote_peer_id,
270 &entry.transport,
271 entry.pubkey,
272 );
273 if let Ok(spec) = rollback_spec {
274 let _ = local_handle
275 .unwire(local_mid, PeerTarget::External(spec))
276 .await;
277 }
278 return Err(CrossMobError::Remote(remote_err));
279 }
280
281 Ok(())
282 }
283 }
284 }
285
286 pub async fn unwire_cross_mob(
292 &self,
293 local_member_id: &str,
294 remote_member_id: &str,
295 remote_mob_id: &str,
296 ) -> Result<(), CrossMobError> {
297 let entry = self.resolve_contact(remote_mob_id)?;
298 let remote = self.dispatch_for(&entry).await?;
299 let local_handle = self.mob_runtime.handle();
300 let local_mob_id = local_handle.mob_id().to_string();
301 let local_mid = crate::member_comms_id::mob_member_id(local_member_id);
302
303 let mut first_error: Option<CrossMobError> = None;
304
305 let local_info_opt = self
306 .get_member_peer_info(&local_handle, &local_mid, &local_mob_id)
307 .await
308 .ok();
309
310 match remote {
311 LocalOrRemote::Local(remote_handle) => {
312 let remote_mid = crate::member_comms_id::mob_member_id(remote_member_id);
313 if let Ok(remote_info) = self
314 .get_member_peer_info(&remote_handle, &remote_mid, remote_mob_id)
315 .await
316 && let Ok(spec) = build_peer_spec(
317 &remote_info.comms_name,
318 &remote_info.peer_id,
319 &entry.transport,
320 Some(remote_info.pubkey),
321 )
322 && let Err(e) = local_handle
323 .unwire(local_mid.clone(), PeerTarget::External(spec))
324 .await
325 {
326 first_error = Some(CrossMobError::Mob(e));
327 }
328
329 if let Some(local_info) = &local_info_opt
330 && let Ok(spec) = build_peer_spec(
331 &local_info.comms_name,
332 &local_info.peer_id,
333 &MobTransport::Inproc,
334 Some(local_info.pubkey),
335 )
336 && let Err(e) = remote_handle
337 .unwire(remote_mid.clone(), PeerTarget::External(spec))
338 .await
339 && first_error.is_none()
340 {
341 first_error = Some(CrossMobError::Mob(e));
342 }
343 }
344 LocalOrRemote::Remote(proxy) => {
345 if let Ok((remote_peer_id, remote_comms_name)) =
346 proxy.lookup_member(remote_member_id).await
347 && let Ok(spec) = build_peer_spec(
348 &remote_comms_name,
349 &remote_peer_id,
350 &entry.transport,
351 entry.pubkey,
352 )
353 && let Err(e) = local_handle
354 .unwire(local_mid.clone(), PeerTarget::External(spec))
355 .await
356 {
357 first_error = Some(CrossMobError::Mob(e));
358 }
359
360 if let Some(local_info) = &local_info_opt {
361 let pubkey_b64 = self
362 .gateway_peer_keys
363 .as_ref()
364 .map(crate::auth::peer_keys::GatewayPeerKeys::pubkey_b64);
365 let local_spec_address = format!("inproc://{}", local_info.comms_name);
366 if let Err(e) = proxy
367 .unwire_remote(
368 remote_member_id,
369 &local_spec_address,
370 &local_info.comms_name,
371 &local_info.peer_id,
372 pubkey_b64,
373 )
374 .await
375 && first_error.is_none()
376 {
377 first_error = Some(CrossMobError::Remote(e));
378 }
379 }
380 }
381 }
382
383 match first_error {
384 Some(e) => Err(e),
385 None => Ok(()),
386 }
387 }
388
389 pub async fn send_cross_mob(
400 &self,
401 from_local_member: &str,
402 remote_member_id: &str,
403 remote_mob_id: &str,
404 content: impl Into<meerkat_core::ContentInput>,
405 ) -> Result<String, CrossMobError> {
406 let entry = self.resolve_contact(remote_mob_id)?;
407 let remote = self.dispatch_for(&entry).await?;
408 let remote_mid = crate::member_comms_id::mob_member_id(remote_member_id);
409 let content = content.into();
410 let _ = from_local_member; match remote {
413 LocalOrRemote::Local(remote_handle) => {
414 let _receipt = remote_handle
415 .member(&remote_mid)
416 .await
417 .map_err(CrossMobError::Mob)?
418 .send(content, HandlingMode::Queue)
419 .await
420 .map_err(CrossMobError::Mob)?;
421 let session_id = remote_handle
425 .resolve_bridge_session_id(&remote_mid)
426 .await
427 .ok_or_else(|| CrossMobError::NoCommsInfo {
428 member_id: remote_member_id.to_string(),
429 mob_id: remote_mob_id.to_string(),
430 })?;
431 Ok(session_id.to_string())
432 }
433 LocalOrRemote::Remote(proxy) => {
434 let content_json = serde_json::to_value(&content).map_err(|err| {
439 CrossMobError::PeerSpec(format!(
440 "failed to serialize content for remote inject: {err}"
441 ))
442 })?;
443 let session_id = proxy
444 .inject_message(remote_member_id, content_json)
445 .await
446 .map_err(CrossMobError::Remote)?;
447 Ok(session_id)
448 }
449 }
450 }
451
452 pub fn list_external_mobs(&self) -> Vec<ContactEntry> {
454 self.contact_directory
455 .as_ref()
456 .map(|d| d.list().into_iter().cloned().collect())
457 .unwrap_or_default()
458 }
459
460 pub fn has_contact_directory(&self) -> bool {
462 self.contact_directory.is_some()
463 }
464
465 pub async fn has_peer_mob_handles(&self) -> bool {
468 !self.peer_mob_handles.read().await.is_empty()
469 }
470
471 pub fn has_inproc_contacts(&self) -> bool {
473 self.contact_directory.as_ref().is_some_and(|d| {
474 d.list()
475 .iter()
476 .any(|e| matches!(e.transport, MobTransport::Inproc))
477 })
478 }
479
480 pub fn has_remote_contacts(&self) -> bool {
483 self.contact_directory.as_ref().is_some_and(|d| {
484 d.list()
485 .iter()
486 .any(|e| matches!(e.transport, MobTransport::Tcp(_) | MobTransport::Uds(_)))
487 })
488 }
489
490 pub fn mob_id(&self) -> String {
492 self.mob_runtime.handle().mob_id().to_string()
493 }
494
495 pub async fn local_member_peer_info(
501 &self,
502 member_id: &str,
503 ) -> Result<(String, String, String), CrossMobError> {
504 let handle = self.mob_runtime.handle();
505 let mob_id = handle.mob_id().to_string();
506 let mid = crate::member_comms_id::mob_member_id(member_id);
507 let info = self.get_member_peer_info(&handle, &mid, &mob_id).await?;
508 let address = format!("inproc://{}", info.comms_name);
509 Ok((info.peer_id, info.comms_name, address))
510 }
511
512 pub async fn wire_local(
525 &self,
526 local_member_id: &str,
527 remote_comms_name: &str,
528 remote_peer_id: &str,
529 remote_address: &str,
530 remote_pubkey: Option<[u8; 32]>,
531 ) -> Result<(), CrossMobError> {
532 let spec = build_external_peer_spec(
533 remote_comms_name,
534 remote_peer_id,
535 remote_address,
536 remote_pubkey,
537 )?;
538 let local_mid = crate::member_comms_id::mob_member_id(local_member_id);
539 self.mob_runtime
540 .handle()
541 .wire(local_mid, PeerTarget::External(spec))
542 .await
543 .map_err(CrossMobError::Mob)
544 }
545
546 pub async fn unwire_local(
549 &self,
550 local_member_id: &str,
551 remote_comms_name: &str,
552 remote_peer_id: &str,
553 remote_address: &str,
554 remote_pubkey: Option<[u8; 32]>,
555 ) -> Result<(), CrossMobError> {
556 let spec = build_external_peer_spec(
557 remote_comms_name,
558 remote_peer_id,
559 remote_address,
560 remote_pubkey,
561 )?;
562 let local_mid = crate::member_comms_id::mob_member_id(local_member_id);
563 self.mob_runtime
564 .handle()
565 .unwire(local_mid, PeerTarget::External(spec))
566 .await
567 .map_err(CrossMobError::Mob)
568 }
569
570 fn resolve_contact(&self, mob_id: &str) -> Result<ContactEntry, CrossMobError> {
573 let dir = self
574 .contact_directory
575 .as_ref()
576 .ok_or(CrossMobError::NoContactDirectory)?;
577 dir.get(mob_id)
578 .cloned()
579 .ok_or_else(|| CrossMobError::UnknownMob(mob_id.to_string()))
580 }
581
582 async fn dispatch_for(&self, entry: &ContactEntry) -> Result<LocalOrRemote, CrossMobError> {
593 if let Some(handle) = self
594 .peer_mob_handles
595 .read()
596 .await
597 .get(&entry.mob_id)
598 .cloned()
599 {
600 return Ok(LocalOrRemote::Local(handle));
601 }
602 match RemoteMobProxy::from_entry(entry)? {
603 Some(proxy) => Ok(LocalOrRemote::Remote(proxy)),
604 None => Err(CrossMobError::NoPeerHandle(entry.mob_id.clone())),
605 }
606 }
607
608 async fn get_member_peer_info(
619 &self,
620 handle: &MobHandle,
621 meerkat_id: &AgentIdentity,
622 mob_id: &str,
623 ) -> Result<MemberPeerInfo, CrossMobError> {
624 let entry = handle
625 .get_member(meerkat_id)
626 .await
627 .map_err(|err| {
628 CrossMobError::PeerSpec(format!(
629 "member lookup for '{meerkat_id}' in mob '{mob_id}' failed: {err}"
630 ))
631 })?
632 .ok_or_else(|| CrossMobError::MemberNotFound {
633 member_id: meerkat_id.to_string(),
634 mob_id: mob_id.to_string(),
635 })?;
636 let peer_id = entry
637 .peer_id()
638 .ok_or_else(|| CrossMobError::NoCommsInfo {
639 member_id: meerkat_id.to_string(),
640 mob_id: mob_id.to_string(),
641 })?
642 .to_string();
643 let pubkey_b64 = entry.transport_public_key().ok_or_else(|| {
644 CrossMobError::PeerSpec(format!(
645 "member '{meerkat_id}' in mob '{mob_id}' has no transport public key"
646 ))
647 })?;
648 let pubkey = crate::auth::peer_keys::decode_pubkey_b64(pubkey_b64).map_err(|err| {
649 CrossMobError::PeerSpec(format!(
650 "member '{meerkat_id}' in mob '{mob_id}' has invalid transport public key: {err}"
651 ))
652 })?;
653 let comms_name = meerkat_core::MemberCommsName::new(
654 mob_id,
655 entry.role.as_str(),
656 meerkat_id.as_str(),
657 )
658 .map_err(|err| {
659 CrossMobError::PeerSpec(format!(
660 "member '{meerkat_id}' in mob '{mob_id}' has an invalid comms name component: {err}"
661 ))
662 })?
663 .to_string();
664 Ok(MemberPeerInfo {
665 peer_id,
666 comms_name,
667 pubkey,
668 })
669 }
670}
671
672fn build_peer_spec(
682 comms_name: &str,
683 peer_id: &str,
684 transport: &MobTransport,
685 pubkey: Option<[u8; 32]>,
686) -> Result<TrustedPeerDescriptor, CrossMobError> {
687 let address = match transport {
688 MobTransport::Inproc => format!("inproc://{comms_name}"),
689 MobTransport::Tcp(addr) => format!("tcp://{addr}"),
690 MobTransport::Uds(path) => format!("uds://{path}"),
691 };
692 build_external_peer_spec(comms_name, peer_id, &address, pubkey)
693}
694
695fn build_external_peer_spec(
706 comms_name: &str,
707 peer_id: &str,
708 address: &str,
709 pubkey: Option<[u8; 32]>,
710) -> Result<TrustedPeerDescriptor, CrossMobError> {
711 let is_inproc = address.starts_with("inproc://");
712 match (is_inproc, pubkey) {
713 (true, None) => TrustedPeerDescriptor::test_only_unsigned(comms_name, peer_id, address)
714 .map_err(CrossMobError::PeerSpec),
715 (true, Some(bytes)) => {
716 TrustedPeerDescriptor::unsigned_with_pubkey(comms_name, peer_id, bytes, address)
717 .map_err(CrossMobError::PeerSpec)
718 }
719 (false, None) => Err(CrossMobError::MissingPeerPubkey { mob_id: None }),
720 (false, Some(bytes)) => {
721 if bytes == [0u8; 32] {
722 return Err(CrossMobError::MissingPeerPubkey { mob_id: None });
723 }
724 TrustedPeerDescriptor::unsigned_with_pubkey(comms_name, peer_id, bytes, address)
725 .map_err(CrossMobError::PeerSpec)
726 }
727 }
728}
729
730pub fn build_tcp_peer_spec(
737 comms_name: &str,
738 peer_id: &str,
739 address: &str,
740) -> Result<TrustedPeerDescriptor, CrossMobError> {
741 TrustedPeerDescriptor::test_only_unsigned(comms_name, peer_id, format!("tcp://{address}"))
742 .map_err(CrossMobError::PeerSpec)
743}
744
745pub fn build_uds_peer_spec(
751 comms_name: &str,
752 peer_id: &str,
753 path: &str,
754) -> Result<TrustedPeerDescriptor, CrossMobError> {
755 let normalized = if let Some(stripped) = path.strip_prefix('/') {
756 stripped
757 } else {
758 path
759 };
760 TrustedPeerDescriptor::test_only_unsigned(comms_name, peer_id, format!("uds:///{normalized}"))
761 .map_err(CrossMobError::PeerSpec)
762}
763
764#[cfg(test)]
765#[allow(clippy::unwrap_used, clippy::expect_used)]
766mod tests {
767 use super::*;
768
769 const TEST_PEER_ID: &str = "00000000-0000-4000-8000-000000000001";
770
771 const TEST_PUBKEY: [u8; 32] = [42u8; 32];
773
774 fn derived_peer_id() -> String {
779 meerkat_core::comms::PeerId::from_ed25519_pubkey(&TEST_PUBKEY).to_string()
780 }
781
782 #[test]
783 fn peer_spec_inproc_uses_comms_name_address() {
784 let spec = build_peer_spec(
785 "authors/coordinator/alice",
786 TEST_PEER_ID,
787 &MobTransport::Inproc,
788 None,
789 )
790 .expect("spec");
791 assert_eq!(spec.address.endpoint(), "authors/coordinator/alice");
792 }
793
794 #[test]
795 fn peer_spec_tcp_uses_tcp_scheme() {
796 let id = derived_peer_id();
797 let spec = build_peer_spec(
798 "authors/coordinator/alice",
799 &id,
800 &MobTransport::Tcp("127.0.0.1:9001".to_string()),
801 Some(TEST_PUBKEY),
802 )
803 .expect("spec");
804 assert_eq!(spec.address.endpoint(), "127.0.0.1:9001");
805 }
806
807 #[test]
808 fn peer_spec_uds_uses_uds_scheme() {
809 let id = derived_peer_id();
810 let spec = build_peer_spec(
811 "authors/coordinator/alice",
812 &id,
813 &MobTransport::Uds("/tmp/x.sock".to_string()),
814 Some(TEST_PUBKEY),
815 )
816 .expect("spec");
817 assert_eq!(spec.address.endpoint(), "/tmp/x.sock");
818 }
819
820 #[test]
824 fn peer_spec_tcp_without_pubkey_rejected() {
825 let result = build_peer_spec(
826 "authors/coordinator/alice",
827 TEST_PEER_ID,
828 &MobTransport::Tcp("127.0.0.1:9001".to_string()),
829 None,
830 );
831 assert!(
832 matches!(result, Err(CrossMobError::MissingPeerPubkey { .. })),
833 "TCP peer spec without pubkey must fail closed, got {result:?}"
834 );
835 }
836
837 #[test]
839 fn peer_spec_uds_without_pubkey_rejected() {
840 let result = build_peer_spec(
841 "authors/coordinator/alice",
842 TEST_PEER_ID,
843 &MobTransport::Uds("/tmp/x.sock".to_string()),
844 None,
845 );
846 assert!(
847 matches!(result, Err(CrossMobError::MissingPeerPubkey { .. })),
848 "UDS peer spec without pubkey must fail closed, got {result:?}"
849 );
850 }
851
852 #[test]
853 fn build_uds_peer_spec_handles_leading_slash() {
854 let with = build_uds_peer_spec("a", "00000000-0000-4000-8000-000000000001", "/tmp/x.sock")
857 .expect("spec");
858 let without =
859 build_uds_peer_spec("a", "00000000-0000-4000-8000-000000000001", "tmp/x.sock")
860 .expect("spec");
861 assert_eq!(with.address.endpoint(), "/tmp/x.sock");
862 assert_eq!(without.address.endpoint(), "/tmp/x.sock");
863 }
864}