use std::collections::HashMap;
use std::sync::OnceLock;
use parking_lot::RwLock;
use uuid::Uuid;
use crate::identity::{Keypair, PubKey, Signature};
use crate::inbox::{AdmissionOutcome, DropReason, InboxSender};
use crate::peer_meta::PeerMeta;
use crate::types::{Envelope, InboxItem, MessageKind};
const DEFAULT_NAMESPACE: &str = "";
#[derive(Debug, Clone)]
pub struct InprocPeerInfo {
pub name: String,
pub pubkey: PubKey,
pub meta: PeerMeta,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RegistrationRejection {
ZeroPubkey,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RegistrationOutcome {
Registered,
ReplacedPubkey { evicted_name: String },
EvictedName { evicted_pubkey: PubKey },
ReplacedPubkeyAndEvictedName {
evicted_name: String,
evicted_pubkey: PubKey,
},
Rejected { reason: RegistrationRejection },
}
impl RegistrationOutcome {
pub fn displaced_existing(&self) -> bool {
matches!(
self,
Self::ReplacedPubkey { .. }
| Self::EvictedName { .. }
| Self::ReplacedPubkeyAndEvictedName { .. }
)
}
pub fn is_rejected(&self) -> bool {
matches!(self, Self::Rejected { .. })
}
}
#[derive(Debug, thiserror::Error, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum InprocPublicationError {
#[error("the prepared and current runtimes do not have the same namespace and name")]
IdentityMismatch,
#[error("the current runtime has no published inproc generation")]
CurrentRuntimeUnpublished,
#[error("the replacement runtime has an invalid zero public key")]
ZeroPubkey,
#[error("the expected current inproc generation is no longer published")]
ExpectedGenerationNotCurrent,
#[error("the replacement public key is already occupied")]
ReplacementPubkeyOccupied,
}
static GLOBAL_REGISTRY: OnceLock<InprocRegistry> = OnceLock::new();
#[derive(Clone)]
struct InprocPeer {
name: String,
pubkey: PubKey,
sender: InboxSender,
meta: PeerMeta,
}
#[derive(Default)]
struct NamespaceState {
peers: HashMap<PubKey, InprocPeer>,
names: HashMap<String, PubKey>,
}
struct RegistryState {
namespaces: HashMap<String, NamespaceState>,
}
impl RegistryState {
fn namespace_mut(&mut self, namespace: &str) -> &mut NamespaceState {
self.namespaces.entry(namespace.to_string()).or_default()
}
fn namespace(&self, namespace: &str) -> Option<&NamespaceState> {
self.namespaces.get(namespace)
}
fn namespace_len(&self, namespace: &str) -> usize {
self.namespace(namespace).map_or(0, |ns| ns.peers.len())
}
fn namespace_is_empty(&self, namespace: &str) -> bool {
self.namespace_len(namespace) == 0
}
}
pub struct InprocRegistry {
state: RwLock<RegistryState>,
}
impl InprocRegistry {
pub fn new() -> Self {
Self {
state: RwLock::new(RegistryState {
namespaces: HashMap::new(),
}),
}
}
pub fn global() -> &'static InprocRegistry {
GLOBAL_REGISTRY.get_or_init(InprocRegistry::new)
}
pub fn register(
&self,
name: impl Into<String>,
pubkey: PubKey,
sender: InboxSender,
) -> RegistrationOutcome {
self.register_with_meta_in_namespace(
DEFAULT_NAMESPACE,
name,
pubkey,
sender,
PeerMeta::default(),
)
}
pub fn register_with_meta_in_namespace(
&self,
namespace: &str,
name: impl Into<String>,
pubkey: PubKey,
sender: InboxSender,
meta: PeerMeta,
) -> RegistrationOutcome {
let name = name.into();
if pubkey.is_zero() {
tracing::warn!(
inproc_namespace = %namespace,
peer_name = %name,
"rejecting zero-pubkey inproc registration"
);
return RegistrationOutcome::Rejected {
reason: RegistrationRejection::ZeroPubkey,
};
}
let peer = InprocPeer {
name: name.clone(),
pubkey,
sender,
meta,
};
let mut state = self.state.write();
let namespace_state = state.namespace_mut(namespace);
let evicted_name = namespace_state
.peers
.get(&pubkey)
.filter(|old_peer| old_peer.name != name)
.map(|old_peer| old_peer.name.clone());
if let Some(old_name) = &evicted_name {
namespace_state.names.remove(old_name);
}
let evicted_pubkey = namespace_state
.names
.get(&name)
.filter(|&&old_pk| old_pk != pubkey)
.copied();
if let Some(old_pubkey) = evicted_pubkey {
namespace_state.peers.remove(&old_pubkey);
}
namespace_state.peers.insert(pubkey, peer);
namespace_state.names.insert(name, pubkey);
match (evicted_name, evicted_pubkey) {
(None, None) => RegistrationOutcome::Registered,
(Some(evicted_name), None) => RegistrationOutcome::ReplacedPubkey { evicted_name },
(None, Some(evicted_pubkey)) => RegistrationOutcome::EvictedName { evicted_pubkey },
(Some(evicted_name), Some(evicted_pubkey)) => {
RegistrationOutcome::ReplacedPubkeyAndEvictedName {
evicted_name,
evicted_pubkey,
}
}
}
}
pub(crate) fn replace_sender_in_namespace(
&self,
namespace: &str,
name: &str,
current: (&PubKey, &InboxSender),
replacement_pubkey: PubKey,
replacement_sender: InboxSender,
) -> Result<(), InprocPublicationError> {
let (current_pubkey, current_sender) = current;
if replacement_pubkey.is_zero() {
return Err(InprocPublicationError::ZeroPubkey);
}
let mut state = self.state.write();
let Some(namespace_state) = state.namespaces.get_mut(namespace) else {
return Err(InprocPublicationError::ExpectedGenerationNotCurrent);
};
let current_meta = namespace_state
.peers
.get(current_pubkey)
.filter(|peer| peer.name == name && peer.sender.same_inbox(current_sender))
.map(|peer| peer.meta.clone());
if namespace_state.names.get(name) != Some(current_pubkey) {
return Err(InprocPublicationError::ExpectedGenerationNotCurrent);
}
let Some(current_meta) = current_meta else {
return Err(InprocPublicationError::ExpectedGenerationNotCurrent);
};
if replacement_pubkey != *current_pubkey
&& namespace_state.peers.contains_key(&replacement_pubkey)
{
return Err(InprocPublicationError::ReplacementPubkeyOccupied);
}
namespace_state.peers.remove(current_pubkey);
namespace_state.peers.insert(
replacement_pubkey,
InprocPeer {
name: name.to_string(),
pubkey: replacement_pubkey,
sender: replacement_sender,
meta: current_meta,
},
);
namespace_state
.names
.insert(name.to_string(), replacement_pubkey);
Ok(())
}
pub fn unregister(&self, pubkey: &PubKey) -> bool {
self.unregister_in_namespace(DEFAULT_NAMESPACE, pubkey)
}
pub fn unregister_in_namespace(&self, namespace: &str, pubkey: &PubKey) -> bool {
let mut state = self.state.write();
if let Some(namespace_state) = state.namespaces.get_mut(namespace)
&& let Some(peer) = namespace_state.peers.remove(pubkey)
{
namespace_state.names.remove(&peer.name);
return true;
}
false
}
pub(crate) fn unregister_sender_in_namespace(
&self,
namespace: &str,
pubkey: &PubKey,
sender: &InboxSender,
) -> bool {
let mut state = self.state.write();
let Some(namespace_state) = state.namespaces.get_mut(namespace) else {
return false;
};
if !namespace_state
.peers
.get(pubkey)
.is_some_and(|peer| peer.sender.same_inbox(sender))
{
return false;
}
let Some(peer) = namespace_state.peers.remove(pubkey) else {
return false;
};
if namespace_state.names.get(&peer.name) == Some(pubkey) {
namespace_state.names.remove(&peer.name);
}
true
}
pub(crate) fn update_meta_for_sender_in_namespace(
&self,
namespace: &str,
name: &str,
pubkey: &PubKey,
sender: &InboxSender,
meta: PeerMeta,
) -> bool {
let mut state = self.state.write();
let Some(namespace_state) = state.namespaces.get_mut(namespace) else {
return false;
};
let Some(peer) = namespace_state.peers.get_mut(pubkey) else {
return false;
};
if peer.name != name
|| !peer.sender.same_inbox(sender)
|| namespace_state.names.get(name) != Some(pubkey)
{
return false;
}
peer.meta = meta;
true
}
pub fn get_by_pubkey(&self, pubkey: &PubKey) -> Option<InboxSender> {
self.get_by_pubkey_in_namespace(DEFAULT_NAMESPACE, pubkey)
}
pub fn get_by_pubkey_in_namespace(
&self,
namespace: &str,
pubkey: &PubKey,
) -> Option<InboxSender> {
if pubkey.is_zero() {
return None;
}
self.state
.read()
.namespace(namespace)?
.peers
.get(pubkey)
.map(|p| p.sender.clone())
}
pub(crate) fn get_by_pubkey_any_namespace(&self, pubkey: &PubKey) -> Option<InboxSender> {
if pubkey.is_zero() {
return None;
}
let state = self.state.read();
let mut found = None;
for namespace_state in state.namespaces.values() {
if let Some(peer) = namespace_state.peers.get(pubkey) {
if found.is_some() {
return None;
}
found = Some(peer.sender.clone());
}
}
found
}
pub fn get_name_by_pubkey(&self, pubkey: &PubKey) -> Option<String> {
self.get_name_by_pubkey_in_namespace(DEFAULT_NAMESPACE, pubkey)
}
pub fn get_name_by_pubkey_in_namespace(
&self,
namespace: &str,
pubkey: &PubKey,
) -> Option<String> {
if pubkey.is_zero() {
return None;
}
self.state
.read()
.namespace(namespace)?
.peers
.get(pubkey)
.map(|peer| peer.name.clone())
}
pub fn contains(&self, pubkey: &PubKey) -> bool {
self.state
.read()
.namespace(DEFAULT_NAMESPACE)
.is_some_and(|ns| ns.peers.contains_key(pubkey))
}
pub fn contains_name(&self, name: &str) -> bool {
self.state
.read()
.namespace(DEFAULT_NAMESPACE)
.is_some_and(|ns| ns.names.contains_key(name))
}
pub fn len(&self) -> usize {
self.state.read().namespace_len(DEFAULT_NAMESPACE)
}
pub fn is_empty(&self) -> bool {
self.state.read().namespace_is_empty(DEFAULT_NAMESPACE)
}
pub fn clear(&self) {
self.state.write().namespaces.clear();
}
pub(crate) async fn send_to_pubkey_any_namespace_with_id_wait(
&self,
from_keypair: &Keypair,
to_pubkey: &PubKey,
envelope_id: Uuid,
kind: MessageKind,
sign_envelope: bool,
) -> Result<uuid::Uuid, InprocSendError> {
let sender = self
.get_by_pubkey_any_namespace(to_pubkey)
.ok_or_else(|| InprocSendError::PeerNotFound(to_pubkey.to_peer_id().to_string()))?;
Self::deliver_to_sender_wait(
from_keypair,
*to_pubkey,
sender,
envelope_id,
kind,
sign_envelope,
)
.await
}
pub(crate) async fn send_to_pubkey_in_namespace_with_id_wait(
&self,
namespace: &str,
from_keypair: &Keypair,
to_pubkey: &PubKey,
envelope_id: Uuid,
kind: MessageKind,
sign_envelope: bool,
) -> Result<uuid::Uuid, InprocSendError> {
let sender = self
.get_by_pubkey_in_namespace(namespace, to_pubkey)
.ok_or_else(|| InprocSendError::PeerNotFound(to_pubkey.to_peer_id().to_string()))?;
Self::deliver_to_sender_wait(
from_keypair,
*to_pubkey,
sender,
envelope_id,
kind,
sign_envelope,
)
.await
}
async fn deliver_to_sender_wait(
from_keypair: &Keypair,
to_pubkey: PubKey,
sender: InboxSender,
envelope_id: Uuid,
kind: MessageKind,
sign_envelope: bool,
) -> Result<uuid::Uuid, InprocSendError> {
let mut envelope = Envelope {
id: envelope_id,
from: from_keypair.public_key(),
to: to_pubkey,
kind,
sig: Signature::new([0u8; 64]),
};
if sign_envelope {
envelope.sign(from_keypair);
}
let envelope_id = envelope.id;
match sender.send_wait(InboxItem::External { envelope }).await {
AdmissionOutcome::Admitted => {}
AdmissionOutcome::Dropped {
reason: DropReason::SessionClosed,
} => return Err(InprocSendError::InboxClosed),
AdmissionOutcome::Dropped {
reason: DropReason::InboxFull,
} => return Err(InprocSendError::InboxFull),
AdmissionOutcome::Dropped { reason } => {
return Err(InprocSendError::IngressDropped(reason));
}
}
Ok(envelope_id)
}
pub fn peer_names_in_namespace(&self, namespace: &str) -> Vec<String> {
self.state
.read()
.namespace(namespace)
.map_or_else(Vec::new, |ns| ns.names.keys().cloned().collect())
}
pub fn peers(&self) -> Vec<InprocPeerInfo> {
self.peers_in_namespace(DEFAULT_NAMESPACE)
}
pub fn peers_in_namespace(&self, namespace: &str) -> Vec<InprocPeerInfo> {
self.state
.read()
.namespace(namespace)
.map_or_else(Vec::new, |ns| {
ns.peers
.values()
.map(|peer| InprocPeerInfo {
name: peer.name.clone(),
pubkey: peer.pubkey,
meta: peer.meta.clone(),
})
.collect()
})
}
}
impl Default for InprocRegistry {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, thiserror::Error)]
pub enum InprocSendError {
#[error("Inproc peer not found: {0}")]
PeerNotFound(String),
#[error("Peer inbox has been closed")]
InboxClosed,
#[error("Peer inbox is full")]
InboxFull,
#[error("Peer inbox dropped ingress: {0:?}")]
IngressDropped(crate::inbox::DropReason),
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::classify::test_support;
use crate::inbox::Inbox;
use crate::trust::TrustStore;
fn classified_inbox() -> (Inbox, crate::InboxSender) {
Inbox::new_classified(test_support::classification_context(
TrustStore::new(),
false,
))
}
fn make_keypair() -> Keypair {
Keypair::generate()
}
#[test]
fn test_registry_new() {
let registry = InprocRegistry::new();
assert!(registry.is_empty());
assert_eq!(registry.len(), 0);
}
#[test]
fn test_registry_register_and_lookup() {
let registry = InprocRegistry::new();
let keypair = make_keypair();
let pubkey = keypair.public_key();
let (_, sender) = classified_inbox();
registry.register("test-agent", pubkey, sender);
assert!(!registry.is_empty());
assert_eq!(registry.len(), 1);
assert!(registry.contains(&pubkey));
assert!(registry.contains_name("test-agent"));
assert_eq!(
registry.get_name_by_pubkey(&pubkey).as_deref(),
Some("test-agent")
);
assert!(registry.get_by_pubkey(&pubkey).is_some());
}
#[test]
fn test_registry_rejects_zero_pubkey_registration() {
let registry = InprocRegistry::new();
let (_, sender) = classified_inbox();
let zero_pubkey = PubKey::new([0u8; 32]);
registry.register("zero-agent", zero_pubkey, sender);
assert!(registry.is_empty());
assert!(!registry.contains_name("zero-agent"));
assert!(registry.get_by_pubkey(&zero_pubkey).is_none());
}
#[test]
fn test_registry_zero_pubkey_registration_does_not_shadow_valid_name() {
let registry = InprocRegistry::new();
let valid_keypair = make_keypair();
let valid_pubkey = valid_keypair.public_key();
let (_, valid_sender) = classified_inbox();
let (_, zero_sender) = classified_inbox();
let zero_pubkey = PubKey::new([0u8; 32]);
registry.register("stable-agent", valid_pubkey, valid_sender);
registry.register("stable-agent", zero_pubkey, zero_sender);
assert_eq!(registry.len(), 1);
assert!(registry.contains(&valid_pubkey));
assert!(registry.contains_name("stable-agent"));
assert!(registry.get_by_pubkey(&valid_pubkey).is_some());
assert!(registry.get_by_pubkey(&zero_pubkey).is_none());
assert_eq!(
registry.get_name_by_pubkey(&valid_pubkey).as_deref(),
Some("stable-agent"),
"valid name mapping should remain"
);
}
#[test]
fn test_registry_unregister() {
let registry = InprocRegistry::new();
let keypair = make_keypair();
let pubkey = keypair.public_key();
let (_, sender) = classified_inbox();
registry.register("test-agent", pubkey, sender);
assert!(registry.contains(&pubkey));
let removed = registry.unregister(&pubkey);
assert!(removed);
assert!(!registry.contains(&pubkey));
assert!(!registry.contains_name("test-agent"));
assert!(registry.is_empty());
let removed_again = registry.unregister(&pubkey);
assert!(!removed_again);
}
#[test]
fn test_registry_replace_on_same_pubkey() {
let registry = InprocRegistry::new();
let keypair = make_keypair();
let pubkey = keypair.public_key();
let (_, sender1) = classified_inbox();
let (_, sender2) = classified_inbox();
registry.register("agent-v1", pubkey, sender1);
assert!(registry.contains_name("agent-v1"));
registry.register("agent-v2", pubkey, sender2);
assert!(!registry.contains_name("agent-v1"));
assert!(registry.contains_name("agent-v2"));
assert_eq!(registry.len(), 1);
}
#[test]
fn test_registry_replace_on_same_name_different_pubkey() {
let registry = InprocRegistry::new();
let keypair1 = make_keypair();
let pubkey1 = keypair1.public_key();
let keypair2 = make_keypair();
let pubkey2 = keypair2.public_key();
let (_, sender1) = classified_inbox();
let (_, sender2) = classified_inbox();
registry.register("my-agent", pubkey1, sender1);
assert!(registry.contains(&pubkey1));
assert!(registry.contains_name("my-agent"));
assert_eq!(registry.len(), 1);
registry.register("my-agent", pubkey2, sender2);
assert!(!registry.contains(&pubkey1), "old pubkey should be evicted");
assert!(registry.contains(&pubkey2));
assert!(registry.contains_name("my-agent"));
assert_eq!(registry.len(), 1);
assert_eq!(
registry.get_name_by_pubkey(&pubkey2).as_deref(),
Some("my-agent")
);
}
#[test]
fn registration_outcome_is_typed_for_displacement_and_rejection() {
let registry = InprocRegistry::new();
let keypair = make_keypair();
let pubkey = keypair.public_key();
let (_, sender1) = classified_inbox();
let (_, sender2) = classified_inbox();
let (_, sender3) = classified_inbox();
assert_eq!(
registry.register("agent-v1", pubkey, sender1),
RegistrationOutcome::Registered
);
assert_eq!(
registry.register("agent-v2", pubkey, sender2),
RegistrationOutcome::ReplacedPubkey {
evicted_name: "agent-v1".to_string()
}
);
let other = make_keypair();
let other_pubkey = other.public_key();
match registry.register("agent-v2", other_pubkey, sender3) {
RegistrationOutcome::EvictedName { evicted_pubkey } => {
assert_eq!(evicted_pubkey, pubkey);
}
other => panic!("expected EvictedName, got {other:?}"),
}
let (_, zero_sender) = classified_inbox();
let zero_pubkey = PubKey::new([0u8; 32]);
assert_eq!(
registry.register("zero", zero_pubkey, zero_sender),
RegistrationOutcome::Rejected {
reason: RegistrationRejection::ZeroPubkey
}
);
assert!(!registry.contains_name("zero"));
}
#[test]
fn test_registry_aba_scenario_safe() {
let registry = InprocRegistry::new();
let keypair_old = make_keypair();
let pubkey_old = keypair_old.public_key();
let keypair_new = make_keypair();
let pubkey_new = keypair_new.public_key();
let (_, sender_old) = classified_inbox();
let (_, sender_new) = classified_inbox();
registry.register("agent", pubkey_old, sender_old);
assert!(registry.contains(&pubkey_old));
registry.register("agent", pubkey_new, sender_new);
assert!(
!registry.contains(&pubkey_old),
"old pubkey should be evicted"
);
assert!(registry.contains(&pubkey_new));
let removed = registry.unregister(&pubkey_old);
assert!(!removed, "unregister of evicted pubkey should return false");
assert!(
registry.contains(&pubkey_new),
"new agent should still be registered"
);
assert!(
registry.contains_name("agent"),
"name should still map to new agent"
);
assert_eq!(
registry.get_name_by_pubkey(&pubkey_new).as_deref(),
Some("agent"),
"name should map to the new pubkey"
);
}
#[test]
fn test_registry_peer_names_in_namespace() {
let registry = InprocRegistry::new();
for i in 0..3 {
let keypair = make_keypair();
let (_, sender) = classified_inbox();
registry.register_with_meta_in_namespace(
"realm-names",
format!("agent-{i}"),
keypair.public_key(),
sender,
PeerMeta::default(),
);
}
let names = registry.peer_names_in_namespace("realm-names");
assert_eq!(names.len(), 3);
assert!(names.contains(&"agent-0".to_string()));
assert!(names.contains(&"agent-1".to_string()));
assert!(names.contains(&"agent-2".to_string()));
assert!(registry.peer_names_in_namespace("realm-other").is_empty());
}
#[test]
fn test_registry_peers_snapshot() {
let registry = InprocRegistry::new();
let keypair = make_keypair();
let pubkey = keypair.public_key();
let (_, sender) = classified_inbox();
registry.register("agent-a", pubkey, sender);
let peers = registry.peers();
assert_eq!(peers.len(), 1);
assert_eq!(peers[0].name, "agent-a");
assert_eq!(peers[0].pubkey, pubkey);
}
#[test]
fn test_registry_clear() {
let registry = InprocRegistry::new();
for i in 0..3 {
let keypair = make_keypair();
let (_, sender) = classified_inbox();
registry.register(format!("agent-{i}"), keypair.public_key(), sender);
}
assert_eq!(registry.len(), 3);
registry.clear();
assert!(registry.is_empty());
}
#[tokio::test]
async fn test_registry_send_delivers_to_inbox() {
let registry = InprocRegistry::new();
let receiver_keypair = make_keypair();
let (mut inbox, sender) = classified_inbox();
registry.register("receiver", receiver_keypair.public_key(), sender);
let sender_keypair = make_keypair();
let result = registry
.send_to_pubkey_in_namespace_with_id_wait(
"",
&sender_keypair,
&receiver_keypair.public_key(),
Uuid::new_v4(),
MessageKind::Message {
objective_id: None,
content_taint: None,
blocks: None,
body: "hello inproc".to_string(),
handling_mode: None,
},
true,
)
.await;
assert!(result.is_ok());
let items = inbox.try_drain_classified();
assert_eq!(items.len(), 1);
match &items[0].item {
InboxItem::External { envelope } => {
assert_eq!(envelope.from, sender_keypair.public_key());
assert_eq!(envelope.to, receiver_keypair.public_key());
match &envelope.kind {
MessageKind::Message {
blocks: None, body, ..
} => {
assert_eq!(body, "hello inproc");
}
_ => panic!("expected Message kind"),
}
assert!(envelope.verify());
}
_ => panic!("expected External inbox item"),
}
}
#[tokio::test]
async fn test_registry_send_peer_not_found() {
let registry = InprocRegistry::new();
let sender_keypair = make_keypair();
let unknown = make_keypair().public_key();
let result = registry
.send_to_pubkey_in_namespace_with_id_wait(
"",
&sender_keypair,
&unknown,
Uuid::new_v4(),
MessageKind::Message {
objective_id: None,
content_taint: None,
blocks: None,
body: "hello".to_string(),
handling_mode: None,
},
true,
)
.await;
assert!(matches!(result, Err(InprocSendError::PeerNotFound(_))));
}
#[tokio::test]
async fn test_registry_send_inbox_closed() {
let registry = InprocRegistry::new();
let receiver_keypair = make_keypair();
let (inbox, sender) = classified_inbox();
registry.register("receiver", receiver_keypair.public_key(), sender);
drop(inbox);
let sender_keypair = make_keypair();
let result = registry
.send_to_pubkey_in_namespace_with_id_wait(
"",
&sender_keypair,
&receiver_keypair.public_key(),
Uuid::new_v4(),
MessageKind::Message {
objective_id: None,
content_taint: None,
blocks: None,
body: "hello".to_string(),
handling_mode: None,
},
true,
)
.await;
assert!(matches!(result, Err(InprocSendError::InboxClosed)));
}
#[tokio::test]
async fn test_registry_namespace_isolation_for_lookup_and_send() {
let registry = InprocRegistry::new();
let receiver_keypair = make_keypair();
let (mut inbox, sender) = classified_inbox();
registry.register_with_meta_in_namespace(
"realm-a",
"receiver",
receiver_keypair.public_key(),
sender,
PeerMeta::default(),
);
assert!(
registry
.get_by_pubkey(&receiver_keypair.public_key())
.is_none()
);
assert!(
registry
.get_by_pubkey_in_namespace("realm-a", &receiver_keypair.public_key())
.is_some()
);
let sender_keypair = make_keypair();
let ok = registry
.send_to_pubkey_in_namespace_with_id_wait(
"realm-a",
&sender_keypair,
&receiver_keypair.public_key(),
Uuid::new_v4(),
MessageKind::Message {
objective_id: None,
content_taint: None,
blocks: None,
body: "hello scoped".to_string(),
handling_mode: None,
},
true,
)
.await;
assert!(ok.is_ok());
let wrong_ns = registry
.send_to_pubkey_in_namespace_with_id_wait(
"realm-b",
&sender_keypair,
&receiver_keypair.public_key(),
Uuid::new_v4(),
MessageKind::Message {
objective_id: None,
content_taint: None,
blocks: None,
body: "should not deliver".to_string(),
handling_mode: None,
},
true,
)
.await;
assert!(matches!(wrong_ns, Err(InprocSendError::PeerNotFound(_))));
let items = inbox.try_drain_classified();
assert_eq!(items.len(), 1);
}
#[tokio::test]
async fn test_send_to_pubkey_in_namespace_ignores_display_name_collision() {
let registry = InprocRegistry::new();
let target_keypair = make_keypair();
let target_pubkey = target_keypair.public_key();
let shadow_keypair = make_keypair();
let shadow_pubkey = shadow_keypair.public_key();
let (mut target_inbox, target_sender) = classified_inbox();
let (mut shadow_inbox, shadow_sender) = classified_inbox();
registry.register_with_meta_in_namespace(
"",
"canonical-target",
target_pubkey,
target_sender,
PeerMeta::default(),
);
registry.register_with_meta_in_namespace(
"",
"shared-display-name",
shadow_pubkey,
shadow_sender,
PeerMeta::default(),
);
let sender_keypair = make_keypair();
let result = registry
.send_to_pubkey_in_namespace_with_id_wait(
"",
&sender_keypair,
&target_pubkey,
Uuid::new_v4(),
MessageKind::Message {
objective_id: None,
content_taint: None,
blocks: None,
body: "hello canonical".to_string(),
handling_mode: None,
},
true,
)
.await;
assert!(result.is_ok());
assert_eq!(shadow_inbox.try_drain_classified().len(), 0);
let items = target_inbox.try_drain_classified();
assert_eq!(items.len(), 1);
let InboxItem::External { envelope } = &items[0].item else {
panic!("expected external envelope");
};
assert_eq!(envelope.to, target_pubkey);
}
#[tokio::test]
async fn test_send_to_pubkey_any_namespace_rejects_ambiguous_identity() {
let registry = InprocRegistry::new();
let sender_keypair = make_keypair();
let target_keypair = make_keypair();
let target_pubkey = target_keypair.public_key();
let (mut alpha_inbox, alpha_sender) = classified_inbox();
let (mut beta_inbox, beta_sender) = classified_inbox();
registry.register_with_meta_in_namespace(
"realm-alpha",
"alpha-target",
target_pubkey,
alpha_sender,
PeerMeta::default(),
);
registry.register_with_meta_in_namespace(
"realm-beta",
"beta-target",
target_pubkey,
beta_sender,
PeerMeta::default(),
);
let result = registry
.send_to_pubkey_any_namespace_with_id_wait(
&sender_keypair,
&target_pubkey,
Uuid::new_v4(),
MessageKind::Message {
objective_id: None,
content_taint: None,
blocks: None,
body: "ambiguous identity".to_string(),
handling_mode: None,
},
true,
)
.await;
assert!(matches!(result, Err(InprocSendError::PeerNotFound(_))));
assert!(alpha_inbox.try_drain_classified().is_empty());
assert!(beta_inbox.try_drain_classified().is_empty());
}
#[test]
fn test_registry_same_name_can_exist_in_different_namespaces() {
let registry = InprocRegistry::new();
let kp_a = make_keypair();
let kp_b = make_keypair();
let (_, sender_a) = classified_inbox();
let (_, sender_b) = classified_inbox();
registry.register_with_meta_in_namespace(
"realm-a",
"shared-name",
kp_a.public_key(),
sender_a,
PeerMeta::default(),
);
registry.register_with_meta_in_namespace(
"realm-b",
"shared-name",
kp_b.public_key(),
sender_b,
PeerMeta::default(),
);
assert_eq!(
registry
.get_name_by_pubkey_in_namespace("realm-a", &kp_a.public_key())
.as_deref(),
Some("shared-name")
);
assert_eq!(
registry
.get_name_by_pubkey_in_namespace("realm-b", &kp_b.public_key())
.as_deref(),
Some("shared-name")
);
assert_ne!(kp_a.public_key(), kp_b.public_key());
assert!(
!registry.contains_name("shared-name"),
"default namespace must not see namespaced registrations"
);
}
#[test]
fn test_global_registry() {
let registry = InprocRegistry::global();
registry.clear();
let keypair = make_keypair();
let (_, sender) = classified_inbox();
registry.register("global-test", keypair.public_key(), sender);
assert!(registry.contains_name("global-test"));
registry.unregister(&keypair.public_key());
}
#[test]
fn test_registry_register_with_meta() {
let registry = InprocRegistry::new();
let keypair = make_keypair();
let pubkey = keypair.public_key();
let (_, sender) = classified_inbox();
let meta = PeerMeta::default()
.with_description("Reviews code for style issues")
.with_label("lang", "rust");
registry.register_with_meta_in_namespace(
DEFAULT_NAMESPACE,
"reviewer",
pubkey,
sender,
meta.clone(),
);
let peers = registry.peers();
assert_eq!(peers.len(), 1);
assert_eq!(peers[0].name, "reviewer");
assert_eq!(peers[0].pubkey, pubkey);
assert_eq!(peers[0].meta, meta);
}
#[test]
fn test_registry_peers_returns_default_meta_for_plain_register() {
let registry = InprocRegistry::new();
let keypair = make_keypair();
let pubkey = keypair.public_key();
let (_, sender) = classified_inbox();
registry.register("plain-agent", pubkey, sender);
let peers = registry.peers();
assert_eq!(peers.len(), 1);
assert_eq!(peers[0].meta, PeerMeta::default());
}
}