use std::collections::HashMap;
use std::sync::OnceLock;
use parking_lot::RwLock;
use uuid::Uuid;
use crate::identity::{Keypair, PubKey, Signature};
use crate::inbox::{InboxError, InboxSender};
use crate::types::{Envelope, InboxItem, MessageKind};
static GLOBAL_REGISTRY: OnceLock<InprocRegistry> = OnceLock::new();
#[derive(Clone)]
struct InprocPeer {
name: String,
pubkey: PubKey,
sender: InboxSender,
}
struct RegistryState {
peers: HashMap<PubKey, InprocPeer>,
names: HashMap<String, PubKey>,
}
impl RegistryState {
fn new() -> Self {
Self {
peers: HashMap::new(),
names: HashMap::new(),
}
}
}
pub struct InprocRegistry {
state: RwLock<RegistryState>,
}
impl InprocRegistry {
pub fn new() -> Self {
Self {
state: RwLock::new(RegistryState::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) {
let name = name.into();
let peer = InprocPeer {
name: name.clone(),
pubkey,
sender,
};
let mut state = self.state.write();
let old_name_to_remove = state
.peers
.get(&pubkey)
.filter(|old_peer| old_peer.name != name)
.map(|old_peer| old_peer.name.clone());
if let Some(old_name) = old_name_to_remove {
state.names.remove(&old_name);
}
let old_pubkey_to_remove = state
.names
.get(&name)
.filter(|&&old_pk| old_pk != pubkey)
.copied();
if let Some(old_pubkey) = old_pubkey_to_remove {
state.peers.remove(&old_pubkey);
}
state.peers.insert(pubkey, peer);
state.names.insert(name, pubkey);
}
pub fn unregister(&self, pubkey: &PubKey) -> bool {
let mut state = self.state.write();
if let Some(peer) = state.peers.remove(pubkey) {
state.names.remove(&peer.name);
true
} else {
false
}
}
pub fn get_by_name(&self, name: &str) -> Option<(PubKey, InboxSender)> {
let state = self.state.read();
let pubkey = state.names.get(name).copied()?;
let peer = state.peers.get(&pubkey)?;
Some((peer.pubkey, peer.sender.clone()))
}
pub fn get_by_pubkey(&self, pubkey: &PubKey) -> Option<InboxSender> {
self.state
.read()
.peers
.get(pubkey)
.map(|p| p.sender.clone())
}
pub fn get_name_by_pubkey(&self, pubkey: &PubKey) -> Option<String> {
self.state
.read()
.peers
.get(pubkey)
.map(|peer| peer.name.clone())
}
pub fn contains(&self, pubkey: &PubKey) -> bool {
self.state.read().peers.contains_key(pubkey)
}
pub fn contains_name(&self, name: &str) -> bool {
self.state.read().names.contains_key(name)
}
pub fn len(&self) -> usize {
self.state.read().peers.len()
}
pub fn is_empty(&self) -> bool {
self.state.read().peers.is_empty()
}
pub fn clear(&self) {
let mut state = self.state.write();
state.peers.clear();
state.names.clear();
}
pub fn send(
&self,
from_keypair: &Keypair,
to_name: &str,
kind: MessageKind,
) -> Result<uuid::Uuid, InprocSendError> {
self.send_with_signature(from_keypair, to_name, kind, true)
}
pub fn send_with_signature(
&self,
from_keypair: &Keypair,
to_name: &str,
kind: MessageKind,
sign_envelope: bool,
) -> Result<uuid::Uuid, InprocSendError> {
let (to_pubkey, sender) = self
.get_by_name(to_name)
.ok_or_else(|| InprocSendError::PeerNotFound(to_name.to_string()))?;
let mut envelope = Envelope {
id: Uuid::new_v4(),
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;
sender
.send(InboxItem::External { envelope })
.map_err(|err| match err {
InboxError::Closed => InprocSendError::InboxClosed,
InboxError::Full => InprocSendError::InboxFull,
})?;
Ok(envelope_id)
}
pub fn peer_names(&self) -> Vec<String> {
self.state.read().names.keys().cloned().collect()
}
pub fn peers(&self) -> Vec<(String, PubKey)> {
self.state
.read()
.peers
.values()
.map(|peer| (peer.name.clone(), peer.pubkey))
.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,
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::inbox::Inbox;
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) = Inbox::new();
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"));
let (found_pubkey, _) = registry.get_by_name("test-agent").unwrap();
assert_eq!(found_pubkey, pubkey);
assert!(registry.get_by_pubkey(&pubkey).is_some());
}
#[test]
fn test_registry_unregister() {
let registry = InprocRegistry::new();
let keypair = make_keypair();
let pubkey = keypair.public_key();
let (_, sender) = Inbox::new();
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) = Inbox::new();
let (_, sender2) = Inbox::new();
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) = Inbox::new();
let (_, sender2) = Inbox::new();
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);
let (found_pubkey, _) = registry.get_by_name("my-agent").unwrap();
assert_eq!(found_pubkey, pubkey2);
}
#[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) = Inbox::new();
let (_, sender_new) = Inbox::new();
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"
);
let (found_pubkey, _) = registry.get_by_name("agent").unwrap();
assert_eq!(found_pubkey, pubkey_new, "lookup should return new pubkey");
}
#[test]
fn test_registry_peer_names() {
let registry = InprocRegistry::new();
for i in 0..3 {
let keypair = make_keypair();
let (_, sender) = Inbox::new();
registry.register(format!("agent-{}", i), keypair.public_key(), sender);
}
let names = registry.peer_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()));
}
#[test]
fn test_registry_peers_snapshot() {
let registry = InprocRegistry::new();
let keypair = make_keypair();
let pubkey = keypair.public_key();
let (_, sender) = Inbox::new();
registry.register("agent-a", pubkey, sender);
let peers = registry.peers();
assert_eq!(peers.len(), 1);
assert_eq!(peers[0].0, "agent-a");
assert_eq!(peers[0].1, pubkey);
}
#[test]
fn test_registry_clear() {
let registry = InprocRegistry::new();
for i in 0..3 {
let keypair = make_keypair();
let (_, sender) = Inbox::new();
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) = Inbox::new();
registry.register("receiver", receiver_keypair.public_key(), sender);
let sender_keypair = make_keypair();
let result = registry.send(
&sender_keypair,
"receiver",
MessageKind::Message {
body: "hello inproc".to_string(),
},
);
assert!(result.is_ok());
let items = inbox.try_drain();
assert_eq!(items.len(), 1);
match &items[0] {
InboxItem::External { envelope } => {
assert_eq!(envelope.from, sender_keypair.public_key());
assert_eq!(envelope.to, receiver_keypair.public_key());
match &envelope.kind {
MessageKind::Message { body } => {
assert_eq!(body, "hello inproc");
}
_ => panic!("expected Message kind"),
}
assert!(envelope.verify());
}
_ => panic!("expected External inbox item"),
}
}
#[test]
fn test_registry_send_peer_not_found() {
let registry = InprocRegistry::new();
let sender_keypair = make_keypair();
let result = registry.send(
&sender_keypair,
"nonexistent",
MessageKind::Message {
body: "hello".to_string(),
},
);
assert!(matches!(result, Err(InprocSendError::PeerNotFound(_))));
}
#[test]
fn test_registry_send_inbox_closed() {
let registry = InprocRegistry::new();
let receiver_keypair = make_keypair();
let (inbox, sender) = Inbox::new();
registry.register("receiver", receiver_keypair.public_key(), sender);
drop(inbox);
let sender_keypair = make_keypair();
let result = registry.send(
&sender_keypair,
"receiver",
MessageKind::Message {
body: "hello".to_string(),
},
);
assert!(matches!(result, Err(InprocSendError::InboxClosed)));
}
#[test]
fn test_global_registry() {
let registry = InprocRegistry::global();
registry.clear();
let keypair = make_keypair();
let (_, sender) = Inbox::new();
registry.register("global-test", keypair.public_key(), sender);
assert!(registry.contains_name("global-test"));
registry.unregister(&keypair.public_key());
}
}