use meerkat_comms::{InprocRegistry, PeerMeta, PubKey};
use meerkat_core::comms::TrustedPeerDescriptor;
use meerkat_core::types::HandlingMode;
use meerkat_mob::ids::AgentIdentity;
use meerkat_mob::{MobHandle, PeerTarget};
use crate::auth::peer_keys::GatewayPeerKeys;
use crate::contact_directory::{ContactDirectory, ContactEntry, MobTransport};
use crate::runtime::cross_mob_remote::{RemoteMemberInfo, RemoteMobError, RemoteMobProxy};
use super::UnifiedRuntime;
enum LocalOrRemote {
Local(Box<PeerMobAuthority>),
Remote(RemoteMobProxy),
}
#[derive(Clone)]
pub(crate) struct PeerMobAuthority {
handle: MobHandle,
identity_runtime: Option<std::sync::Arc<crate::identity_first::IdentityRuntime>>,
}
struct MemberPeerInfo {
peer_id: String,
comms_name: String,
pubkey: [u8; 32],
pubkey_b64: String,
}
struct BilateralAliasRollback<'a> {
local_namespace: &'a str,
peer_namespace: &'a str,
local_pubkey: [u8; 32],
peer_pubkey: [u8; 32],
local_alias_preexisting: bool,
peer_alias_preexisting: bool,
}
struct BilateralWireRollback<'a> {
local_handle: MobHandle,
peer_handle: MobHandle,
local_member_id: &'a str,
peer_member_id: &'a str,
peer_spec: &'a TrustedPeerDescriptor,
local_spec: &'a TrustedPeerDescriptor,
local_mutated: bool,
peer_mutated: bool,
aliases: Option<BilateralAliasRollback<'a>>,
}
async fn rollback_bilateral_wire_attempt(attempt: BilateralWireRollback<'_>) -> Vec<String> {
let BilateralWireRollback {
local_handle,
peer_handle,
local_member_id,
peer_member_id,
peer_spec,
local_spec,
local_mutated,
peer_mutated,
aliases,
} = attempt;
let mut failures = Vec::new();
if let Some(aliases) = aliases {
let registry = InprocRegistry::global();
if !aliases.local_alias_preexisting {
registry.unregister_in_namespace(
aliases.local_namespace,
&PubKey::new(aliases.peer_pubkey),
);
}
if !aliases.peer_alias_preexisting {
registry.unregister_in_namespace(
aliases.peer_namespace,
&PubKey::new(aliases.local_pubkey),
);
}
}
if peer_mutated
&& let Err(error) = mutate_member_unchecked(
peer_handle,
peer_member_id,
PeerTarget::External(local_spec.clone()),
false,
)
.await
{
failures.push(format!("target trust rollback failed: {error}"));
}
if local_mutated
&& let Err(error) = mutate_member_unchecked(
local_handle,
local_member_id,
PeerTarget::External(peer_spec.clone()),
false,
)
.await
{
failures.push(format!("source trust rollback failed: {error}"));
}
failures
}
#[derive(Debug)]
pub enum CrossMobError {
NoContactDirectory,
UnknownMob(String),
NoPeerHandle(String),
MissingPeerPubkey { mob_id: Option<String> },
MemberNotFound { member_id: String, mob_id: String },
NoCommsInfo { member_id: String, mob_id: String },
Mob(meerkat_mob::MobError),
IdentityAuthority {
member_id: String,
mob_id: String,
message: String,
},
PeerSpec(String),
InprocAlias(String),
Remote(RemoteMobError),
ControlListener(String),
LocalMemberNotRemotelyAddressable {
member_id: String,
mob_id: String,
advertised: Option<String>,
},
RemoteMemberNotRemotelyAddressable {
member_id: String,
mob_id: String,
detail: String,
},
}
impl std::fmt::Display for CrossMobError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::NoContactDirectory => write!(f, "no contact directory configured"),
Self::UnknownMob(id) => write!(f, "unknown mob: {id}"),
Self::NoPeerHandle(id) => write!(f, "no peer mob handle registered for: {id}"),
Self::Remote(e) => write!(f, "cross-process cross-mob: {e}"),
Self::MemberNotFound { member_id, mob_id } => {
write!(f, "member '{member_id}' not found in mob '{mob_id}'")
}
Self::NoCommsInfo { member_id, mob_id } => {
write!(
f,
"member '{member_id}' in mob '{mob_id}' has no comms runtime"
)
}
Self::Mob(err) => write!(f, "mob error: {err}"),
Self::IdentityAuthority {
member_id,
mob_id,
message,
} => write!(
f,
"identity authority rejected member '{member_id}' in mob '{mob_id}': {message}"
),
Self::PeerSpec(reason) => write!(f, "peer spec error: {reason}"),
Self::InprocAlias(reason) => write!(f, "inproc alias error: {reason}"),
Self::ControlListener(reason) => {
write!(f, "cross-mob control listener error: {reason}")
}
Self::LocalMemberNotRemotelyAddressable {
member_id,
mob_id,
advertised,
} => write!(
f,
"local member '{member_id}' in mob '{mob_id}' is not remotely addressable \
(its comms runtime advertises {}); remote cross-mob wiring needs a member \
socket transport - set runtime_options.member_comms_address (rpc_gateway) or \
comms.mode=\"tcp\" with comms.address in the meerkat config, and bind \
control_listen / --control-listen on the peer gateways",
match advertised {
Some(address) => format!("'{address}'"),
None => "no address (no comms runtime)".to_string(),
}
),
Self::RemoteMemberNotRemotelyAddressable {
member_id,
mob_id,
detail,
} => write!(
f,
"remote member '{member_id}' in mob '{mob_id}' is not remotely addressable: \
{detail}; the peer gateway must run its members with a socket comms transport \
and serve LookupMember with session-service access"
),
Self::MissingPeerPubkey { mob_id } => match mob_id {
Some(id) => write!(
f,
"non-inproc peer for mob '{id}' has no signing pubkey; \
bootstrap via mobkit/peer_pubkey or populate the contact \
directory's pubkey field before wiring"
),
None => write!(
f,
"non-inproc peer has no signing pubkey; supply a 32-byte \
Ed25519 pubkey or use inproc transport"
),
},
}
}
}
async fn member_lifecycle_target_with_authority(
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
member_id: &str,
mob_id: &str,
) -> Result<Option<crate::identity_first::runtime::MemberAliasLifecycleTarget>, CrossMobError> {
if let Some(runtime) = identity_runtime {
return runtime
.member_alias_lifecycle_target(member_id)
.await
.map_err(|error| CrossMobError::IdentityAuthority {
member_id: member_id.to_string(),
mob_id: mob_id.to_string(),
message: error.to_string(),
});
}
if crate::member_comms_id::is_reserved_generated_alias(member_id) {
return Err(CrossMobError::IdentityAuthority {
member_id: member_id.to_string(),
mob_id: mob_id.to_string(),
message: "generated aliases require the owning IdentityRuntime".to_string(),
});
}
Ok(None)
}
async fn resolve_member_alias_under_authority(
handle: &MobHandle,
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
identity_authoritative: bool,
member_alias: &str,
mob_id: &str,
) -> Result<String, CrossMobError> {
let member_alias = crate::member_comms_id::runtime_alias_str(member_alias).into_owned();
if identity_authoritative {
let runtime = identity_runtime.ok_or_else(|| CrossMobError::IdentityAuthority {
member_id: member_alias.clone(),
mob_id: mob_id.to_string(),
message: "durable member target lost its IdentityRuntime authority".to_string(),
})?;
let identity = crate::identity_first::IdentityRuntime::identity_for_generated_member_alias(
&member_alias,
)
.or_else(|| crate::identity_first::AgentIdentity::parse(&member_alias).ok())
.ok_or_else(|| CrossMobError::IdentityAuthority {
member_id: member_alias.clone(),
mob_id: mob_id.to_string(),
message: "durable member target is not a valid identity alias".to_string(),
})?;
let status =
runtime
.status(&identity)
.await
.map_err(|error| CrossMobError::IdentityAuthority {
member_id: member_alias.clone(),
mob_id: mob_id.to_string(),
message: error.to_string(),
})?;
return status
.agent_runtime_id
.map(|runtime_id| runtime_id.to_string())
.ok_or_else(|| CrossMobError::IdentityAuthority {
member_id: member_alias,
mob_id: mob_id.to_string(),
message: format!("identity {identity} has no current runtime member"),
});
}
let direct = crate::member_comms_id::mob_member_id(&member_alias);
if handle
.get_member(&direct)
.await
.map_err(CrossMobError::Mob)?
.is_some()
{
return Ok(member_alias);
}
let candidates = handle
.list_members_including_retiring()
.await
.into_iter()
.filter(|entry| {
crate::member_comms_id::durable_identity_label(&entry.labels)
.is_some_and(|identity| identity == member_alias)
})
.map(|entry| {
crate::member_comms_id::runtime_alias_str(entry.agent_identity.as_str()).into_owned()
})
.collect::<std::collections::BTreeSet<_>>();
match candidates.len() {
0 => Ok(member_alias),
1 => candidates.into_iter().next().ok_or_else(|| {
CrossMobError::PeerSpec("member alias candidate disappeared".to_string())
}),
_ => Err(CrossMobError::IdentityAuthority {
member_id: member_alias.clone(),
mob_id: mob_id.to_string(),
message: format!(
"ambiguous durable member alias {member_alias}: candidates [{}]",
candidates.into_iter().collect::<Vec<_>>().join(", ")
),
}),
}
}
async fn run_member_authority_transaction<T, F, Fut>(
targets: impl IntoIterator<
Item = Option<crate::identity_first::runtime::MemberAliasLifecycleTarget>,
>,
member_context: impl Into<String>,
mob_context: impl Into<String>,
operation: F,
) -> Result<T, CrossMobError>
where
T: Send + 'static,
F: FnOnce() -> Fut + Send + 'static,
Fut: std::future::Future<Output = Result<T, CrossMobError>> + Send + 'static,
{
let targets = targets.into_iter().flatten().collect::<Vec<_>>();
if targets.is_empty() {
return operation().await;
}
let member_context = member_context.into();
let mob_context = mob_context.into();
let operation_error = std::sync::Arc::new(std::sync::Mutex::new(None));
let task_operation_error = std::sync::Arc::clone(&operation_error);
let result =
crate::identity_first::IdentityRuntime::run_member_alias_targets_operation_tracked(
targets,
move || async move {
match operation().await {
Ok(value) => Ok(value),
Err(error) => {
*task_operation_error
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(error);
Err("cross-mob transaction failed".to_string())
}
}
},
)
.await;
match result {
Ok(value) => Ok(value),
Err(authority_error) => {
if let Some(operation_error) = operation_error
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
{
return Err(operation_error);
}
Err(CrossMobError::IdentityAuthority {
member_id: member_context,
mob_id: mob_context,
message: authority_error.to_string(),
})
}
}
}
async fn mutate_member_unchecked(
handle: MobHandle,
member_id: &str,
peer: PeerTarget,
wire: bool,
) -> Result<(), CrossMobError> {
let member_id = crate::member_comms_id::runtime_alias_str(member_id).into_owned();
let member = crate::member_comms_id::mob_member_id(&member_id);
let result = if wire {
handle.wire(member, peer).await
} else {
handle.unwire(member, peer).await
};
result.map_err(CrossMobError::Mob)
}
async fn send_member_unchecked(
handle: MobHandle,
member_id: &str,
mob_id: &str,
content: meerkat_core::ContentInput,
) -> Result<String, CrossMobError> {
let member_id = crate::member_comms_id::runtime_alias_str(member_id).into_owned();
let member = crate::member_comms_id::mob_member_id(&member_id);
handle
.member(&member)
.await
.map_err(CrossMobError::Mob)?
.send(content, HandlingMode::Queue)
.await
.map_err(CrossMobError::Mob)?;
handle
.resolve_bridge_session_id(&member)
.await
.map(|session_id| session_id.to_string())
.ok_or_else(|| CrossMobError::NoCommsInfo {
member_id,
mob_id: mob_id.to_string(),
})
}
async fn member_peer_info(
handle: &MobHandle,
meerkat_id: &AgentIdentity,
mob_id: &str,
) -> Result<MemberPeerInfo, CrossMobError> {
let entry = handle
.get_member(meerkat_id)
.await
.map_err(|err| {
CrossMobError::PeerSpec(format!(
"member lookup for '{meerkat_id}' in mob '{mob_id}' failed: {err}"
))
})?
.ok_or_else(|| CrossMobError::MemberNotFound {
member_id: meerkat_id.to_string(),
mob_id: mob_id.to_string(),
})?;
let peer_id = entry
.peer_id()
.ok_or_else(|| CrossMobError::NoCommsInfo {
member_id: meerkat_id.to_string(),
mob_id: mob_id.to_string(),
})?
.to_string();
let pubkey_b64 = entry.transport_public_key().ok_or_else(|| {
CrossMobError::PeerSpec(format!(
"member '{meerkat_id}' in mob '{mob_id}' has no transport public key"
))
})?;
let pubkey = crate::auth::peer_keys::decode_pubkey_b64(pubkey_b64).map_err(|err| {
CrossMobError::PeerSpec(format!(
"member '{meerkat_id}' in mob '{mob_id}' has invalid transport public key: {err}"
))
})?;
let comms_name = meerkat_core::MemberCommsName::new(
mob_id,
entry.role.as_str(),
meerkat_id.as_str(),
)
.map_err(|err| {
CrossMobError::PeerSpec(format!(
"member '{meerkat_id}' in mob '{mob_id}' has an invalid comms name component: {err}"
))
})?
.to_string();
Ok(MemberPeerInfo {
peer_id,
comms_name,
pubkey,
pubkey_b64: pubkey_b64.to_string(),
})
}
async fn member_comms_advertised_address(
mob_runtime: &crate::MobRuntime,
handle: &MobHandle,
member: &AgentIdentity,
) -> Option<String> {
let session_id = handle.resolve_bridge_session_id(member).await?;
let service = mob_runtime.session_service()?;
let comms = service.comms_runtime(&session_id).await?;
comms.advertised_address()
}
fn build_remote_member_spec(
info: &RemoteMemberInfo,
remote_member_id: &str,
remote_mob_id: &str,
) -> Result<TrustedPeerDescriptor, CrossMobError> {
let not_addressable = |detail: String| CrossMobError::RemoteMemberNotRemotelyAddressable {
member_id: remote_member_id.to_string(),
mob_id: remote_mob_id.to_string(),
detail,
};
let address = info
.advertised_address
.as_deref()
.filter(|address| !address.starts_with("inproc://"))
.ok_or_else(|| {
not_addressable(match &info.advertised_address {
Some(address) => {
format!("advertised address '{address}' is not dialable across processes")
}
None => "the peer gateway reported no advertised comms address for this member"
.to_string(),
})
})?;
let pubkey_b64 = info.pubkey_b64.as_deref().ok_or_else(|| {
not_addressable("the peer gateway reported no transport pubkey for this member".to_string())
})?;
let pubkey = crate::auth::peer_keys::decode_pubkey_b64(pubkey_b64).map_err(|err| {
CrossMobError::PeerSpec(format!(
"remote member '{remote_member_id}' in mob '{remote_mob_id}' has an invalid \
transport pubkey: {err}"
))
})?;
build_external_peer_spec(&info.comms_name, &info.peer_id, address, Some(pubkey))
}
async fn member_can_address_peer(
mob_runtime: &crate::MobRuntime,
handle: &MobHandle,
local_member: &AgentIdentity,
expected_peer: &MemberPeerInfo,
) -> Result<bool, CrossMobError> {
let session_id = handle
.resolve_bridge_session_id(local_member)
.await
.ok_or_else(|| CrossMobError::NoCommsInfo {
member_id: local_member.to_string(),
mob_id: handle.mob_id().to_string(),
})?;
let service = mob_runtime
.session_service()
.ok_or_else(|| CrossMobError::NoCommsInfo {
member_id: local_member.to_string(),
mob_id: handle.mob_id().to_string(),
})?;
let comms =
service
.comms_runtime(&session_id)
.await
.ok_or_else(|| CrossMobError::NoCommsInfo {
member_id: local_member.to_string(),
mob_id: handle.mob_id().to_string(),
})?;
let peers = comms.peers().await;
Ok(peers.iter().any(|peer| {
peer.peer_id.to_string() == expected_peer.peer_id
&& peer.name.as_str() == expected_peer.comms_name
&& peer
.sendable_kinds
.contains(&meerkat_core::comms::PeerSendability::PeerMessage)
}))
}
async fn wire_member_with_authority(
handle: MobHandle,
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
member_id: &str,
mob_id: &str,
peer: PeerTarget,
wire: bool,
) -> Result<(), CrossMobError> {
let operation_member_id = crate::member_comms_id::runtime_alias_str(member_id).into_owned();
let operation_mob_id = mob_id.to_string();
let identity_runtime = identity_runtime.cloned();
let target = member_lifecycle_target_with_authority(
identity_runtime.as_ref(),
&operation_member_id,
&operation_mob_id,
)
.await?;
let identity_authoritative = target.is_some();
run_member_authority_transaction(
[target],
operation_member_id.clone(),
operation_mob_id.clone(),
move || async move {
let current_member_id = resolve_member_alias_under_authority(
&handle,
identity_runtime.as_ref(),
identity_authoritative,
&operation_member_id,
&operation_mob_id,
)
.await?;
mutate_member_unchecked(handle, ¤t_member_id, peer, wire).await
},
)
.await
}
async fn wire_bilateral_transaction(
local_runtime: crate::MobRuntime,
peer_runtime: crate::MobRuntime,
local_member_id: String,
peer_member_id: String,
) -> Result<(), CrossMobError> {
let local_handle = local_runtime.handle();
let peer_handle = peer_runtime.handle();
let local_mob_id = local_handle.mob_id().to_string();
let peer_mob_id = peer_handle.mob_id().to_string();
let local_mid = crate::member_comms_id::mob_member_id(&local_member_id);
let peer_mid = crate::member_comms_id::mob_member_id(&peer_member_id);
let local_info = member_peer_info(&local_handle, &local_mid, &local_mob_id).await?;
let peer_info = member_peer_info(&peer_handle, &peer_mid, &peer_mob_id).await?;
let peer_spec = build_peer_spec(
&peer_info.comms_name,
&peer_info.peer_id,
&MobTransport::Inproc,
Some(peer_info.pubkey),
)?;
let local_spec = build_peer_spec(
&local_info.comms_name,
&local_info.peer_id,
&MobTransport::Inproc,
Some(local_info.pubkey),
)?;
let local_member = local_handle.get_member(&local_mid).await?.ok_or_else(|| {
CrossMobError::MemberNotFound {
member_id: local_member_id.clone(),
mob_id: local_mob_id.clone(),
}
})?;
let peer_member =
peer_handle
.get_member(&peer_mid)
.await?
.ok_or_else(|| CrossMobError::MemberNotFound {
member_id: peer_member_id.clone(),
mob_id: peer_mob_id.clone(),
})?;
let local_wired = local_member
.wired_to
.iter()
.any(|identity| identity.as_str() == peer_info.comms_name);
let peer_wired = peer_member
.wired_to
.iter()
.any(|identity| identity.as_str() == local_info.comms_name);
let local_agent_ready =
member_can_address_peer(&local_runtime, &local_handle, &local_mid, &peer_info).await?;
let peer_agent_ready =
member_can_address_peer(&peer_runtime, &peer_handle, &peer_mid, &local_info).await?;
let local_namespace = mob_inproc_namespace_for_id(&local_mob_id)?;
let peer_namespace = mob_inproc_namespace_for_id(&peer_mob_id)?;
let local_alias_preexisting = alias_already_installed(
&local_namespace,
&peer_info.comms_name,
PubKey::new(peer_info.pubkey),
);
let peer_alias_preexisting = alias_already_installed(
&peer_namespace,
&local_info.comms_name,
PubKey::new(local_info.pubkey),
);
let mut local_mutated = false;
let mut peer_mutated = false;
if !local_wired || !local_agent_ready {
mutate_member_unchecked(
local_handle.clone(),
&local_member_id,
PeerTarget::External(peer_spec.clone()),
true,
)
.await?;
local_mutated = true;
}
if (!peer_wired || !peer_agent_ready)
&& let Err(error) = mutate_member_unchecked(
peer_handle.clone(),
&peer_member_id,
PeerTarget::External(local_spec.clone()),
true,
)
.await
{
let rollback_failures = rollback_bilateral_wire_attempt(BilateralWireRollback {
local_handle: local_handle.clone(),
peer_handle: peer_handle.clone(),
local_member_id: &local_member_id,
peer_member_id: &peer_member_id,
peer_spec: &peer_spec,
local_spec: &local_spec,
local_mutated,
peer_mutated: false,
aliases: None,
})
.await;
if rollback_failures.is_empty() {
return Err(error);
}
return Err(CrossMobError::InprocAlias(format!(
"target wire failed: {error}; rollback failures: {}",
rollback_failures.join("; ")
)));
} else if !peer_wired || !peer_agent_ready {
peer_mutated = true;
}
if let Err(error) = register_cross_namespace_aliases(
&local_namespace,
&local_info.comms_name,
local_info.pubkey,
&peer_namespace,
&peer_info.comms_name,
peer_info.pubkey,
) {
let rollback_failures = rollback_bilateral_wire_attempt(BilateralWireRollback {
local_handle: local_handle.clone(),
peer_handle: peer_handle.clone(),
local_member_id: &local_member_id,
peer_member_id: &peer_member_id,
peer_spec: &peer_spec,
local_spec: &local_spec,
local_mutated,
peer_mutated,
aliases: Some(BilateralAliasRollback {
local_namespace: &local_namespace,
peer_namespace: &peer_namespace,
local_pubkey: local_info.pubkey,
peer_pubkey: peer_info.pubkey,
local_alias_preexisting,
peer_alias_preexisting,
}),
})
.await;
if rollback_failures.is_empty() {
return Err(CrossMobError::InprocAlias(error));
}
return Err(CrossMobError::InprocAlias(format!(
"{error}; rollback failures: {}",
rollback_failures.join("; ")
)));
}
let readiness = async {
let local_ready =
member_can_address_peer(&local_runtime, &local_handle, &local_mid, &peer_info).await?;
let peer_ready =
member_can_address_peer(&peer_runtime, &peer_handle, &peer_mid, &local_info).await?;
Ok::<_, CrossMobError>((local_ready, peer_ready))
}
.await;
let readiness_error = match readiness {
Ok((true, true)) => return Ok(()),
Ok((local_ready, peer_ready)) => format!(
"wire completed without agent-facing trust directory convergence: \
local_ready={local_ready}, peer_ready={peer_ready}"
),
Err(error) => format!("wire readiness inspection failed: {error}"),
};
let rollback_failures = rollback_bilateral_wire_attempt(BilateralWireRollback {
local_handle,
peer_handle,
local_member_id: &local_member_id,
peer_member_id: &peer_member_id,
peer_spec: &peer_spec,
local_spec: &local_spec,
local_mutated,
peer_mutated,
aliases: Some(BilateralAliasRollback {
local_namespace: &local_namespace,
peer_namespace: &peer_namespace,
local_pubkey: local_info.pubkey,
peer_pubkey: peer_info.pubkey,
local_alias_preexisting,
peer_alias_preexisting,
}),
})
.await;
if rollback_failures.is_empty() {
return Err(CrossMobError::InprocAlias(readiness_error));
}
Err(CrossMobError::InprocAlias(format!(
"{readiness_error}; rollback failures: {}",
rollback_failures.join("; ")
)))
}
async fn unwire_bilateral_transaction(
local_runtime: crate::MobRuntime,
peer_runtime: crate::MobRuntime,
local_member_id: String,
peer_member_id: String,
) -> Result<(), CrossMobError> {
let local_handle = local_runtime.handle();
let peer_handle = peer_runtime.handle();
let local_mob_id = local_handle.mob_id().to_string();
let peer_mob_id = peer_handle.mob_id().to_string();
let local_mid = crate::member_comms_id::mob_member_id(&local_member_id);
let peer_mid = crate::member_comms_id::mob_member_id(&peer_member_id);
let local_info = member_peer_info(&local_handle, &local_mid, &local_mob_id).await?;
let peer_info = member_peer_info(&peer_handle, &peer_mid, &peer_mob_id).await?;
let peer_spec = build_peer_spec(
&peer_info.comms_name,
&peer_info.peer_id,
&MobTransport::Inproc,
Some(peer_info.pubkey),
)?;
let local_spec = build_peer_spec(
&local_info.comms_name,
&local_info.peer_id,
&MobTransport::Inproc,
Some(local_info.pubkey),
)?;
let local_member = local_handle.get_member(&local_mid).await?.ok_or_else(|| {
CrossMobError::MemberNotFound {
member_id: local_member_id.clone(),
mob_id: local_mob_id.clone(),
}
})?;
let peer_member =
peer_handle
.get_member(&peer_mid)
.await?
.ok_or_else(|| CrossMobError::MemberNotFound {
member_id: peer_member_id.clone(),
mob_id: peer_mob_id.clone(),
})?;
let local_wired = local_member
.wired_to
.iter()
.any(|identity| identity.as_str() == peer_info.comms_name);
let peer_wired = peer_member
.wired_to
.iter()
.any(|identity| identity.as_str() == local_info.comms_name);
if local_wired {
mutate_member_unchecked(
local_handle.clone(),
&local_member_id,
PeerTarget::External(peer_spec.clone()),
false,
)
.await?;
}
if peer_wired
&& let Err(error) = mutate_member_unchecked(
peer_handle.clone(),
&peer_member_id,
PeerTarget::External(local_spec),
false,
)
.await
{
let rollback = if local_wired {
mutate_member_unchecked(
local_handle.clone(),
&local_member_id,
PeerTarget::External(peer_spec),
true,
)
.await
} else {
Ok(())
};
return match rollback {
Ok(()) => Err(error),
Err(rollback_error) => Err(CrossMobError::InprocAlias(format!(
"target unwire failed: {error}; source rollback failed: {rollback_error}"
))),
};
}
let local_namespace = mob_inproc_namespace_for_id(&local_mob_id)?;
let peer_namespace = mob_inproc_namespace_for_id(&peer_mob_id)?;
let local_still_references_peer =
handle_references_peer(&local_handle, &peer_info.comms_name).await;
let peer_still_references_local =
handle_references_peer(&peer_handle, &local_info.comms_name).await;
unregister_cross_namespace_aliases(
&local_namespace,
local_info.pubkey,
peer_info.pubkey,
local_still_references_peer,
&peer_namespace,
peer_still_references_local,
);
Ok(())
}
#[allow(clippy::too_many_arguments)]
async fn wire_cross_mob_transaction(
entry: ContactEntry,
remote: LocalOrRemote,
local_runtime: crate::MobRuntime,
local_mid: AgentIdentity,
local_mob_id: String,
local_member_id: String,
remote_member_id: String,
remote_mob_id: String,
) -> Result<(), CrossMobError> {
let local_handle = local_runtime.handle();
let local_info = member_peer_info(&local_handle, &local_mid, &local_mob_id).await?;
match remote {
LocalOrRemote::Local(remote_authority) => {
let remote_handle = remote_authority.handle;
let remote_mid = crate::member_comms_id::mob_member_id(&remote_member_id);
let remote_info = member_peer_info(&remote_handle, &remote_mid, &remote_mob_id).await?;
let remote_spec = build_peer_spec(
&remote_info.comms_name,
&remote_info.peer_id,
&entry.transport,
Some(remote_info.pubkey),
)?;
let local_spec = build_peer_spec(
&local_info.comms_name,
&local_info.peer_id,
&MobTransport::Inproc,
Some(local_info.pubkey),
)?;
mutate_member_unchecked(
local_handle.clone(),
&local_member_id,
PeerTarget::External(remote_spec),
true,
)
.await?;
if let Err(error) = mutate_member_unchecked(
remote_handle,
&remote_member_id,
PeerTarget::External(local_spec),
true,
)
.await
{
if let Ok(rollback_spec) = build_peer_spec(
&remote_info.comms_name,
&remote_info.peer_id,
&entry.transport,
Some(remote_info.pubkey),
) {
let _ = mutate_member_unchecked(
local_handle,
&local_member_id,
PeerTarget::External(rollback_spec),
false,
)
.await;
}
return Err(error);
}
Ok(())
}
LocalOrRemote::Remote(proxy) => {
let local_address =
member_comms_advertised_address(&local_runtime, &local_handle, &local_mid).await;
let local_address = match local_address {
Some(address) if !address.starts_with("inproc://") => address,
other => {
return Err(CrossMobError::LocalMemberNotRemotelyAddressable {
member_id: local_member_id.clone(),
mob_id: local_mob_id.clone(),
advertised: other,
});
}
};
let remote_info = proxy
.lookup_member(&remote_member_id)
.await
.map_err(CrossMobError::Remote)?;
let remote_spec =
build_remote_member_spec(&remote_info, &remote_member_id, &remote_mob_id)?;
mutate_member_unchecked(
local_handle.clone(),
&local_member_id,
PeerTarget::External(remote_spec),
true,
)
.await?;
if let Err(remote_error) = proxy
.wire_remote(
&remote_member_id,
&local_address,
&local_info.comms_name,
&local_info.peer_id,
Some(local_info.pubkey_b64.clone()),
)
.await
{
if let Ok(spec) =
build_remote_member_spec(&remote_info, &remote_member_id, &remote_mob_id)
{
let _ = mutate_member_unchecked(
local_handle,
&local_member_id,
PeerTarget::External(spec),
false,
)
.await;
}
return Err(CrossMobError::Remote(remote_error));
}
Ok(())
}
}
}
#[allow(clippy::too_many_arguments)]
async fn unwire_cross_mob_transaction(
entry: ContactEntry,
remote: LocalOrRemote,
local_runtime: crate::MobRuntime,
local_mid: AgentIdentity,
local_mob_id: String,
local_member_id: String,
remote_member_id: String,
remote_mob_id: String,
) -> Result<(), CrossMobError> {
let local_handle = local_runtime.handle();
let mut first_error = None;
let local_info = member_peer_info(&local_handle, &local_mid, &local_mob_id)
.await
.ok();
match remote {
LocalOrRemote::Local(remote_authority) => {
let remote_handle = remote_authority.handle;
let remote_mid = crate::member_comms_id::mob_member_id(&remote_member_id);
if let Ok(remote_info) =
member_peer_info(&remote_handle, &remote_mid, &remote_mob_id).await
&& let Ok(spec) = build_peer_spec(
&remote_info.comms_name,
&remote_info.peer_id,
&entry.transport,
Some(remote_info.pubkey),
)
&& let Err(error) = mutate_member_unchecked(
local_handle,
&local_member_id,
PeerTarget::External(spec),
false,
)
.await
{
first_error = Some(error);
}
if let Some(local_info) = &local_info
&& let Ok(spec) = build_peer_spec(
&local_info.comms_name,
&local_info.peer_id,
&MobTransport::Inproc,
Some(local_info.pubkey),
)
&& let Err(error) = mutate_member_unchecked(
remote_handle,
&remote_member_id,
PeerTarget::External(spec),
false,
)
.await
&& first_error.is_none()
{
first_error = Some(error);
}
}
LocalOrRemote::Remote(proxy) => {
if let Ok(remote_info) = proxy.lookup_member(&remote_member_id).await
&& let Ok(spec) =
build_remote_member_spec(&remote_info, &remote_member_id, &remote_mob_id)
&& let Err(error) = mutate_member_unchecked(
local_handle.clone(),
&local_member_id,
PeerTarget::External(spec),
false,
)
.await
{
first_error = Some(error);
}
if let Some(local_info) = &local_info {
match member_comms_advertised_address(&local_runtime, &local_handle, &local_mid)
.await
{
Some(local_address) if !local_address.starts_with("inproc://") => {
if let Err(error) = proxy
.unwire_remote(
&remote_member_id,
&local_address,
&local_info.comms_name,
&local_info.peer_id,
Some(local_info.pubkey_b64.clone()),
)
.await
&& first_error.is_none()
{
first_error = Some(CrossMobError::Remote(error));
}
}
other => {
if first_error.is_none() {
first_error = Some(CrossMobError::LocalMemberNotRemotelyAddressable {
member_id: local_member_id.clone(),
mob_id: local_mob_id.clone(),
advertised: other,
});
}
}
}
}
}
}
first_error.map_or(Ok(()), Err)
}
impl std::error::Error for CrossMobError {}
impl From<meerkat_mob::MobError> for CrossMobError {
fn from(err: meerkat_mob::MobError) -> Self {
Self::Mob(err)
}
}
impl From<RemoteMobError> for CrossMobError {
fn from(err: RemoteMobError) -> Self {
Self::Remote(err)
}
}
impl UnifiedRuntime {
pub(crate) async fn wire_bilateral_same_process(
&self,
peer: &UnifiedRuntime,
local_member_id: &str,
peer_member_id: &str,
) -> Result<(), CrossMobError> {
let local_member_id = self.topology_concrete_member_id(local_member_id).await?;
let peer_member_id = peer.topology_concrete_member_id(peer_member_id).await?;
let local_target = member_lifecycle_target_with_authority(
self.identity_runtime(),
&local_member_id,
&self.mob_id(),
)
.await?;
let peer_target = member_lifecycle_target_with_authority(
peer.identity_runtime(),
&peer_member_id,
&peer.mob_id(),
)
.await?;
let local_runtime = self.mob_runtime.clone();
let peer_runtime = peer.mob_runtime.clone();
run_member_authority_transaction(
[local_target, peer_target],
format!("{local_member_id} <-> {peer_member_id}"),
format!("{} <-> {}", self.mob_id(), peer.mob_id()),
move || async move {
wire_bilateral_transaction(
local_runtime,
peer_runtime,
local_member_id,
peer_member_id,
)
.await
},
)
.await
}
pub(crate) async fn unwire_bilateral_same_process(
&self,
peer: &UnifiedRuntime,
local_member_id: &str,
peer_member_id: &str,
) -> Result<(), CrossMobError> {
let local_member_id = self.topology_concrete_member_id(local_member_id).await?;
let peer_member_id = peer.topology_concrete_member_id(peer_member_id).await?;
let local_target = member_lifecycle_target_with_authority(
self.identity_runtime(),
&local_member_id,
&self.mob_id(),
)
.await?;
let peer_target = member_lifecycle_target_with_authority(
peer.identity_runtime(),
&peer_member_id,
&peer.mob_id(),
)
.await?;
let local_runtime = self.mob_runtime.clone();
let peer_runtime = peer.mob_runtime.clone();
run_member_authority_transaction(
[local_target, peer_target],
format!("{local_member_id} <-> {peer_member_id}"),
format!("{} <-> {}", self.mob_id(), peer.mob_id()),
move || async move {
unwire_bilateral_transaction(
local_runtime,
peer_runtime,
local_member_id,
peer_member_id,
)
.await
},
)
.await
}
pub(crate) async fn bilateral_same_process_state(
&self,
peer: &UnifiedRuntime,
local_member_id: &str,
peer_member_id: &str,
) -> Result<(bool, bool), CrossMobError> {
let local_member_id = self.topology_concrete_member_id(local_member_id).await?;
let peer_member_id = peer.topology_concrete_member_id(peer_member_id).await?;
let local_handle = self.mob_runtime.handle();
let peer_handle = peer.mob_runtime.handle();
let local_mid = crate::member_comms_id::mob_member_id(&local_member_id);
let peer_mid = crate::member_comms_id::mob_member_id(&peer_member_id);
let local_info = self
.get_member_peer_info(&local_handle, &local_mid, &self.mob_id())
.await?;
let peer_info = self
.get_member_peer_info(&peer_handle, &peer_mid, &peer.mob_id())
.await?;
let local = local_handle.get_member(&local_mid).await?.ok_or_else(|| {
CrossMobError::MemberNotFound {
member_id: local_member_id.clone(),
mob_id: self.mob_id(),
}
})?;
let remote = peer_handle.get_member(&peer_mid).await?.ok_or_else(|| {
CrossMobError::MemberNotFound {
member_id: peer_member_id.clone(),
mob_id: peer.mob_id(),
}
})?;
let local_structural = local
.wired_to
.iter()
.any(|identity| identity.as_str() == peer_info.comms_name);
let remote_structural = remote
.wired_to
.iter()
.any(|identity| identity.as_str() == local_info.comms_name);
let local_namespace = mob_inproc_namespace(self)?;
let peer_namespace = mob_inproc_namespace(peer)?;
let local_route = alias_already_installed(
&local_namespace,
&peer_info.comms_name,
PubKey::new(peer_info.pubkey),
);
let remote_route = alias_already_installed(
&peer_namespace,
&local_info.comms_name,
PubKey::new(local_info.pubkey),
);
let local_agent = self
.agent_can_address_peer(&local_handle, &local_mid, &peer_info)
.await?;
let remote_agent = peer
.agent_can_address_peer(&peer_handle, &peer_mid, &local_info)
.await?;
Ok((
local_structural && local_route && local_agent,
remote_structural && remote_route && remote_agent,
))
}
async fn topology_concrete_member_id(
&self,
logical_identity: &str,
) -> Result<String, CrossMobError> {
let Some(context) = self.identity_first_context.as_ref() else {
return Ok(logical_identity.to_string());
};
let identity = crate::identity_first::AgentIdentity::parse(logical_identity)
.map_err(|error| CrossMobError::PeerSpec(error.to_string()))?;
context
.runtime
.runtime_id_for(&identity)
.await
.map(|runtime_id| runtime_id.to_string())
.map_err(|error| CrossMobError::MemberNotFound {
member_id: format!("{logical_identity}: {error}"),
mob_id: self.mob_id(),
})
}
pub async fn register_peer_mob(&self, mob_id: &str, handle: MobHandle) {
self.peer_mob_handles.write().await.insert(
mob_id.to_string(),
PeerMobAuthority {
handle,
identity_runtime: None,
},
);
}
pub async fn register_peer_mob_with_identity_runtime(
&self,
mob_id: &str,
handle: MobHandle,
identity_runtime: std::sync::Arc<crate::identity_first::IdentityRuntime>,
) {
self.peer_mob_handles.write().await.insert(
mob_id.to_string(),
PeerMobAuthority {
handle,
identity_runtime: Some(identity_runtime),
},
);
}
pub async fn register_peer_runtime(&self, peer: &UnifiedRuntime) {
self.peer_mob_handles.write().await.insert(
peer.mob_id(),
PeerMobAuthority {
handle: peer.mob_handle(),
identity_runtime: peer.identity_runtime().cloned(),
},
);
}
pub fn set_contact_directory(&mut self, directory: ContactDirectory) {
self.contact_directory = Some(directory);
}
pub async fn start_control_listener(
&self,
addr: &crate::runtime::cross_mob_control::ControlListenAddr,
) -> Result<String, CrossMobError> {
tracing::warn!(%addr, "cross-mob control listener has no caller grants; refusing every request");
self.start_control_listener_with_authorizer(
addr,
std::sync::Arc::new(
crate::runtime::cross_mob_control::ControlAuthorizer::with_grants_for_audience(
crate::runtime::cross_mob_control::ControlGrantTable::new(),
self.mob_id(),
),
),
)
.await
}
pub async fn start_control_listener_with_authorizer(
&self,
addr: &crate::runtime::cross_mob_control::ControlListenAddr,
authorizer: std::sync::Arc<crate::runtime::cross_mob_control::ControlAuthorizer>,
) -> Result<String, CrossMobError> {
let mut task_slot = self.cross_mob_control_task.lock().await;
if task_slot.is_some() {
return Err(CrossMobError::ControlListener(
"control listener already running".to_string(),
));
}
let bound = crate::runtime::cross_mob_control::BoundControlListener::bind(addr)
.await
.map_err(|error| CrossMobError::ControlListener(format!("bind {addr}: {error}")))?;
let advertised = bound.advertised_address().to_string();
let mut handler =
crate::runtime::cross_mob_control::MobHandleControlHandler::with_shared_identity_authority(
self.mob_handle(),
std::sync::Arc::clone(&self.implicit_delegate_identity_runtime),
);
if let Some(service) = self.mob_runtime.session_service() {
handler = handler.with_session_service(std::sync::Arc::clone(service));
}
self.install_remote_host_facts(&advertised);
handler = handler.with_host_facts_slot(std::sync::Arc::clone(&self.remote_host_facts));
let handler: std::sync::Arc<dyn crate::runtime::cross_mob_control::ControlHandler> =
std::sync::Arc::new(handler);
*task_slot = Some(tokio::spawn(bound.serve_with_authorizer(
handler,
std::sync::Arc::clone(&self.gateway_peer_keys),
authorizer,
)));
*self
.cross_mob_control_advertised
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(advertised.clone());
Ok(advertised)
}
pub fn control_listener_advertised_address(&self) -> Option<String> {
self.cross_mob_control_advertised
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
fn install_remote_host_facts(&self, advertised: &str) {
let Some(keys) = self.gateway_peer_keys() else {
return;
};
let mut options = meerkat::surface::RuntimeHostSurfaceOptions::process(
"meerkat-mobkit",
env!("CARGO_PKG_VERSION"),
);
options.runtime_backed_sessions = self.mob_runtime.session_service().is_some();
options.mobs = true;
options.multi_host_mobs = true;
options.comms = true;
options.session_events = true;
options.session_streams = true;
options.rpc_transport = Some("mobkit_control".to_string());
options.rpc_methods = vec!["host/describe".to_string(), "host/health".to_string()];
let capabilities = meerkat::surface::build_runtime_host_capabilities(&options);
let facts = crate::runtime::remote_host::HostFacts::new(
self.mob_id(),
keys.pubkey_b64(),
advertised,
capabilities,
)
.with_endpoints(meerkat_contracts::RuntimeHostEndpointProjection {
rpc_transport: Some("mobkit_control".to_string()),
rest_base_url: None,
rpc_methods: vec!["host/describe".to_string(), "host/health".to_string()],
rest_paths: Vec::new(),
});
let health = meerkat_contracts::RuntimeHostHealth {
contract_version: meerkat_contracts::ContractVersion::CURRENT,
status: meerkat_contracts::RuntimeHostHealthStatus::Ok,
checks: std::collections::BTreeMap::new(),
};
*self
.remote_host_facts
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(std::sync::Arc::new(
crate::runtime::remote_host::StaticHostFacts::new(facts, health),
));
}
pub fn set_gateway_peer_keys(&mut self, keys: GatewayPeerKeys) {
if self.gateway_peer_keys().is_some() {
tracing::warn!("gateway peer keys are already installed; refusing live replacement");
return;
}
let state_root = keys.state_root().map(std::path::Path::to_path_buf);
*self
.gateway_peer_keys
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(std::sync::Arc::new(keys));
if let Some(advertised) = self.control_listener_advertised_address() {
self.install_remote_host_facts(&advertised);
}
let Some(state_root) = state_root else {
return;
};
let path = state_root.join(crate::runtime::remote_host::HOST_PAIRING_FILE_NAME);
match crate::runtime::remote_host::RemoteHostLifecycle::load(
path,
crate::runtime::remote_host::HostReconnectPolicy::default(),
) {
Ok(lifecycle) => {
let lifecycle = std::sync::Arc::new(lifecycle);
*self
.remote_host_lifecycle
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) =
Some(std::sync::Arc::clone(&lifecycle));
*self
.remote_host_lifecycle_error
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
self.start_remote_host_reconnect_task(lifecycle);
}
Err(error) => {
*self
.remote_host_lifecycle
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
*self
.remote_host_lifecycle_error
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(error);
}
}
}
fn start_remote_host_reconnect_task(
&mut self,
lifecycle: std::sync::Arc<crate::runtime::remote_host::RemoteHostLifecycle>,
) {
if self.remote_host_reconnect_task.get_mut().is_some() {
return;
}
let Ok(handle) = tokio::runtime::Handle::try_current() else {
return;
};
let caller_keys = self.gateway_peer_keys();
let task = handle.spawn(async move {
loop {
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
for outcome in lifecycle
.probe_due(
caller_keys.clone(),
crate::runtime::remote_host::unix_now_secs(),
)
.await
{
if let Err(error) = outcome.result {
tracing::debug!(
host_key_b64 = %outcome.host_key_b64,
%error,
"runtime-host reconnect probe failed"
);
}
}
}
});
*self.remote_host_reconnect_task.get_mut() = Some(task);
}
fn remote_host_lifecycle(
&self,
) -> Result<
std::sync::Arc<crate::runtime::remote_host::RemoteHostLifecycle>,
crate::runtime::remote_host::RemoteHostLifecycleError,
> {
if let Some(error) = self
.remote_host_lifecycle_error
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
{
return Err(error.into());
}
self.remote_host_lifecycle
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
.ok_or_else(|| {
crate::runtime::remote_host::RemoteHostLifecycleError::Pairing(
crate::runtime::remote_host::HostPairingError::Io {
path: crate::runtime::remote_host::HOST_PAIRING_FILE_NAME.to_string(),
reason: "gateway has no durable state root".to_string(),
},
)
})
}
pub async fn pair_runtime_host(
&self,
contact_id: &str,
) -> Result<
crate::runtime::remote_host::RuntimeHostPlacement,
crate::runtime::remote_host::RemoteHostLifecycleError,
> {
let contact = self
.contact_directory
.as_ref()
.and_then(|directory| directory.get(contact_id))
.cloned()
.ok_or_else(|| {
crate::runtime::remote_host::RemoteHostLifecycleError::UnknownContact {
contact: contact_id.to_string(),
}
})?;
self.remote_host_lifecycle()?
.pair_contact(
&contact,
self.gateway_peer_keys(),
crate::runtime::remote_host::unix_now_secs(),
)
.await
}
pub async fn runtime_host_placement(
&self,
host_key_b64: &str,
) -> Result<
crate::runtime::remote_host::RuntimeHostPlacement,
crate::runtime::remote_host::RemoteHostLifecycleError,
> {
self.remote_host_lifecycle()?.placement(host_key_b64).await
}
pub fn gateway_peer_keys(&self) -> Option<std::sync::Arc<GatewayPeerKeys>> {
self.gateway_peer_keys
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub async fn wire_cross_mob(
&self,
local_member_id: &str,
remote_member_id: &str,
remote_mob_id: &str,
) -> Result<(), CrossMobError> {
self.wire_cross_mob_with_identity_runtime(
local_member_id,
remote_member_id,
remote_mob_id,
self.identity_runtime(),
)
.await
}
pub(crate) async fn wire_cross_mob_with_identity_runtime(
&self,
local_member_id: &str,
remote_member_id: &str,
remote_mob_id: &str,
local_identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
) -> Result<(), CrossMobError> {
let local_member_id =
crate::member_comms_id::runtime_alias_str(local_member_id).into_owned();
let remote_member_id =
crate::member_comms_id::runtime_alias_str(remote_member_id).into_owned();
let local_handle = self.mob_runtime.handle();
let local_mob_id = local_handle.mob_id().to_string();
let local_identity_runtime = local_identity_runtime.cloned();
let local_target = member_lifecycle_target_with_authority(
local_identity_runtime.as_ref(),
&local_member_id,
&local_mob_id,
)
.await?;
let local_identity_authoritative = local_target.is_some();
let entry = self.resolve_contact(remote_mob_id)?;
let remote = self.dispatch_for(&entry).await?;
let remote_target = match &remote {
LocalOrRemote::Local(authority) => {
member_lifecycle_target_with_authority(
authority.identity_runtime.as_ref(),
&remote_member_id,
remote_mob_id,
)
.await?
}
LocalOrRemote::Remote(_) => None,
};
let remote_identity_authoritative = remote_target.is_some();
let remote_mob_id = remote_mob_id.to_string();
let local_runtime = self.mob_runtime.clone();
run_member_authority_transaction(
[local_target, remote_target],
format!("{local_member_id} <-> {remote_member_id}"),
format!("{local_mob_id} <-> {remote_mob_id}"),
move || async move {
let local_member_id = resolve_member_alias_under_authority(
&local_handle,
local_identity_runtime.as_ref(),
local_identity_authoritative,
&local_member_id,
&local_mob_id,
)
.await?;
let local_mid = crate::member_comms_id::mob_member_id(&local_member_id);
let remote_member_id = match &remote {
LocalOrRemote::Local(authority) => {
resolve_member_alias_under_authority(
&authority.handle,
authority.identity_runtime.as_ref(),
remote_identity_authoritative,
&remote_member_id,
&remote_mob_id,
)
.await?
}
LocalOrRemote::Remote(_) => remote_member_id,
};
wire_cross_mob_transaction(
entry,
remote,
local_runtime,
local_mid,
local_mob_id,
local_member_id,
remote_member_id,
remote_mob_id,
)
.await
},
)
.await
}
pub async fn unwire_cross_mob(
&self,
local_member_id: &str,
remote_member_id: &str,
remote_mob_id: &str,
) -> Result<(), CrossMobError> {
self.unwire_cross_mob_with_identity_runtime(
local_member_id,
remote_member_id,
remote_mob_id,
self.identity_runtime(),
)
.await
}
pub(crate) async fn unwire_cross_mob_with_identity_runtime(
&self,
local_member_id: &str,
remote_member_id: &str,
remote_mob_id: &str,
local_identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
) -> Result<(), CrossMobError> {
let local_member_id =
crate::member_comms_id::runtime_alias_str(local_member_id).into_owned();
let remote_member_id =
crate::member_comms_id::runtime_alias_str(remote_member_id).into_owned();
let local_handle = self.mob_runtime.handle();
let local_mob_id = local_handle.mob_id().to_string();
let local_identity_runtime = local_identity_runtime.cloned();
let local_target = member_lifecycle_target_with_authority(
local_identity_runtime.as_ref(),
&local_member_id,
&local_mob_id,
)
.await?;
let local_identity_authoritative = local_target.is_some();
let entry = self.resolve_contact(remote_mob_id)?;
let remote = self.dispatch_for(&entry).await?;
let remote_target = match &remote {
LocalOrRemote::Local(authority) => {
member_lifecycle_target_with_authority(
authority.identity_runtime.as_ref(),
&remote_member_id,
remote_mob_id,
)
.await?
}
LocalOrRemote::Remote(_) => None,
};
let remote_identity_authoritative = remote_target.is_some();
let remote_mob_id = remote_mob_id.to_string();
let local_runtime = self.mob_runtime.clone();
run_member_authority_transaction(
[local_target, remote_target],
format!("{local_member_id} <-> {remote_member_id}"),
format!("{local_mob_id} <-> {remote_mob_id}"),
move || async move {
let local_member_id = resolve_member_alias_under_authority(
&local_handle,
local_identity_runtime.as_ref(),
local_identity_authoritative,
&local_member_id,
&local_mob_id,
)
.await?;
let local_mid = crate::member_comms_id::mob_member_id(&local_member_id);
let remote_member_id = match &remote {
LocalOrRemote::Local(authority) => {
resolve_member_alias_under_authority(
&authority.handle,
authority.identity_runtime.as_ref(),
remote_identity_authoritative,
&remote_member_id,
&remote_mob_id,
)
.await?
}
LocalOrRemote::Remote(_) => remote_member_id,
};
unwire_cross_mob_transaction(
entry,
remote,
local_runtime,
local_mid,
local_mob_id,
local_member_id,
remote_member_id,
remote_mob_id,
)
.await
},
)
.await
}
pub async fn send_cross_mob(
&self,
from_local_member: &str,
remote_member_id: &str,
remote_mob_id: &str,
content: impl Into<meerkat_core::ContentInput>,
) -> Result<String, CrossMobError> {
self.send_cross_mob_with_identity_runtime(
from_local_member,
remote_member_id,
remote_mob_id,
content,
self.identity_runtime(),
)
.await
}
pub(crate) async fn send_cross_mob_with_identity_runtime(
&self,
from_local_member: &str,
remote_member_id: &str,
remote_mob_id: &str,
content: impl Into<meerkat_core::ContentInput>,
local_identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
) -> Result<String, CrossMobError> {
let content = content.into();
let from_local_member =
crate::member_comms_id::runtime_alias_str(from_local_member).into_owned();
let remote_member_id =
crate::member_comms_id::runtime_alias_str(remote_member_id).into_owned();
let local_mob_id = self.mob_id();
let local_identity_runtime = local_identity_runtime.cloned();
let local_target = member_lifecycle_target_with_authority(
local_identity_runtime.as_ref(),
&from_local_member,
&local_mob_id,
)
.await?;
let entry = self.resolve_contact(remote_mob_id)?;
let remote = self.dispatch_for(&entry).await?;
let remote_target = match &remote {
LocalOrRemote::Local(authority) => {
member_lifecycle_target_with_authority(
authority.identity_runtime.as_ref(),
&remote_member_id,
remote_mob_id,
)
.await?
}
LocalOrRemote::Remote(_) => None,
};
let remote_identity_authoritative = remote_target.is_some();
let remote_mob_id = remote_mob_id.to_string();
run_member_authority_transaction(
[local_target, remote_target],
format!("{from_local_member} -> {remote_member_id}"),
format!("{local_mob_id} -> {remote_mob_id}"),
move || async move {
match remote {
LocalOrRemote::Local(authority) => {
let remote_member_id = resolve_member_alias_under_authority(
&authority.handle,
authority.identity_runtime.as_ref(),
remote_identity_authoritative,
&remote_member_id,
&remote_mob_id,
)
.await?;
send_member_unchecked(
authority.handle,
&remote_member_id,
&remote_mob_id,
content,
)
.await
}
LocalOrRemote::Remote(proxy) => {
let content_json = serde_json::to_value(&content).map_err(|err| {
CrossMobError::PeerSpec(format!(
"failed to serialize content for remote inject: {err}"
))
})?;
proxy
.inject_message(&remote_member_id, content_json)
.await
.map_err(CrossMobError::Remote)
}
}
},
)
.await
}
pub fn list_external_mobs(&self) -> Vec<ContactEntry> {
self.contact_directory
.as_ref()
.map(|d| d.list().into_iter().cloned().collect())
.unwrap_or_default()
}
pub fn has_contact_directory(&self) -> bool {
self.contact_directory.is_some()
}
pub async fn has_peer_mob_handles(&self) -> bool {
!self.peer_mob_handles.read().await.is_empty()
}
pub fn has_inproc_contacts(&self) -> bool {
self.contact_directory.as_ref().is_some_and(|d| {
d.list()
.iter()
.any(|e| matches!(e.transport, MobTransport::Inproc))
})
}
pub fn has_remote_contacts(&self) -> bool {
self.contact_directory.as_ref().is_some_and(|d| {
d.list()
.iter()
.any(|e| matches!(e.transport, MobTransport::Tcp(_) | MobTransport::Uds(_)))
})
}
pub fn mob_id(&self) -> String {
self.mob_runtime.handle().mob_id().to_string()
}
pub async fn local_member_peer_info(
&self,
member_id: &str,
) -> Result<(String, String, String), CrossMobError> {
let handle = self.mob_runtime.handle();
let mob_id = handle.mob_id().to_string();
let member_alias = crate::member_comms_id::runtime_alias_str(member_id).into_owned();
let identity_runtime = self.identity_runtime().cloned();
let target = member_lifecycle_target_with_authority(
identity_runtime.as_ref(),
&member_alias,
&mob_id,
)
.await?;
let identity_authoritative = target.is_some();
let operation_alias = member_alias.clone();
let operation_mob_id = mob_id.clone();
run_member_authority_transaction([target], member_alias, mob_id, move || async move {
let current_alias = resolve_member_alias_under_authority(
&handle,
identity_runtime.as_ref(),
identity_authoritative,
&operation_alias,
&operation_mob_id,
)
.await?;
let mid = crate::member_comms_id::mob_member_id(¤t_alias);
let info = member_peer_info(&handle, &mid, &operation_mob_id).await?;
let address = format!("inproc://{}", info.comms_name);
Ok((info.peer_id, info.comms_name, address))
})
.await
}
pub async fn wire_local(
&self,
local_member_id: &str,
remote_comms_name: &str,
remote_peer_id: &str,
remote_address: &str,
remote_pubkey: Option<[u8; 32]>,
) -> Result<(), CrossMobError> {
self.wire_local_with_identity_runtime(
local_member_id,
remote_comms_name,
remote_peer_id,
remote_address,
remote_pubkey,
None,
)
.await
}
pub(crate) async fn wire_local_with_identity_runtime(
&self,
local_member_id: &str,
remote_comms_name: &str,
remote_peer_id: &str,
remote_address: &str,
remote_pubkey: Option<[u8; 32]>,
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
) -> Result<(), CrossMobError> {
let spec = build_external_peer_spec(
remote_comms_name,
remote_peer_id,
remote_address,
remote_pubkey,
)?;
wire_member_with_authority(
self.mob_runtime.handle(),
identity_runtime.or_else(|| self.identity_runtime()),
local_member_id,
&self.mob_id(),
PeerTarget::External(spec),
true,
)
.await
}
pub async fn unwire_local(
&self,
local_member_id: &str,
remote_comms_name: &str,
remote_peer_id: &str,
remote_address: &str,
remote_pubkey: Option<[u8; 32]>,
) -> Result<(), CrossMobError> {
self.unwire_local_with_identity_runtime(
local_member_id,
remote_comms_name,
remote_peer_id,
remote_address,
remote_pubkey,
None,
)
.await
}
pub(crate) async fn unwire_local_with_identity_runtime(
&self,
local_member_id: &str,
remote_comms_name: &str,
remote_peer_id: &str,
remote_address: &str,
remote_pubkey: Option<[u8; 32]>,
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
) -> Result<(), CrossMobError> {
let spec = build_external_peer_spec(
remote_comms_name,
remote_peer_id,
remote_address,
remote_pubkey,
)?;
wire_member_with_authority(
self.mob_runtime.handle(),
identity_runtime.or_else(|| self.identity_runtime()),
local_member_id,
&self.mob_id(),
PeerTarget::External(spec),
false,
)
.await
}
fn resolve_contact(&self, mob_id: &str) -> Result<ContactEntry, CrossMobError> {
let dir = self
.contact_directory
.as_ref()
.ok_or(CrossMobError::NoContactDirectory)?;
dir.get(mob_id)
.cloned()
.ok_or_else(|| CrossMobError::UnknownMob(mob_id.to_string()))
}
async fn dispatch_for(&self, entry: &ContactEntry) -> Result<LocalOrRemote, CrossMobError> {
if let Some(authority) = self
.peer_mob_handles
.read()
.await
.get(&entry.mob_id)
.cloned()
{
return Ok(LocalOrRemote::Local(Box::new(authority)));
}
match RemoteMobProxy::from_entry_with_caller(entry, self.gateway_peer_keys())? {
Some(proxy) => Ok(LocalOrRemote::Remote(proxy)),
None => Err(CrossMobError::NoPeerHandle(entry.mob_id.clone())),
}
}
async fn get_member_peer_info(
&self,
handle: &MobHandle,
meerkat_id: &AgentIdentity,
mob_id: &str,
) -> Result<MemberPeerInfo, CrossMobError> {
let entry = handle
.get_member(meerkat_id)
.await
.map_err(|err| {
CrossMobError::PeerSpec(format!(
"member lookup for '{meerkat_id}' in mob '{mob_id}' failed: {err}"
))
})?
.ok_or_else(|| CrossMobError::MemberNotFound {
member_id: meerkat_id.to_string(),
mob_id: mob_id.to_string(),
})?;
let peer_id = entry
.peer_id()
.ok_or_else(|| CrossMobError::NoCommsInfo {
member_id: meerkat_id.to_string(),
mob_id: mob_id.to_string(),
})?
.to_string();
let pubkey_b64 = entry.transport_public_key().ok_or_else(|| {
CrossMobError::PeerSpec(format!(
"member '{meerkat_id}' in mob '{mob_id}' has no transport public key"
))
})?;
let pubkey = crate::auth::peer_keys::decode_pubkey_b64(pubkey_b64).map_err(|err| {
CrossMobError::PeerSpec(format!(
"member '{meerkat_id}' in mob '{mob_id}' has invalid transport public key: {err}"
))
})?;
let comms_name = meerkat_core::MemberCommsName::new(
mob_id,
entry.role.as_str(),
meerkat_id.as_str(),
)
.map_err(|err| {
CrossMobError::PeerSpec(format!(
"member '{meerkat_id}' in mob '{mob_id}' has an invalid comms name component: {err}"
))
})?
.to_string();
Ok(MemberPeerInfo {
peer_id,
comms_name,
pubkey,
pubkey_b64: pubkey_b64.to_string(),
})
}
async fn agent_can_address_peer(
&self,
handle: &MobHandle,
local_member: &AgentIdentity,
expected_peer: &MemberPeerInfo,
) -> Result<bool, CrossMobError> {
let session_id = handle
.resolve_bridge_session_id(local_member)
.await
.ok_or_else(|| CrossMobError::NoCommsInfo {
member_id: local_member.to_string(),
mob_id: handle.mob_id().to_string(),
})?;
let service =
self.mob_runtime
.session_service()
.ok_or_else(|| CrossMobError::NoCommsInfo {
member_id: local_member.to_string(),
mob_id: handle.mob_id().to_string(),
})?;
let comms =
service
.comms_runtime(&session_id)
.await
.ok_or_else(|| CrossMobError::NoCommsInfo {
member_id: local_member.to_string(),
mob_id: handle.mob_id().to_string(),
})?;
let peers = comms.peers().await;
Ok(peers.iter().any(|peer| {
peer.peer_id.to_string() == expected_peer.peer_id
&& peer.name.as_str() == expected_peer.comms_name
&& peer
.sendable_kinds
.contains(&meerkat_core::comms::PeerSendability::PeerMessage)
}))
}
}
fn build_peer_spec(
comms_name: &str,
peer_id: &str,
transport: &MobTransport,
pubkey: Option<[u8; 32]>,
) -> Result<TrustedPeerDescriptor, CrossMobError> {
let address = match transport {
MobTransport::Inproc => format!("inproc://{comms_name}"),
MobTransport::Tcp(addr) => format!("tcp://{addr}"),
MobTransport::Uds(path) => format!("uds://{path}"),
};
build_external_peer_spec(comms_name, peer_id, &address, pubkey)
}
fn build_external_peer_spec(
comms_name: &str,
peer_id: &str,
address: &str,
pubkey: Option<[u8; 32]>,
) -> Result<TrustedPeerDescriptor, CrossMobError> {
let is_inproc = address.starts_with("inproc://");
match (is_inproc, pubkey) {
(true, None) => TrustedPeerDescriptor::test_only_unsigned(comms_name, peer_id, address)
.map_err(CrossMobError::PeerSpec),
(true, Some(bytes)) => {
TrustedPeerDescriptor::unsigned_with_pubkey(comms_name, peer_id, bytes, address)
.map_err(CrossMobError::PeerSpec)
}
(false, None) => Err(CrossMobError::MissingPeerPubkey { mob_id: None }),
(false, Some(bytes)) => {
if bytes == [0u8; 32] {
return Err(CrossMobError::MissingPeerPubkey { mob_id: None });
}
TrustedPeerDescriptor::unsigned_with_pubkey(comms_name, peer_id, bytes, address)
.map_err(CrossMobError::PeerSpec)
}
}
}
pub fn build_tcp_peer_spec(
comms_name: &str,
peer_id: &str,
address: &str,
) -> Result<TrustedPeerDescriptor, CrossMobError> {
TrustedPeerDescriptor::test_only_unsigned(comms_name, peer_id, format!("tcp://{address}"))
.map_err(CrossMobError::PeerSpec)
}
pub fn build_uds_peer_spec(
comms_name: &str,
peer_id: &str,
path: &str,
) -> Result<TrustedPeerDescriptor, CrossMobError> {
let normalized = if let Some(stripped) = path.strip_prefix('/') {
stripped
} else {
path
};
TrustedPeerDescriptor::test_only_unsigned(comms_name, peer_id, format!("uds:///{normalized}"))
.map_err(CrossMobError::PeerSpec)
}
fn mob_inproc_namespace(runtime: &UnifiedRuntime) -> Result<String, CrossMobError> {
mob_inproc_namespace_for_id(&runtime.mob_id())
}
fn mob_inproc_namespace_for_id(mob_id: &str) -> Result<String, CrossMobError> {
meerkat_core::mob_realm_id(mob_id)
.map(|realm| realm.as_str().to_string())
.map_err(|error| CrossMobError::InprocAlias(error.to_string()))
}
fn alias_already_installed(namespace: &str, name: &str, pubkey: PubKey) -> bool {
InprocRegistry::global()
.peers_in_namespace(namespace)
.into_iter()
.any(|peer| peer.name == name && peer.pubkey == pubkey)
}
fn ensure_alias_slot(
namespace: &str,
canonical_namespace: &str,
name: &str,
pubkey: PubKey,
) -> Result<(), String> {
let registry = InprocRegistry::global();
for peer in InprocRegistry::global().peers_in_namespace(namespace) {
if peer.name == name && peer.pubkey != pubkey {
let old_route_is_live = registry
.peers_in_namespace(canonical_namespace)
.into_iter()
.any(|canonical| canonical.name == name && canonical.pubkey == peer.pubkey);
if old_route_is_live {
return Err(format!(
"namespace {namespace:?} already binds name {name:?} to another live peer"
));
}
registry.unregister_in_namespace(namespace, &peer.pubkey);
continue;
}
if peer.pubkey == pubkey && peer.name != name {
let canonical_name_matches = registry
.peers_in_namespace(canonical_namespace)
.into_iter()
.any(|canonical| canonical.name == name && canonical.pubkey == pubkey);
if !canonical_name_matches {
return Err(format!(
"namespace {namespace:?} already binds peer {pubkey:?} as {:?}",
peer.name
));
}
registry.unregister_in_namespace(namespace, &peer.pubkey);
}
}
Ok(())
}
fn register_cross_namespace_aliases(
source_namespace: &str,
source_comms_name: &str,
source_pubkey: [u8; 32],
target_namespace: &str,
target_comms_name: &str,
target_pubkey: [u8; 32],
) -> Result<(), String> {
let registry = InprocRegistry::global();
let source_pubkey = PubKey::new(source_pubkey);
let target_pubkey = PubKey::new(target_pubkey);
ensure_alias_slot(
source_namespace,
target_namespace,
target_comms_name,
target_pubkey,
)?;
ensure_alias_slot(
target_namespace,
source_namespace,
source_comms_name,
source_pubkey,
)?;
let target_sender = registry
.get_by_pubkey_in_namespace(target_namespace, &target_pubkey)
.ok_or_else(|| {
format!("target peer {target_comms_name:?} is not registered in {target_namespace:?}")
})?;
let source_sender = registry
.get_by_pubkey_in_namespace(source_namespace, &source_pubkey)
.ok_or_else(|| {
format!("source peer {source_comms_name:?} is not registered in {source_namespace:?}")
})?;
let target_preexisting =
alias_already_installed(source_namespace, target_comms_name, target_pubkey);
let source_preexisting =
alias_already_installed(target_namespace, source_comms_name, source_pubkey);
if !target_preexisting {
let outcome = registry.register_with_meta_in_namespace(
source_namespace,
target_comms_name,
target_pubkey,
target_sender,
PeerMeta::default(),
);
if outcome.is_rejected() || outcome.displaced_existing() {
return Err(format!("target alias installation failed: {outcome:?}"));
}
}
if !source_preexisting {
let outcome = registry.register_with_meta_in_namespace(
target_namespace,
source_comms_name,
source_pubkey,
source_sender,
PeerMeta::default(),
);
if outcome.is_rejected() || outcome.displaced_existing() {
if !target_preexisting {
registry.unregister_in_namespace(source_namespace, &target_pubkey);
}
return Err(format!("source alias installation failed: {outcome:?}"));
}
}
Ok(())
}
async fn handle_references_peer(handle: &MobHandle, peer_id: &str) -> bool {
handle
.list_members_including_retiring()
.await
.iter()
.any(|member| member.wired_to.iter().any(|peer| peer.as_str() == peer_id))
}
fn unregister_cross_namespace_aliases(
source_namespace: &str,
source_pubkey: [u8; 32],
target_pubkey: [u8; 32],
source_still_references_target: bool,
target_namespace: &str,
target_still_references_source: bool,
) {
let registry = InprocRegistry::global();
if !source_still_references_target {
registry.unregister_in_namespace(source_namespace, &PubKey::new(target_pubkey));
}
if !target_still_references_source {
registry.unregister_in_namespace(target_namespace, &PubKey::new(source_pubkey));
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
const TEST_PEER_ID: &str = "00000000-0000-4000-8000-000000000001";
const TEST_PUBKEY: [u8; 32] = [42u8; 32];
fn derived_peer_id() -> String {
meerkat_core::comms::PeerId::from_ed25519_pubkey(&TEST_PUBKEY).to_string()
}
#[test]
fn peer_spec_inproc_uses_comms_name_address() {
let spec = build_peer_spec(
"authors/coordinator/alice",
TEST_PEER_ID,
&MobTransport::Inproc,
None,
)
.expect("spec");
assert_eq!(spec.address.endpoint(), "authors/coordinator/alice");
}
#[test]
fn peer_spec_tcp_uses_tcp_scheme() {
let id = derived_peer_id();
let spec = build_peer_spec(
"authors/coordinator/alice",
&id,
&MobTransport::Tcp("127.0.0.1:9001".to_string()),
Some(TEST_PUBKEY),
)
.expect("spec");
assert_eq!(spec.address.endpoint(), "127.0.0.1:9001");
}
#[test]
fn peer_spec_uds_uses_uds_scheme() {
let id = derived_peer_id();
let spec = build_peer_spec(
"authors/coordinator/alice",
&id,
&MobTransport::Uds("/tmp/x.sock".to_string()),
Some(TEST_PUBKEY),
)
.expect("spec");
assert_eq!(spec.address.endpoint(), "/tmp/x.sock");
}
#[test]
fn peer_spec_tcp_without_pubkey_rejected() {
let result = build_peer_spec(
"authors/coordinator/alice",
TEST_PEER_ID,
&MobTransport::Tcp("127.0.0.1:9001".to_string()),
None,
);
assert!(
matches!(result, Err(CrossMobError::MissingPeerPubkey { .. })),
"TCP peer spec without pubkey must fail closed, got {result:?}"
);
}
#[test]
fn peer_spec_uds_without_pubkey_rejected() {
let result = build_peer_spec(
"authors/coordinator/alice",
TEST_PEER_ID,
&MobTransport::Uds("/tmp/x.sock".to_string()),
None,
);
assert!(
matches!(result, Err(CrossMobError::MissingPeerPubkey { .. })),
"UDS peer spec without pubkey must fail closed, got {result:?}"
);
}
#[test]
fn build_uds_peer_spec_handles_leading_slash() {
let with = build_uds_peer_spec("a", "00000000-0000-4000-8000-000000000001", "/tmp/x.sock")
.expect("spec");
let without =
build_uds_peer_spec("a", "00000000-0000-4000-8000-000000000001", "tmp/x.sock")
.expect("spec");
assert_eq!(with.address.endpoint(), "/tmp/x.sock");
assert_eq!(without.address.endpoint(), "/tmp/x.sock");
}
}