Skip to main content

meerkat_comms/
inproc.rs

1//! In-process message transport for peer communication within one runtime.
2//!
3//! This module provides a process-global registry that allows agents within
4//! the same process to communicate without network sockets. Messages are
5//! delivered directly via in-memory channels.
6//!
7//! # Usage
8//!
9//! ```text
10//! // Register an agent's inbox
11//! let (inbox, sender) = Inbox::new();
12//! InprocRegistry::global().register("my-agent", pubkey, sender);
13//!
14//! // Delivery is pubkey-keyed: the Router resolves a trusted peer's
15//! // signing key and delivers through the namespace-scoped
16//! // send_to_pubkey_*_wait owners.
17//!
18//! // Unregister when done
19//! InprocRegistry::global().unregister(&pubkey);
20//! ```
21
22use std::collections::HashMap;
23use std::sync::OnceLock;
24
25use parking_lot::RwLock;
26use uuid::Uuid;
27
28use crate::identity::{Keypair, PubKey, Signature};
29use crate::inbox::{AdmissionOutcome, DropReason, InboxSender};
30use crate::peer_meta::PeerMeta;
31use crate::types::{Envelope, InboxItem, MessageKind};
32
33const DEFAULT_NAMESPACE: &str = "";
34
35/// Snapshot of an inproc peer returned by [`InprocRegistry::peers()`].
36#[derive(Debug, Clone)]
37pub struct InprocPeerInfo {
38    pub name: String,
39    pub pubkey: PubKey,
40    pub meta: PeerMeta,
41}
42
43/// Why a registration was rejected without mutating the registry.
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
45pub enum RegistrationRejection {
46    /// The supplied pubkey was the all-zero key, which can never identify a
47    /// distinct peer and is refused fail-closed.
48    ZeroPubkey,
49}
50
51/// Typed result of registering an inproc peer.
52///
53/// Registration is not always a clean insert: re-registering an existing
54/// pubkey under a new name evicts the old name mapping, and re-registering an
55/// existing name with a new pubkey evicts the old pubkey entry. Both evictions
56/// can also happen at once. Callers (runtime constructors, metadata refresh)
57/// must observe these facts rather than assume a clean success.
58#[derive(Debug, Clone, PartialEq, Eq)]
59pub enum RegistrationOutcome {
60    /// The peer was inserted without displacing any existing route.
61    Registered,
62    /// This pubkey was already registered under a different name; the old name
63    /// mapping was removed and replaced with the new name.
64    ReplacedPubkey { evicted_name: String },
65    /// This name was already bound to a different pubkey; the old pubkey entry
66    /// was evicted so the stale key is no longer reachable.
67    EvictedName { evicted_pubkey: PubKey },
68    /// Both evictions happened: this pubkey's old name was removed AND this
69    /// name's old pubkey was evicted in the same registration.
70    ReplacedPubkeyAndEvictedName {
71        evicted_name: String,
72        evicted_pubkey: PubKey,
73    },
74    /// The registration was refused without mutating the registry.
75    Rejected { reason: RegistrationRejection },
76}
77
78impl RegistrationOutcome {
79    /// Whether the registration displaced an existing route (either eviction).
80    pub fn displaced_existing(&self) -> bool {
81        matches!(
82            self,
83            Self::ReplacedPubkey { .. }
84                | Self::EvictedName { .. }
85                | Self::ReplacedPubkeyAndEvictedName { .. }
86        )
87    }
88
89    /// Whether the registration was rejected (no mutation occurred).
90    pub fn is_rejected(&self) -> bool {
91        matches!(self, Self::Rejected { .. })
92    }
93}
94
95/// Why publishing an exact inproc runtime replacement failed before mutation.
96#[derive(Debug, thiserror::Error, Clone, Copy, PartialEq, Eq)]
97#[non_exhaustive]
98pub enum InprocPublicationError {
99    #[error("the prepared and current runtimes do not have the same namespace and name")]
100    IdentityMismatch,
101    #[error("the current runtime has no published inproc generation")]
102    CurrentRuntimeUnpublished,
103    #[error("the replacement runtime has an invalid zero public key")]
104    ZeroPubkey,
105    #[error("the expected current inproc generation is no longer published")]
106    ExpectedGenerationNotCurrent,
107    #[error("the replacement public key is already occupied")]
108    ReplacementPubkeyOccupied,
109}
110
111/// Global inproc registry instance.
112static GLOBAL_REGISTRY: OnceLock<InprocRegistry> = OnceLock::new();
113
114/// Registry entry for an inproc peer.
115#[derive(Clone)]
116struct InprocPeer {
117    name: String,
118    pubkey: PubKey,
119    sender: InboxSender,
120    meta: PeerMeta,
121}
122
123/// Internal namespace state protected by a single lock to prevent deadlocks.
124#[derive(Default)]
125struct NamespaceState {
126    /// Map from pubkey to peer entry.
127    peers: HashMap<PubKey, InprocPeer>,
128    /// Map from name to pubkey for name-based lookup.
129    names: HashMap<String, PubKey>,
130}
131
132/// Internal registry state keyed by namespace.
133struct RegistryState {
134    namespaces: HashMap<String, NamespaceState>,
135}
136
137impl RegistryState {
138    fn namespace_mut(&mut self, namespace: &str) -> &mut NamespaceState {
139        self.namespaces.entry(namespace.to_string()).or_default()
140    }
141
142    fn namespace(&self, namespace: &str) -> Option<&NamespaceState> {
143        self.namespaces.get(namespace)
144    }
145
146    fn namespace_len(&self, namespace: &str) -> usize {
147        self.namespace(namespace).map_or(0, |ns| ns.peers.len())
148    }
149
150    fn namespace_is_empty(&self, namespace: &str) -> bool {
151        self.namespace_len(namespace) == 0
152    }
153}
154
155/// Process-global registry for in-process peer communication.
156///
157/// This registry maps agent pubkeys to their inbox senders, allowing
158/// direct message delivery without network transport.
159///
160/// # Thread Safety
161///
162/// All operations are protected by a single RwLock to ensure consistent
163/// state and prevent deadlocks.
164pub struct InprocRegistry {
165    state: RwLock<RegistryState>,
166}
167
168impl InprocRegistry {
169    /// Create a new empty registry.
170    pub fn new() -> Self {
171        Self {
172            state: RwLock::new(RegistryState {
173                namespaces: HashMap::new(),
174            }),
175        }
176    }
177
178    /// Get the global registry instance.
179    ///
180    /// This creates the registry on first access.
181    pub fn global() -> &'static InprocRegistry {
182        GLOBAL_REGISTRY.get_or_init(InprocRegistry::new)
183    }
184
185    /// Register an agent's inbox for inproc communication.
186    ///
187    /// Returns a typed [`RegistrationOutcome`] describing whether the insert was
188    /// clean or displaced an existing route (see
189    /// [`register_with_meta_in_namespace`](Self::register_with_meta_in_namespace)).
190    pub fn register(
191        &self,
192        name: impl Into<String>,
193        pubkey: PubKey,
194        sender: InboxSender,
195    ) -> RegistrationOutcome {
196        self.register_with_meta_in_namespace(
197            DEFAULT_NAMESPACE,
198            name,
199            pubkey,
200            sender,
201            PeerMeta::default(),
202        )
203    }
204
205    /// Register an agent's inbox within an explicit namespace.
206    ///
207    /// Returns a typed [`RegistrationOutcome`] that surfaces route displacement
208    /// explicitly: re-registering an existing pubkey under a new name evicts
209    /// the old name mapping ([`RegistrationOutcome::ReplacedPubkey`]); a name
210    /// rebound to a new pubkey evicts the old pubkey entry
211    /// ([`RegistrationOutcome::EvictedName`]); both can happen at once
212    /// ([`RegistrationOutcome::ReplacedPubkeyAndEvictedName`]). A zero pubkey is
213    /// refused without mutation ([`RegistrationOutcome::Rejected`]). Callers
214    /// must observe displacement/rejection rather than assume a clean success.
215    pub fn register_with_meta_in_namespace(
216        &self,
217        namespace: &str,
218        name: impl Into<String>,
219        pubkey: PubKey,
220        sender: InboxSender,
221        meta: PeerMeta,
222    ) -> RegistrationOutcome {
223        let name = name.into();
224        if pubkey.is_zero() {
225            tracing::warn!(
226                inproc_namespace = %namespace,
227                peer_name = %name,
228                "rejecting zero-pubkey inproc registration"
229            );
230            return RegistrationOutcome::Rejected {
231                reason: RegistrationRejection::ZeroPubkey,
232            };
233        }
234        let peer = InprocPeer {
235            name: name.clone(),
236            pubkey,
237            sender,
238            meta,
239        };
240
241        let mut state = self.state.write();
242        let namespace_state = state.namespace_mut(namespace);
243
244        // If this pubkey was registered under a different name, remove old name mapping
245        let evicted_name = namespace_state
246            .peers
247            .get(&pubkey)
248            .filter(|old_peer| old_peer.name != name)
249            .map(|old_peer| old_peer.name.clone());
250        if let Some(old_name) = &evicted_name {
251            namespace_state.names.remove(old_name);
252        }
253
254        // If this name was registered to a different pubkey, remove the old pubkey entry
255        // This prevents stale pubkeys from remaining reachable
256        let evicted_pubkey = namespace_state
257            .names
258            .get(&name)
259            .filter(|&&old_pk| old_pk != pubkey)
260            .copied();
261        if let Some(old_pubkey) = evicted_pubkey {
262            namespace_state.peers.remove(&old_pubkey);
263        }
264
265        namespace_state.peers.insert(pubkey, peer);
266        namespace_state.names.insert(name, pubkey);
267
268        match (evicted_name, evicted_pubkey) {
269            (None, None) => RegistrationOutcome::Registered,
270            (Some(evicted_name), None) => RegistrationOutcome::ReplacedPubkey { evicted_name },
271            (None, Some(evicted_pubkey)) => RegistrationOutcome::EvictedName { evicted_pubkey },
272            (Some(evicted_name), Some(evicted_pubkey)) => {
273                RegistrationOutcome::ReplacedPubkeyAndEvictedName {
274                    evicted_name,
275                    evicted_pubkey,
276                }
277            }
278        }
279    }
280
281    /// Atomically replace one exact inbox generation.
282    ///
283    /// All checks happen under the registry write lock. A stale predecessor or
284    /// occupied replacement key therefore leaves the live route unchanged.
285    pub(crate) fn replace_sender_in_namespace(
286        &self,
287        namespace: &str,
288        name: &str,
289        current: (&PubKey, &InboxSender),
290        replacement_pubkey: PubKey,
291        replacement_sender: InboxSender,
292    ) -> Result<(), InprocPublicationError> {
293        let (current_pubkey, current_sender) = current;
294        if replacement_pubkey.is_zero() {
295            return Err(InprocPublicationError::ZeroPubkey);
296        }
297
298        let mut state = self.state.write();
299        let Some(namespace_state) = state.namespaces.get_mut(namespace) else {
300            return Err(InprocPublicationError::ExpectedGenerationNotCurrent);
301        };
302        let current_meta = namespace_state
303            .peers
304            .get(current_pubkey)
305            .filter(|peer| peer.name == name && peer.sender.same_inbox(current_sender))
306            .map(|peer| peer.meta.clone());
307        if namespace_state.names.get(name) != Some(current_pubkey) {
308            return Err(InprocPublicationError::ExpectedGenerationNotCurrent);
309        }
310        let Some(current_meta) = current_meta else {
311            return Err(InprocPublicationError::ExpectedGenerationNotCurrent);
312        };
313        if replacement_pubkey != *current_pubkey
314            && namespace_state.peers.contains_key(&replacement_pubkey)
315        {
316            return Err(InprocPublicationError::ReplacementPubkeyOccupied);
317        }
318
319        namespace_state.peers.remove(current_pubkey);
320        namespace_state.peers.insert(
321            replacement_pubkey,
322            InprocPeer {
323                name: name.to_string(),
324                pubkey: replacement_pubkey,
325                sender: replacement_sender,
326                meta: current_meta,
327            },
328        );
329        namespace_state
330            .names
331            .insert(name.to_string(), replacement_pubkey);
332        Ok(())
333    }
334
335    /// Unregister an agent by pubkey.
336    ///
337    /// Returns true if the agent was found and removed.
338    pub fn unregister(&self, pubkey: &PubKey) -> bool {
339        self.unregister_in_namespace(DEFAULT_NAMESPACE, pubkey)
340    }
341
342    /// Unregister an agent by pubkey from an explicit namespace.
343    pub fn unregister_in_namespace(&self, namespace: &str, pubkey: &PubKey) -> bool {
344        let mut state = self.state.write();
345        if let Some(namespace_state) = state.namespaces.get_mut(namespace)
346            && let Some(peer) = namespace_state.peers.remove(pubkey)
347        {
348            namespace_state.names.remove(&peer.name);
349            return true;
350        }
351        false
352    }
353
354    /// Remove a route only when its exact inbox generation is still current.
355    pub(crate) fn unregister_sender_in_namespace(
356        &self,
357        namespace: &str,
358        pubkey: &PubKey,
359        sender: &InboxSender,
360    ) -> bool {
361        let mut state = self.state.write();
362        let Some(namespace_state) = state.namespaces.get_mut(namespace) else {
363            return false;
364        };
365        if !namespace_state
366            .peers
367            .get(pubkey)
368            .is_some_and(|peer| peer.sender.same_inbox(sender))
369        {
370            return false;
371        }
372        let Some(peer) = namespace_state.peers.remove(pubkey) else {
373            return false;
374        };
375        if namespace_state.names.get(&peer.name) == Some(pubkey) {
376            namespace_state.names.remove(&peer.name);
377        }
378        true
379    }
380
381    /// Update metadata only when the exact inbox generation is still current.
382    pub(crate) fn update_meta_for_sender_in_namespace(
383        &self,
384        namespace: &str,
385        name: &str,
386        pubkey: &PubKey,
387        sender: &InboxSender,
388        meta: PeerMeta,
389    ) -> bool {
390        let mut state = self.state.write();
391        let Some(namespace_state) = state.namespaces.get_mut(namespace) else {
392            return false;
393        };
394        let Some(peer) = namespace_state.peers.get_mut(pubkey) else {
395            return false;
396        };
397        if peer.name != name
398            || !peer.sender.same_inbox(sender)
399            || namespace_state.names.get(name) != Some(pubkey)
400        {
401            return false;
402        }
403        peer.meta = meta;
404        true
405    }
406
407    /// Look up an inproc peer by pubkey.
408    pub fn get_by_pubkey(&self, pubkey: &PubKey) -> Option<InboxSender> {
409        self.get_by_pubkey_in_namespace(DEFAULT_NAMESPACE, pubkey)
410    }
411
412    /// Look up an inproc peer by pubkey in an explicit namespace.
413    pub fn get_by_pubkey_in_namespace(
414        &self,
415        namespace: &str,
416        pubkey: &PubKey,
417    ) -> Option<InboxSender> {
418        if pubkey.is_zero() {
419            return None;
420        }
421        self.state
422            .read()
423            .namespace(namespace)?
424            .peers
425            .get(pubkey)
426            .map(|p| p.sender.clone())
427    }
428
429    /// Look up an inproc peer by pubkey across all namespaces.
430    ///
431    /// Cross-namespace delivery has no typed target namespace. If the same
432    /// canonical identity is live in more than one namespace, fail closed
433    /// rather than choosing whichever namespace the map happens to yield first.
434    pub(crate) fn get_by_pubkey_any_namespace(&self, pubkey: &PubKey) -> Option<InboxSender> {
435        if pubkey.is_zero() {
436            return None;
437        }
438        let state = self.state.read();
439        let mut found = None;
440        for namespace_state in state.namespaces.values() {
441            if let Some(peer) = namespace_state.peers.get(pubkey) {
442                if found.is_some() {
443                    return None;
444                }
445                found = Some(peer.sender.clone());
446            }
447        }
448        found
449    }
450
451    /// Look up an inproc peer name by public key.
452    pub fn get_name_by_pubkey(&self, pubkey: &PubKey) -> Option<String> {
453        self.get_name_by_pubkey_in_namespace(DEFAULT_NAMESPACE, pubkey)
454    }
455
456    /// Look up an inproc peer name by public key in an explicit namespace.
457    pub fn get_name_by_pubkey_in_namespace(
458        &self,
459        namespace: &str,
460        pubkey: &PubKey,
461    ) -> Option<String> {
462        if pubkey.is_zero() {
463            return None;
464        }
465        self.state
466            .read()
467            .namespace(namespace)?
468            .peers
469            .get(pubkey)
470            .map(|peer| peer.name.clone())
471    }
472
473    /// Check if a peer is registered.
474    pub fn contains(&self, pubkey: &PubKey) -> bool {
475        self.state
476            .read()
477            .namespace(DEFAULT_NAMESPACE)
478            .is_some_and(|ns| ns.peers.contains_key(pubkey))
479    }
480
481    /// Check if a peer name is registered.
482    pub fn contains_name(&self, name: &str) -> bool {
483        self.state
484            .read()
485            .namespace(DEFAULT_NAMESPACE)
486            .is_some_and(|ns| ns.names.contains_key(name))
487    }
488
489    /// Get the number of registered peers.
490    pub fn len(&self) -> usize {
491        self.state.read().namespace_len(DEFAULT_NAMESPACE)
492    }
493
494    /// Check if the registry is empty.
495    pub fn is_empty(&self) -> bool {
496        self.state.read().namespace_is_empty(DEFAULT_NAMESPACE)
497    }
498
499    /// Clear all registrations (primarily for testing).
500    pub fn clear(&self) {
501        self.state.write().namespaces.clear();
502    }
503
504    /// Backpressured pubkey-keyed delivery across all namespaces.
505    ///
506    /// Runtime-originated peer sends should await receiver capacity instead of
507    /// turning a transient full inbox into semantic message loss.
508    pub(crate) async fn send_to_pubkey_any_namespace_with_id_wait(
509        &self,
510        from_keypair: &Keypair,
511        to_pubkey: &PubKey,
512        envelope_id: Uuid,
513        kind: MessageKind,
514        sign_envelope: bool,
515    ) -> Result<uuid::Uuid, InprocSendError> {
516        let sender = self
517            .get_by_pubkey_any_namespace(to_pubkey)
518            .ok_or_else(|| InprocSendError::PeerNotFound(to_pubkey.to_peer_id().to_string()))?;
519
520        Self::deliver_to_sender_wait(
521            from_keypair,
522            *to_pubkey,
523            sender,
524            envelope_id,
525            kind,
526            sign_envelope,
527        )
528        .await
529    }
530
531    /// Namespace-scoped variant of
532    /// [`Self::send_to_pubkey_any_namespace_with_id_wait`]: the destination is
533    /// resolved exactly once, *inside* `namespace`, and that resolved sender is
534    /// the delivery target.
535    ///
536    /// This is the single-resolution send for namespace-isolated routers. The
537    /// namespace is the delivery authority, so the destination must not be
538    /// re-derived from the global registry between an isolation check and the
539    /// inbox handoff — a second any-namespace lookup would open a window where
540    /// the peer re-registers elsewhere and delivery crosses the namespace
541    /// boundary.
542    pub(crate) async fn send_to_pubkey_in_namespace_with_id_wait(
543        &self,
544        namespace: &str,
545        from_keypair: &Keypair,
546        to_pubkey: &PubKey,
547        envelope_id: Uuid,
548        kind: MessageKind,
549        sign_envelope: bool,
550    ) -> Result<uuid::Uuid, InprocSendError> {
551        let sender = self
552            .get_by_pubkey_in_namespace(namespace, to_pubkey)
553            .ok_or_else(|| InprocSendError::PeerNotFound(to_pubkey.to_peer_id().to_string()))?;
554
555        Self::deliver_to_sender_wait(
556            from_keypair,
557            *to_pubkey,
558            sender,
559            envelope_id,
560            kind,
561            sign_envelope,
562        )
563        .await
564    }
565
566    async fn deliver_to_sender_wait(
567        from_keypair: &Keypair,
568        to_pubkey: PubKey,
569        sender: InboxSender,
570        envelope_id: Uuid,
571        kind: MessageKind,
572        sign_envelope: bool,
573    ) -> Result<uuid::Uuid, InprocSendError> {
574        let mut envelope = Envelope {
575            id: envelope_id,
576            from: from_keypair.public_key(),
577            to: to_pubkey,
578            kind,
579            sig: Signature::new([0u8; 64]),
580        };
581        if sign_envelope {
582            envelope.sign(from_keypair);
583        }
584
585        let envelope_id = envelope.id;
586        match sender.send_wait(InboxItem::External { envelope }).await {
587            AdmissionOutcome::Admitted => {}
588            AdmissionOutcome::Dropped {
589                reason: DropReason::SessionClosed,
590            } => return Err(InprocSendError::InboxClosed),
591            AdmissionOutcome::Dropped {
592                reason: DropReason::InboxFull,
593            } => return Err(InprocSendError::InboxFull),
594            AdmissionOutcome::Dropped { reason } => {
595                return Err(InprocSendError::IngressDropped(reason));
596            }
597        }
598
599        Ok(envelope_id)
600    }
601
602    /// List all registered peer names in an explicit namespace.
603    pub fn peer_names_in_namespace(&self, namespace: &str) -> Vec<String> {
604        self.state
605            .read()
606            .namespace(namespace)
607            .map_or_else(Vec::new, |ns| ns.names.keys().cloned().collect())
608    }
609
610    /// List all registered peers.
611    pub fn peers(&self) -> Vec<InprocPeerInfo> {
612        self.peers_in_namespace(DEFAULT_NAMESPACE)
613    }
614
615    /// List all registered peers in an explicit namespace.
616    pub fn peers_in_namespace(&self, namespace: &str) -> Vec<InprocPeerInfo> {
617        self.state
618            .read()
619            .namespace(namespace)
620            .map_or_else(Vec::new, |ns| {
621                ns.peers
622                    .values()
623                    .map(|peer| InprocPeerInfo {
624                        name: peer.name.clone(),
625                        pubkey: peer.pubkey,
626                        meta: peer.meta.clone(),
627                    })
628                    .collect()
629            })
630    }
631}
632
633impl Default for InprocRegistry {
634    fn default() -> Self {
635        Self::new()
636    }
637}
638
639/// Errors that can occur during inproc send operations.
640#[derive(Debug, thiserror::Error)]
641pub enum InprocSendError {
642    #[error("Inproc peer not found: {0}")]
643    PeerNotFound(String),
644    #[error("Peer inbox has been closed")]
645    InboxClosed,
646    #[error("Peer inbox is full")]
647    InboxFull,
648    #[error("Peer inbox dropped ingress: {0:?}")]
649    IngressDropped(crate::inbox::DropReason),
650}
651
652#[cfg(test)]
653#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
654mod tests {
655    use super::*;
656    use crate::classify::test_support;
657    use crate::inbox::Inbox;
658    use crate::trust::TrustStore;
659
660    fn classified_inbox() -> (Inbox, crate::InboxSender) {
661        Inbox::new_classified(test_support::classification_context(
662            TrustStore::new(),
663            false,
664        ))
665    }
666
667    fn make_keypair() -> Keypair {
668        Keypair::generate()
669    }
670
671    #[test]
672    fn test_registry_new() {
673        let registry = InprocRegistry::new();
674        assert!(registry.is_empty());
675        assert_eq!(registry.len(), 0);
676    }
677
678    #[test]
679    fn test_registry_register_and_lookup() {
680        let registry = InprocRegistry::new();
681        let keypair = make_keypair();
682        let pubkey = keypair.public_key();
683        let (_, sender) = classified_inbox();
684
685        registry.register("test-agent", pubkey, sender);
686
687        assert!(!registry.is_empty());
688        assert_eq!(registry.len(), 1);
689        assert!(registry.contains(&pubkey));
690        assert!(registry.contains_name("test-agent"));
691
692        // Name is display metadata; routing lookups are pubkey-keyed.
693        assert_eq!(
694            registry.get_name_by_pubkey(&pubkey).as_deref(),
695            Some("test-agent")
696        );
697        assert!(registry.get_by_pubkey(&pubkey).is_some());
698    }
699
700    #[test]
701    fn test_registry_rejects_zero_pubkey_registration() {
702        let registry = InprocRegistry::new();
703        let (_, sender) = classified_inbox();
704        let zero_pubkey = PubKey::new([0u8; 32]);
705
706        registry.register("zero-agent", zero_pubkey, sender);
707
708        assert!(registry.is_empty());
709        assert!(!registry.contains_name("zero-agent"));
710        assert!(registry.get_by_pubkey(&zero_pubkey).is_none());
711    }
712
713    #[test]
714    fn test_registry_zero_pubkey_registration_does_not_shadow_valid_name() {
715        let registry = InprocRegistry::new();
716        let valid_keypair = make_keypair();
717        let valid_pubkey = valid_keypair.public_key();
718        let (_, valid_sender) = classified_inbox();
719        let (_, zero_sender) = classified_inbox();
720        let zero_pubkey = PubKey::new([0u8; 32]);
721
722        registry.register("stable-agent", valid_pubkey, valid_sender);
723        registry.register("stable-agent", zero_pubkey, zero_sender);
724
725        assert_eq!(registry.len(), 1);
726        assert!(registry.contains(&valid_pubkey));
727        assert!(registry.contains_name("stable-agent"));
728        assert!(registry.get_by_pubkey(&valid_pubkey).is_some());
729        assert!(registry.get_by_pubkey(&zero_pubkey).is_none());
730
731        assert_eq!(
732            registry.get_name_by_pubkey(&valid_pubkey).as_deref(),
733            Some("stable-agent"),
734            "valid name mapping should remain"
735        );
736    }
737
738    #[test]
739    fn test_registry_unregister() {
740        let registry = InprocRegistry::new();
741        let keypair = make_keypair();
742        let pubkey = keypair.public_key();
743        let (_, sender) = classified_inbox();
744
745        registry.register("test-agent", pubkey, sender);
746        assert!(registry.contains(&pubkey));
747
748        let removed = registry.unregister(&pubkey);
749        assert!(removed);
750        assert!(!registry.contains(&pubkey));
751        assert!(!registry.contains_name("test-agent"));
752        assert!(registry.is_empty());
753
754        // Unregister non-existent returns false
755        let removed_again = registry.unregister(&pubkey);
756        assert!(!removed_again);
757    }
758
759    #[test]
760    fn test_registry_replace_on_same_pubkey() {
761        let registry = InprocRegistry::new();
762        let keypair = make_keypair();
763        let pubkey = keypair.public_key();
764        let (_, sender1) = classified_inbox();
765        let (_, sender2) = classified_inbox();
766
767        // Register with first name
768        registry.register("agent-v1", pubkey, sender1);
769        assert!(registry.contains_name("agent-v1"));
770
771        // Re-register same pubkey with different name
772        registry.register("agent-v2", pubkey, sender2);
773
774        // Old name should be removed, new name should exist
775        assert!(!registry.contains_name("agent-v1"));
776        assert!(registry.contains_name("agent-v2"));
777        assert_eq!(registry.len(), 1);
778    }
779
780    #[test]
781    fn test_registry_replace_on_same_name_different_pubkey() {
782        let registry = InprocRegistry::new();
783        let keypair1 = make_keypair();
784        let pubkey1 = keypair1.public_key();
785        let keypair2 = make_keypair();
786        let pubkey2 = keypair2.public_key();
787        let (_, sender1) = classified_inbox();
788        let (_, sender2) = classified_inbox();
789
790        // Register first agent
791        registry.register("my-agent", pubkey1, sender1);
792        assert!(registry.contains(&pubkey1));
793        assert!(registry.contains_name("my-agent"));
794        assert_eq!(registry.len(), 1);
795
796        // Re-register same name with different pubkey
797        registry.register("my-agent", pubkey2, sender2);
798
799        // Old pubkey should be evicted, new pubkey should exist
800        assert!(!registry.contains(&pubkey1), "old pubkey should be evicted");
801        assert!(registry.contains(&pubkey2));
802        assert!(registry.contains_name("my-agent"));
803        assert_eq!(registry.len(), 1);
804
805        // The name maps to the new identity.
806        assert_eq!(
807            registry.get_name_by_pubkey(&pubkey2).as_deref(),
808            Some("my-agent")
809        );
810    }
811
812    /// ROW #292 gate: registration returns a typed [`RegistrationOutcome`] that
813    /// surfaces route displacement and zero-pubkey rejection, instead of
814    /// silently evicting and returning `()`.
815    #[test]
816    fn registration_outcome_is_typed_for_displacement_and_rejection() {
817        let registry = InprocRegistry::new();
818        let keypair = make_keypair();
819        let pubkey = keypair.public_key();
820        let (_, sender1) = classified_inbox();
821        let (_, sender2) = classified_inbox();
822        let (_, sender3) = classified_inbox();
823
824        // Clean first insert.
825        assert_eq!(
826            registry.register("agent-v1", pubkey, sender1),
827            RegistrationOutcome::Registered
828        );
829
830        // Re-registering the SAME pubkey under a NEW name evicts the old name
831        // and reports it typed.
832        assert_eq!(
833            registry.register("agent-v2", pubkey, sender2),
834            RegistrationOutcome::ReplacedPubkey {
835                evicted_name: "agent-v1".to_string()
836            }
837        );
838
839        // Re-registering an existing NAME with a NEW pubkey evicts the old
840        // pubkey and reports it typed.
841        let other = make_keypair();
842        let other_pubkey = other.public_key();
843        match registry.register("agent-v2", other_pubkey, sender3) {
844            RegistrationOutcome::EvictedName { evicted_pubkey } => {
845                assert_eq!(evicted_pubkey, pubkey);
846            }
847            other => panic!("expected EvictedName, got {other:?}"),
848        }
849
850        // A zero pubkey is refused fail-closed with a typed rejection, no
851        // mutation.
852        let (_, zero_sender) = classified_inbox();
853        let zero_pubkey = PubKey::new([0u8; 32]);
854        assert_eq!(
855            registry.register("zero", zero_pubkey, zero_sender),
856            RegistrationOutcome::Rejected {
857                reason: RegistrationRejection::ZeroPubkey
858            }
859        );
860        assert!(!registry.contains_name("zero"));
861    }
862
863    /// Test that the ABA scenario is handled correctly:
864    /// When a new agent registers with the same name, the old agent's
865    /// unregister call (on Drop) should be a safe no-op.
866    #[test]
867    fn test_registry_aba_scenario_safe() {
868        let registry = InprocRegistry::new();
869        let keypair_old = make_keypair();
870        let pubkey_old = keypair_old.public_key();
871        let keypair_new = make_keypair();
872        let pubkey_new = keypair_new.public_key();
873        let (_, sender_old) = classified_inbox();
874        let (_, sender_new) = classified_inbox();
875
876        // Step 1: Old runtime registers
877        registry.register("agent", pubkey_old, sender_old);
878        assert!(registry.contains(&pubkey_old));
879
880        // Step 2: New runtime registers same name (evicts old)
881        registry.register("agent", pubkey_new, sender_new);
882        assert!(
883            !registry.contains(&pubkey_old),
884            "old pubkey should be evicted"
885        );
886        assert!(registry.contains(&pubkey_new));
887
888        // Step 3: Old runtime drops and calls unregister(pubkey_old)
889        // This should be a no-op since pubkey_old was already evicted
890        let removed = registry.unregister(&pubkey_old);
891        assert!(!removed, "unregister of evicted pubkey should return false");
892
893        // New agent should still be registered (not affected by old unregister)
894        assert!(
895            registry.contains(&pubkey_new),
896            "new agent should still be registered"
897        );
898        assert!(
899            registry.contains_name("agent"),
900            "name should still map to new agent"
901        );
902
903        // The name maps to the new identity.
904        assert_eq!(
905            registry.get_name_by_pubkey(&pubkey_new).as_deref(),
906            Some("agent"),
907            "name should map to the new pubkey"
908        );
909    }
910
911    #[test]
912    fn test_registry_peer_names_in_namespace() {
913        let registry = InprocRegistry::new();
914
915        for i in 0..3 {
916            let keypair = make_keypair();
917            let (_, sender) = classified_inbox();
918            registry.register_with_meta_in_namespace(
919                "realm-names",
920                format!("agent-{i}"),
921                keypair.public_key(),
922                sender,
923                PeerMeta::default(),
924            );
925        }
926
927        let names = registry.peer_names_in_namespace("realm-names");
928        assert_eq!(names.len(), 3);
929        assert!(names.contains(&"agent-0".to_string()));
930        assert!(names.contains(&"agent-1".to_string()));
931        assert!(names.contains(&"agent-2".to_string()));
932        assert!(registry.peer_names_in_namespace("realm-other").is_empty());
933    }
934
935    #[test]
936    fn test_registry_peers_snapshot() {
937        let registry = InprocRegistry::new();
938        let keypair = make_keypair();
939        let pubkey = keypair.public_key();
940        let (_, sender) = classified_inbox();
941        registry.register("agent-a", pubkey, sender);
942
943        let peers = registry.peers();
944        assert_eq!(peers.len(), 1);
945        assert_eq!(peers[0].name, "agent-a");
946        assert_eq!(peers[0].pubkey, pubkey);
947    }
948
949    #[test]
950    fn test_registry_clear() {
951        let registry = InprocRegistry::new();
952
953        for i in 0..3 {
954            let keypair = make_keypair();
955            let (_, sender) = classified_inbox();
956            registry.register(format!("agent-{i}"), keypair.public_key(), sender);
957        }
958
959        assert_eq!(registry.len(), 3);
960        registry.clear();
961        assert!(registry.is_empty());
962    }
963
964    #[tokio::test]
965    async fn test_registry_send_delivers_to_inbox() {
966        let registry = InprocRegistry::new();
967
968        // Set up receiver
969        let receiver_keypair = make_keypair();
970        let (mut inbox, sender) = classified_inbox();
971        registry.register("receiver", receiver_keypair.public_key(), sender);
972
973        // Set up sender
974        let sender_keypair = make_keypair();
975
976        // Send a message (pubkey-keyed delivery)
977        let result = registry
978            .send_to_pubkey_in_namespace_with_id_wait(
979                "",
980                &sender_keypair,
981                &receiver_keypair.public_key(),
982                Uuid::new_v4(),
983                MessageKind::Message {
984                    objective_id: None,
985                    content_taint: None,
986                    blocks: None,
987                    body: "hello inproc".to_string(),
988                    handling_mode: None,
989                },
990                true,
991            )
992            .await;
993        assert!(result.is_ok());
994
995        // Verify message was received
996        let items = inbox.try_drain_classified();
997        assert_eq!(items.len(), 1);
998
999        match &items[0].item {
1000            InboxItem::External { envelope } => {
1001                assert_eq!(envelope.from, sender_keypair.public_key());
1002                assert_eq!(envelope.to, receiver_keypair.public_key());
1003                match &envelope.kind {
1004                    MessageKind::Message {
1005                        blocks: None, body, ..
1006                    } => {
1007                        assert_eq!(body, "hello inproc");
1008                    }
1009                    _ => panic!("expected Message kind"),
1010                }
1011                // Verify signature
1012                assert!(envelope.verify());
1013            }
1014            _ => panic!("expected External inbox item"),
1015        }
1016    }
1017
1018    #[tokio::test]
1019    async fn test_registry_send_peer_not_found() {
1020        let registry = InprocRegistry::new();
1021        let sender_keypair = make_keypair();
1022        let unknown = make_keypair().public_key();
1023
1024        let result = registry
1025            .send_to_pubkey_in_namespace_with_id_wait(
1026                "",
1027                &sender_keypair,
1028                &unknown,
1029                Uuid::new_v4(),
1030                MessageKind::Message {
1031                    objective_id: None,
1032                    content_taint: None,
1033                    blocks: None,
1034                    body: "hello".to_string(),
1035                    handling_mode: None,
1036                },
1037                true,
1038            )
1039            .await;
1040
1041        assert!(matches!(result, Err(InprocSendError::PeerNotFound(_))));
1042    }
1043
1044    #[tokio::test]
1045    async fn test_registry_send_inbox_closed() {
1046        let registry = InprocRegistry::new();
1047
1048        // Set up receiver but drop the inbox
1049        let receiver_keypair = make_keypair();
1050        let (inbox, sender) = classified_inbox();
1051        registry.register("receiver", receiver_keypair.public_key(), sender);
1052        drop(inbox); // Close the inbox
1053
1054        let sender_keypair = make_keypair();
1055
1056        let result = registry
1057            .send_to_pubkey_in_namespace_with_id_wait(
1058                "",
1059                &sender_keypair,
1060                &receiver_keypair.public_key(),
1061                Uuid::new_v4(),
1062                MessageKind::Message {
1063                    objective_id: None,
1064                    content_taint: None,
1065                    blocks: None,
1066                    body: "hello".to_string(),
1067                    handling_mode: None,
1068                },
1069                true,
1070            )
1071            .await;
1072
1073        assert!(matches!(result, Err(InprocSendError::InboxClosed)));
1074    }
1075
1076    #[tokio::test]
1077    async fn test_registry_namespace_isolation_for_lookup_and_send() {
1078        let registry = InprocRegistry::new();
1079        let receiver_keypair = make_keypair();
1080        let (mut inbox, sender) = classified_inbox();
1081        registry.register_with_meta_in_namespace(
1082            "realm-a",
1083            "receiver",
1084            receiver_keypair.public_key(),
1085            sender,
1086            PeerMeta::default(),
1087        );
1088
1089        // Default namespace cannot see realm-a registrations.
1090        assert!(
1091            registry
1092                .get_by_pubkey(&receiver_keypair.public_key())
1093                .is_none()
1094        );
1095        assert!(
1096            registry
1097                .get_by_pubkey_in_namespace("realm-a", &receiver_keypair.public_key())
1098                .is_some()
1099        );
1100
1101        let sender_keypair = make_keypair();
1102
1103        // Matching namespace succeeds.
1104        let ok = registry
1105            .send_to_pubkey_in_namespace_with_id_wait(
1106                "realm-a",
1107                &sender_keypair,
1108                &receiver_keypair.public_key(),
1109                Uuid::new_v4(),
1110                MessageKind::Message {
1111                    objective_id: None,
1112                    content_taint: None,
1113                    blocks: None,
1114                    body: "hello scoped".to_string(),
1115                    handling_mode: None,
1116                },
1117                true,
1118            )
1119            .await;
1120        assert!(ok.is_ok());
1121
1122        // Different namespace cannot route to receiver.
1123        let wrong_ns = registry
1124            .send_to_pubkey_in_namespace_with_id_wait(
1125                "realm-b",
1126                &sender_keypair,
1127                &receiver_keypair.public_key(),
1128                Uuid::new_v4(),
1129                MessageKind::Message {
1130                    objective_id: None,
1131                    content_taint: None,
1132                    blocks: None,
1133                    body: "should not deliver".to_string(),
1134                    handling_mode: None,
1135                },
1136                true,
1137            )
1138            .await;
1139        assert!(matches!(wrong_ns, Err(InprocSendError::PeerNotFound(_))));
1140
1141        let items = inbox.try_drain_classified();
1142        assert_eq!(items.len(), 1);
1143    }
1144
1145    #[tokio::test]
1146    async fn test_send_to_pubkey_in_namespace_ignores_display_name_collision() {
1147        let registry = InprocRegistry::new();
1148        let target_keypair = make_keypair();
1149        let target_pubkey = target_keypair.public_key();
1150        let shadow_keypair = make_keypair();
1151        let shadow_pubkey = shadow_keypair.public_key();
1152        let (mut target_inbox, target_sender) = classified_inbox();
1153        let (mut shadow_inbox, shadow_sender) = classified_inbox();
1154
1155        registry.register_with_meta_in_namespace(
1156            "",
1157            "canonical-target",
1158            target_pubkey,
1159            target_sender,
1160            PeerMeta::default(),
1161        );
1162        registry.register_with_meta_in_namespace(
1163            "",
1164            "shared-display-name",
1165            shadow_pubkey,
1166            shadow_sender,
1167            PeerMeta::default(),
1168        );
1169
1170        let sender_keypair = make_keypair();
1171        let result = registry
1172            .send_to_pubkey_in_namespace_with_id_wait(
1173                "",
1174                &sender_keypair,
1175                &target_pubkey,
1176                Uuid::new_v4(),
1177                MessageKind::Message {
1178                    objective_id: None,
1179                    content_taint: None,
1180                    blocks: None,
1181                    body: "hello canonical".to_string(),
1182                    handling_mode: None,
1183                },
1184                true,
1185            )
1186            .await;
1187        assert!(result.is_ok());
1188
1189        assert_eq!(shadow_inbox.try_drain_classified().len(), 0);
1190        let items = target_inbox.try_drain_classified();
1191        assert_eq!(items.len(), 1);
1192        let InboxItem::External { envelope } = &items[0].item else {
1193            panic!("expected external envelope");
1194        };
1195        assert_eq!(envelope.to, target_pubkey);
1196    }
1197
1198    #[tokio::test]
1199    async fn test_send_to_pubkey_any_namespace_rejects_ambiguous_identity() {
1200        let registry = InprocRegistry::new();
1201        let sender_keypair = make_keypair();
1202        let target_keypair = make_keypair();
1203        let target_pubkey = target_keypair.public_key();
1204        let (mut alpha_inbox, alpha_sender) = classified_inbox();
1205        let (mut beta_inbox, beta_sender) = classified_inbox();
1206
1207        registry.register_with_meta_in_namespace(
1208            "realm-alpha",
1209            "alpha-target",
1210            target_pubkey,
1211            alpha_sender,
1212            PeerMeta::default(),
1213        );
1214        registry.register_with_meta_in_namespace(
1215            "realm-beta",
1216            "beta-target",
1217            target_pubkey,
1218            beta_sender,
1219            PeerMeta::default(),
1220        );
1221
1222        let result = registry
1223            .send_to_pubkey_any_namespace_with_id_wait(
1224                &sender_keypair,
1225                &target_pubkey,
1226                Uuid::new_v4(),
1227                MessageKind::Message {
1228                    objective_id: None,
1229                    content_taint: None,
1230                    blocks: None,
1231                    body: "ambiguous identity".to_string(),
1232                    handling_mode: None,
1233                },
1234                true,
1235            )
1236            .await;
1237
1238        assert!(matches!(result, Err(InprocSendError::PeerNotFound(_))));
1239        assert!(alpha_inbox.try_drain_classified().is_empty());
1240        assert!(beta_inbox.try_drain_classified().is_empty());
1241    }
1242
1243    #[test]
1244    fn test_registry_same_name_can_exist_in_different_namespaces() {
1245        let registry = InprocRegistry::new();
1246        let kp_a = make_keypair();
1247        let kp_b = make_keypair();
1248        let (_, sender_a) = classified_inbox();
1249        let (_, sender_b) = classified_inbox();
1250
1251        registry.register_with_meta_in_namespace(
1252            "realm-a",
1253            "shared-name",
1254            kp_a.public_key(),
1255            sender_a,
1256            PeerMeta::default(),
1257        );
1258        registry.register_with_meta_in_namespace(
1259            "realm-b",
1260            "shared-name",
1261            kp_b.public_key(),
1262            sender_b,
1263            PeerMeta::default(),
1264        );
1265
1266        assert_eq!(
1267            registry
1268                .get_name_by_pubkey_in_namespace("realm-a", &kp_a.public_key())
1269                .as_deref(),
1270            Some("shared-name")
1271        );
1272        assert_eq!(
1273            registry
1274                .get_name_by_pubkey_in_namespace("realm-b", &kp_b.public_key())
1275                .as_deref(),
1276            Some("shared-name")
1277        );
1278        assert_ne!(kp_a.public_key(), kp_b.public_key());
1279        assert!(
1280            !registry.contains_name("shared-name"),
1281            "default namespace must not see namespaced registrations"
1282        );
1283    }
1284
1285    #[test]
1286    fn test_global_registry() {
1287        // Access global registry
1288        let registry = InprocRegistry::global();
1289
1290        // Clear any existing state (from other tests)
1291        registry.clear();
1292
1293        // Register a peer
1294        let keypair = make_keypair();
1295        let (_, sender) = classified_inbox();
1296        registry.register("global-test", keypair.public_key(), sender);
1297
1298        // Verify it's accessible
1299        assert!(registry.contains_name("global-test"));
1300
1301        // Clean up
1302        registry.unregister(&keypair.public_key());
1303    }
1304
1305    #[test]
1306    fn test_registry_register_with_meta() {
1307        let registry = InprocRegistry::new();
1308        let keypair = make_keypair();
1309        let pubkey = keypair.public_key();
1310        let (_, sender) = classified_inbox();
1311
1312        let meta = PeerMeta::default()
1313            .with_description("Reviews code for style issues")
1314            .with_label("lang", "rust");
1315
1316        registry.register_with_meta_in_namespace(
1317            DEFAULT_NAMESPACE,
1318            "reviewer",
1319            pubkey,
1320            sender,
1321            meta.clone(),
1322        );
1323
1324        let peers = registry.peers();
1325        assert_eq!(peers.len(), 1);
1326        assert_eq!(peers[0].name, "reviewer");
1327        assert_eq!(peers[0].pubkey, pubkey);
1328        assert_eq!(peers[0].meta, meta);
1329    }
1330
1331    #[test]
1332    fn test_registry_peers_returns_default_meta_for_plain_register() {
1333        let registry = InprocRegistry::new();
1334        let keypair = make_keypair();
1335        let pubkey = keypair.public_key();
1336        let (_, sender) = classified_inbox();
1337
1338        registry.register("plain-agent", pubkey, sender);
1339
1340        let peers = registry.peers();
1341        assert_eq!(peers.len(), 1);
1342        assert_eq!(peers[0].meta, PeerMeta::default());
1343    }
1344}