#[cfg(not(target_arch = "wasm32"))]
use super::comms_config::ResolvedCommsConfig;
#[cfg(not(target_arch = "wasm32"))]
use crate::InboxSender;
use crate::agent::types::CommsMessage;
#[cfg(not(target_arch = "wasm32"))]
use crate::handle_connection;
#[cfg(target_arch = "wasm32")]
use crate::tokio;
use crate::{InprocRegistry, Keypair, PubKey, Router, TrustedPeer, TrustedPeers};
use async_trait::async_trait;
use futures::Stream;
use futures::task::{Context, Poll};
use meerkat_core::agent::CommsRuntime as CoreCommsRuntime;
use meerkat_core::comms::{
CommsCommand, EventStream, InputStreamMode, PeerDirectoryEntry, PeerDirectorySource, PeerName,
SendAndStreamError, SendError, SendReceipt, StreamError, StreamScope, TrustedPeerSpec,
};
use meerkat_core::config::PlainEventSource;
use meerkat_core::time_compat::Instant;
use parking_lot::Mutex;
use std::collections::{HashMap, HashSet};
#[cfg(unix)]
use std::path::Path;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use thiserror::Error;
#[cfg(not(target_arch = "wasm32"))]
use tokio::net::TcpListener;
use tokio::sync::Mutex as AsyncMutex;
use tokio::sync::mpsc;
use tokio::sync::mpsc::Receiver;
#[cfg(not(target_arch = "wasm32"))]
use tokio::task::JoinHandle;
use uuid::Uuid;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ReservationState {
Reserved,
Attached,
Completed,
Expired,
ClosedEarly,
}
impl ReservationState {
fn is_terminal(self) -> bool {
matches!(self, Self::Completed | Self::Expired | Self::ClosedEarly)
}
}
const RESERVATION_TTL: meerkat_core::time_compat::Duration =
meerkat_core::time_compat::Duration::from_secs(30);
#[derive(Debug)]
struct StreamRegistryEntry {
state: ReservationState,
_sender: mpsc::Sender<meerkat_core::AgentEvent>,
receiver: Option<Receiver<meerkat_core::AgentEvent>>,
created_at: Instant,
}
impl StreamRegistryEntry {
fn reserved(
sender: mpsc::Sender<meerkat_core::AgentEvent>,
receiver: Receiver<meerkat_core::AgentEvent>,
) -> Self {
Self {
state: ReservationState::Reserved,
_sender: sender,
receiver: Some(receiver),
created_at: Instant::now(),
}
}
fn transition(&mut self, from: ReservationState, to: ReservationState) -> bool {
if self.state == from {
self.state = to;
true
} else {
false
}
}
}
type InteractionStreamRegistry = Arc<Mutex<HashMap<Uuid, StreamRegistryEntry>>>;
struct InteractionStream {
id: Uuid,
receiver: Option<Receiver<meerkat_core::AgentEvent>>,
registry: InteractionStreamRegistry,
source_id: String,
seq: u64,
}
struct ResolvedPeer {
name: PeerName,
peer_id: String,
address: String,
source: PeerDirectorySource,
meta: crate::PeerMeta,
}
impl InteractionStream {
fn finish(&mut self) {
if let Some(mut receiver) = self.receiver.take() {
receiver.close();
}
let mut registry = self.registry.lock();
if let Some(entry) = registry.get_mut(&self.id) {
if entry.transition(ReservationState::Attached, ReservationState::ClosedEarly) {
tracing::debug!(interaction_id = %self.id, "stream closed early by consumer");
}
if entry.state.is_terminal() {
registry.remove(&self.id);
}
}
}
}
impl Stream for InteractionStream {
type Item = meerkat_core::EventEnvelope<meerkat_core::AgentEvent>;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let this = self.get_mut();
match this.receiver.as_mut() {
Some(receiver) => match Pin::new(receiver).poll_recv(cx) {
Poll::Ready(None) => {
this.finish();
Poll::Ready(None)
}
Poll::Ready(Some(event)) => {
this.seq = this.seq.saturating_add(1);
let envelope =
meerkat_core::EventEnvelope::new(&this.source_id, this.seq, None, event);
Poll::Ready(Some(envelope))
}
Poll::Pending => Poll::Pending,
},
None => Poll::Ready(None),
}
}
}
impl Drop for InteractionStream {
fn drop(&mut self) {
self.finish();
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl CoreCommsRuntime for CommsRuntime {
async fn drain_messages(&self) -> Vec<String> {
self.drain_inbox_interactions()
.await
.into_iter()
.map(|i| i.rendered_text)
.collect()
}
fn inbox_notify(&self) -> Arc<tokio::sync::Notify> {
self.inbox_notify.clone()
}
fn public_key(&self) -> Option<String> {
Some(self.public_key.to_peer_id())
}
async fn add_trusted_peer(&self, peer: TrustedPeerSpec) -> Result<(), SendError> {
let public_key = PubKey::from_peer_id(&peer.peer_id)
.map_err(|err| SendError::Validation(err.to_string()))?;
let trusted_peer = TrustedPeer {
name: peer.name,
pubkey: public_key,
addr: peer.address,
meta: crate::PeerMeta::default(),
};
self.router.add_trusted_peer(trusted_peer);
Ok(())
}
async fn remove_trusted_peer(&self, peer_id: &str) -> Result<bool, SendError> {
let public_key =
PubKey::from_peer_id(peer_id).map_err(|err| SendError::Validation(err.to_string()))?;
let removed = self.router.remove_trusted_peer(&public_key);
Ok(removed)
}
fn dismiss_received(&self) -> bool {
self.dismiss_flag.swap(false, Ordering::SeqCst)
}
fn event_injector(&self) -> Option<Arc<dyn meerkat_core::EventInjector>> {
Some(self.event_injector())
}
fn interaction_event_injector(
&self,
) -> Option<Arc<dyn meerkat_core::event_injector::SubscribableInjector>> {
Some(self.interaction_event_injector())
}
fn stream(&self, scope: StreamScope) -> Result<EventStream, StreamError> {
match scope {
StreamScope::Session(session_id) => {
Err(StreamError::NotFound(format!("session {session_id}")))
}
StreamScope::Interaction(interaction_id) => {
let id = interaction_id.0;
let mut registry = self.interaction_stream_registry.lock();
let entry = registry
.get_mut(&id)
.ok_or(StreamError::NotReserved(interaction_id))?;
match entry.state {
ReservationState::Reserved => {
if entry.created_at.elapsed() > RESERVATION_TTL {
entry.state = ReservationState::Expired;
registry.remove(&id);
return Err(StreamError::Timeout(format!(
"reservation expired for interaction {}",
interaction_id.0
)));
}
let receiver = entry.receiver.take().ok_or_else(|| {
StreamError::Internal("interaction stream receiver missing".to_string())
})?;
entry.state = ReservationState::Attached;
Ok(Box::pin(InteractionStream {
id,
receiver: Some(receiver),
registry: self.interaction_stream_registry.clone(),
source_id: format!("interaction:{id}"),
seq: 0,
}))
}
ReservationState::Attached => Err(StreamError::AlreadyAttached(interaction_id)),
state if state.is_terminal() => Err(StreamError::NotReserved(interaction_id)),
_ => Err(StreamError::Internal(format!(
"unexpected reservation state for {}",
interaction_id.0
))),
}
}
}
}
async fn send(&self, cmd: CommsCommand) -> Result<SendReceipt, SendError> {
match cmd {
CommsCommand::Input {
body,
source,
stream,
allow_self_session,
session_id: _,
} => {
if !allow_self_session {
return Err(SendError::Validation(
"self-session input rejected: set allow_self_session=true to override"
.into(),
));
}
match stream {
InputStreamMode::None => {
let injector = CoreCommsRuntime::event_injector(self).ok_or_else(|| {
SendError::Unsupported("event injector unavailable".into())
})?;
injector
.inject(body, PlainEventSource::from(source))
.map_err(map_event_injector_error)?;
Ok(SendReceipt::InputAccepted {
interaction_id: meerkat_core::InteractionId(Uuid::new_v4()),
stream_reserved: false,
})
}
InputStreamMode::ReserveInteraction => {
let interaction_id = Uuid::new_v4();
self.register_interaction_stream(interaction_id, body, source)?;
Ok(SendReceipt::InputAccepted {
interaction_id: meerkat_core::InteractionId(interaction_id),
stream_reserved: true,
})
}
}
}
CommsCommand::PeerMessage { to, body } => self
.send_peer_command(to.as_str(), crate::types::MessageKind::Message { body })
.await
.map(|envelope_id| SendReceipt::PeerMessageSent {
envelope_id,
acked: false,
}),
CommsCommand::PeerRequest {
to,
intent,
params,
stream,
} => {
let interaction_id = Uuid::new_v4();
let stream_reserved = stream == InputStreamMode::ReserveInteraction;
if stream_reserved {
let (sender, receiver) = mpsc::channel::<meerkat_core::AgentEvent>(4096);
self.subscriber_registry
.lock()
.insert(interaction_id, sender.clone());
self.interaction_stream_registry.lock().insert(
interaction_id,
StreamRegistryEntry::reserved(sender, receiver),
);
}
let envelope_id = match self
.send_peer_command(
to.as_str(),
crate::types::MessageKind::Request { intent, params },
)
.await
{
Ok(id) => id,
Err(e) => {
if stream_reserved {
self.interaction_stream_registry
.lock()
.remove(&interaction_id);
self.subscriber_registry.lock().remove(&interaction_id);
}
return Err(e);
}
};
Ok(SendReceipt::PeerRequestSent {
envelope_id,
interaction_id: meerkat_core::InteractionId(interaction_id),
stream_reserved,
})
}
CommsCommand::PeerResponse {
to,
in_reply_to,
status,
result,
} => {
let status = match status {
meerkat_core::ResponseStatus::Accepted => crate::Status::Accepted,
meerkat_core::ResponseStatus::Completed => crate::Status::Completed,
meerkat_core::ResponseStatus::Failed => crate::Status::Failed,
};
self.send_peer_command(
to.as_str(),
crate::types::MessageKind::Response {
in_reply_to: in_reply_to.0,
status,
result,
},
)
.await
.map(|envelope_id| SendReceipt::PeerResponseSent {
envelope_id,
in_reply_to,
})
}
}
}
async fn send_and_stream(
&self,
cmd: CommsCommand,
) -> Result<(SendReceipt, EventStream), SendAndStreamError> {
match cmd {
CommsCommand::Input {
body,
source,
stream: InputStreamMode::ReserveInteraction,
allow_self_session,
session_id: _,
} => {
if !allow_self_session {
return Err(SendAndStreamError::Send(SendError::Validation(
"self-session input rejected: set allow_self_session=true to override"
.into(),
)));
}
let interaction_id = Uuid::new_v4();
self.register_interaction_stream(interaction_id, body, source)?;
let receipt = SendReceipt::InputAccepted {
interaction_id: meerkat_core::InteractionId(interaction_id),
stream_reserved: true,
};
let stream = self
.stream(StreamScope::Interaction(meerkat_core::InteractionId(
interaction_id,
)))
.map_err(|error| SendAndStreamError::StreamAttach {
receipt: receipt.clone(),
error,
})?;
Ok((receipt, stream))
}
CommsCommand::PeerRequest {
to,
intent,
params,
stream: InputStreamMode::ReserveInteraction,
} => {
let interaction_id = Uuid::new_v4();
let (sender, receiver) = mpsc::channel::<meerkat_core::AgentEvent>(4096);
self.subscriber_registry
.lock()
.insert(interaction_id, sender.clone());
self.interaction_stream_registry.lock().insert(
interaction_id,
StreamRegistryEntry::reserved(sender, receiver),
);
let envelope_id = match self
.send_peer_command(
to.as_str(),
crate::types::MessageKind::Request { intent, params },
)
.await
{
Ok(id) => id,
Err(e) => {
self.interaction_stream_registry
.lock()
.remove(&interaction_id);
self.subscriber_registry.lock().remove(&interaction_id);
return Err(SendAndStreamError::Send(e));
}
};
let receipt = SendReceipt::PeerRequestSent {
envelope_id,
interaction_id: meerkat_core::InteractionId(interaction_id),
stream_reserved: true,
};
let event_stream = self
.stream(StreamScope::Interaction(meerkat_core::InteractionId(
interaction_id,
)))
.map_err(|error| SendAndStreamError::StreamAttach {
receipt: receipt.clone(),
error,
})?;
Ok((receipt, event_stream))
}
other => {
let receipt = self.send(other).await?;
Err(SendAndStreamError::StreamAttach {
receipt,
error: StreamError::NotFound("command is not streamable".to_string()),
})
}
}
}
async fn peers(&self) -> Vec<PeerDirectoryEntry> {
self.resolve_peer_directory().await
}
async fn peer_count(&self) -> usize {
self.resolve_peer_count().await
}
async fn drain_inbox_interactions(&self) -> Vec<meerkat_core::InboxInteraction> {
self.drain_classified_inbox_interactions()
.await
.unwrap_or_default()
.into_iter()
.map(|ci| ci.interaction)
.collect()
}
fn interaction_subscriber(
&self,
id: &meerkat_core::InteractionId,
) -> Option<tokio::sync::mpsc::Sender<meerkat_core::AgentEvent>> {
let sender = self.subscriber_registry.lock().remove(&id.0);
sender.as_ref()?;
let mut registry = self.interaction_stream_registry.lock();
match registry.get(&id.0) {
Some(entry) => {
if entry.state == ReservationState::Reserved {
registry.remove(&id.0);
}
sender
}
None => sender,
}
}
fn mark_interaction_complete(&self, id: &meerkat_core::InteractionId) {
self.mark_interaction_complete(id.0);
}
async fn drain_classified_inbox_interactions(
&self,
) -> Result<Vec<meerkat_core::ClassifiedInboxInteraction>, meerkat_core::CommsCapabilityError>
{
use crate::agent::types::MessageIntent;
use crate::types::MessageKind;
let mut inbox = self.inbox.lock().await;
let classified_entries = inbox.try_drain_classified();
if classified_entries.is_empty() {
return Ok(Vec::new());
}
Ok(classified_entries
.into_iter()
.filter_map(|entry| {
let from_peer = entry.from_peer.unwrap_or_else(|| "unknown".to_string());
match entry.item {
crate::types::InboxItem::External { envelope } => {
if let MessageKind::Message { ref body } = envelope.kind
&& body.trim().eq_ignore_ascii_case("DISMISS") {
self.dismiss_flag.store(true, Ordering::SeqCst);
return None;
}
let (content, rendered_text) = match envelope.kind {
MessageKind::Message { body } => {
let rendered = format!(
"[COMMS MESSAGE from {from_peer}]\n{body}"
);
(
meerkat_core::InteractionContent::Message { body },
rendered,
)
}
MessageKind::Request { intent, params } => {
let typed_intent = MessageIntent::from(intent.as_str());
let params_str = if params.is_null()
|| params == serde_json::Value::Object(Default::default())
{
String::new()
} else {
format!(
"\nParams: {}",
serde_json::to_string_pretty(¶ms).unwrap_or_default()
)
};
let rendered = format!(
"[COMMS REQUEST from {} (id: {})]\n\
Intent: {}{}\n\
\n\
To respond, use send_response with peer=\"{}\", request_id=\"{}\"",
from_peer, envelope.id, typed_intent, params_str,
from_peer, envelope.id
);
(
meerkat_core::InteractionContent::Request {
intent: typed_intent.to_string(),
params,
},
rendered,
)
}
MessageKind::Response {
in_reply_to,
status,
result,
} => {
let core_status = match status {
crate::types::Status::Accepted => {
meerkat_core::ResponseStatus::Accepted
}
crate::types::Status::Completed => {
meerkat_core::ResponseStatus::Completed
}
crate::types::Status::Failed => {
meerkat_core::ResponseStatus::Failed
}
};
let status_str = match status {
crate::types::Status::Accepted => "accepted",
crate::types::Status::Completed => "completed",
crate::types::Status::Failed => "failed",
};
let result_str = if result.is_null()
|| result == serde_json::Value::Object(Default::default())
{
String::new()
} else {
format!(
"\nResult: {}",
serde_json::to_string_pretty(&result).unwrap_or_default()
)
};
let rendered = format!(
"[COMMS RESPONSE from {from_peer} (to request: {in_reply_to})]\n\
Status: {status_str}{result_str}"
);
(
meerkat_core::InteractionContent::Response {
in_reply_to: meerkat_core::InteractionId(in_reply_to),
status: core_status,
result,
},
rendered,
)
}
MessageKind::Ack { .. } => {
return None;
}
};
Some(meerkat_core::ClassifiedInboxInteraction {
interaction: meerkat_core::InboxInteraction {
id: meerkat_core::InteractionId(envelope.id),
from: from_peer,
content,
rendered_text,
},
class: entry.class,
lifecycle_peer: entry.lifecycle_peer,
})
}
crate::types::InboxItem::PlainEvent {
body,
source,
interaction_id,
} => {
let rendered = format!("[EVENT via {source}] {body}");
Some(meerkat_core::ClassifiedInboxInteraction {
interaction: meerkat_core::InboxInteraction {
id: meerkat_core::InteractionId(
interaction_id.unwrap_or_else(uuid::Uuid::new_v4),
),
from: format!("event:{source}"),
content: meerkat_core::InteractionContent::Message { body },
rendered_text: rendered,
},
class: entry.class,
lifecycle_peer: entry.lifecycle_peer,
})
}
crate::types::InboxItem::SubagentResult { .. } => {
None
}
}
})
.collect())
}
fn actionable_input_notify(
&self,
) -> Result<Arc<tokio::sync::Notify>, meerkat_core::CommsCapabilityError> {
self.actionable_notify.clone().ok_or_else(|| {
meerkat_core::CommsCapabilityError::Unsupported(
"actionable_input_notify: classified inbox not initialized".to_string(),
)
})
}
}
fn map_event_injector_error(error: meerkat_core::event_injector::EventInjectorError) -> SendError {
match error {
meerkat_core::event_injector::EventInjectorError::Full => {
SendError::Validation("input queue full".into())
}
meerkat_core::event_injector::EventInjectorError::Closed => SendError::InputClosed,
}
}
fn map_router_send_error(err: crate::router::SendError) -> SendError {
match err {
crate::router::SendError::PeerNotFound(peer) => SendError::PeerNotFound(peer),
crate::router::SendError::PeerOffline => SendError::PeerOffline,
crate::router::SendError::Transport(_) | crate::router::SendError::Io(_) => {
SendError::Internal(err.to_string())
}
}
}
#[derive(Debug, Error)]
pub enum CommsRuntimeError {
#[error("Identity error: {0}")]
IdentityError(String),
#[error("Trust load error: {0}")]
TrustLoadError(String),
#[error("Listener error: {0}")]
ListenerError(#[from] std::io::Error),
#[error("Listeners already started")]
AlreadyStarted,
#[error("Unsafe binding: {0}")]
UnsafeBinding(String),
}
#[cfg_attr(target_arch = "wasm32", allow(dead_code))]
pub struct CommsRuntime {
public_key: PubKey,
router: Arc<Router>,
trusted_peers: Arc<parking_lot::RwLock<TrustedPeers>>,
inbox: Arc<AsyncMutex<crate::Inbox>>,
inbox_notify: Arc<tokio::sync::Notify>,
#[cfg(not(target_arch = "wasm32"))]
config: ResolvedCommsConfig,
#[cfg(target_arch = "wasm32")]
name: String,
#[cfg(target_arch = "wasm32")]
inproc_namespace: Option<String>,
#[cfg(not(target_arch = "wasm32"))]
listener_handles: Vec<ListenerHandle>,
#[cfg(not(target_arch = "wasm32"))]
listeners_started: bool,
keypair: Arc<Keypair>,
require_peer_auth: bool,
dismiss_flag: AtomicBool,
subscriber_registry: crate::event_injector::SubscriberRegistry,
interaction_stream_registry: InteractionStreamRegistry,
#[allow(dead_code)] silent_intents: Arc<HashSet<String>>,
actionable_notify: Option<Arc<tokio::sync::Notify>>,
}
impl CommsRuntime {
#[cfg(not(target_arch = "wasm32"))]
pub async fn new(config: ResolvedCommsConfig) -> Result<Self, CommsRuntimeError> {
Self::new_with_silent_intents(config, Arc::new(HashSet::new())).await
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn new_with_silent_intents(
config: ResolvedCommsConfig,
silent_intents: Arc<HashSet<String>>,
) -> Result<Self, CommsRuntimeError> {
let keypair = Keypair::load_or_generate(&config.identity_dir)
.await
.map_err(|e| CommsRuntimeError::IdentityError(e.to_string()))?;
let trusted_peers = TrustedPeers::load_or_default(&config.trusted_peers_path)
.map_err(|e| CommsRuntimeError::TrustLoadError(e.to_string()))?;
let public_key = keypair.public_key();
let trusted_peers = Arc::new(parking_lot::RwLock::new(trusted_peers));
let classification_context = Arc::new(crate::classify::IngressClassificationContext {
require_peer_auth: config.require_peer_auth,
trusted_peers: trusted_peers.clone(),
silent_intents: silent_intents.clone(),
});
let (inbox, inbox_sender) = crate::Inbox::new_classified(classification_context);
let inbox_notify = inbox.notify();
let actionable_notify = inbox.classified_actionable_notify();
let router = Router::with_shared_peers(
keypair.clone(),
trusted_peers.clone(),
config.comms_config.clone(),
inbox_sender.clone(),
config.require_peer_auth,
)
.with_inproc_namespace(config.inproc_namespace.clone());
let runtime = Self {
public_key,
router: Arc::new(router),
trusted_peers,
inbox: Arc::new(AsyncMutex::new(inbox)),
inbox_notify,
config: config.clone(),
listener_handles: Vec::new(),
listeners_started: false,
keypair: Arc::new(keypair),
require_peer_auth: config.require_peer_auth,
dismiss_flag: AtomicBool::new(false),
subscriber_registry: crate::event_injector::new_subscriber_registry(),
interaction_stream_registry: Arc::new(Mutex::new(HashMap::new())),
silent_intents,
actionable_notify,
};
InprocRegistry::global().register_with_meta_in_namespace(
config.inproc_namespace.as_deref().unwrap_or(""),
config.name,
runtime.public_key,
inbox_sender,
crate::PeerMeta::default(),
);
Ok(runtime)
}
pub fn inproc_only(name: &str) -> Result<Self, CommsRuntimeError> {
Self::inproc_only_with_silent_intents(name, None, Arc::new(HashSet::new()))
}
pub fn inproc_only_scoped(
name: &str,
namespace: Option<String>,
) -> Result<Self, CommsRuntimeError> {
Self::inproc_only_with_silent_intents(name, namespace, Arc::new(HashSet::new()))
}
pub fn inproc_only_with_silent_intents(
name: &str,
namespace: Option<String>,
silent_intents: Arc<HashSet<String>>,
) -> Result<Self, CommsRuntimeError> {
let keypair = Keypair::generate();
let public_key = keypair.public_key();
let trusted_peers = Arc::new(parking_lot::RwLock::new(TrustedPeers::new()));
let classification_context = Arc::new(crate::classify::IngressClassificationContext {
require_peer_auth: true,
trusted_peers: trusted_peers.clone(),
silent_intents: silent_intents.clone(),
});
let (inbox, inbox_sender) = crate::Inbox::new_classified(classification_context);
let inbox_notify = inbox.notify();
let actionable_notify = inbox.classified_actionable_notify();
let comms_config = crate::CommsConfig::default();
#[cfg(not(target_arch = "wasm32"))]
let config = ResolvedCommsConfig {
enabled: true,
name: name.to_string(),
inproc_namespace: namespace.clone(),
identity_dir: std::path::PathBuf::new(),
trusted_peers_path: std::path::PathBuf::new(),
listen_uds: None,
listen_tcp: None,
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
comms_config: comms_config.clone(),
auth: meerkat_core::CommsAuthMode::Open,
require_peer_auth: true,
allow_external_unauthenticated: false,
};
let router = Router::with_shared_peers(
keypair.clone(),
trusted_peers.clone(),
comms_config,
inbox_sender.clone(),
true,
)
.with_inproc_namespace(namespace.clone());
let runtime = Self {
public_key,
router: Arc::new(router),
trusted_peers,
inbox: Arc::new(AsyncMutex::new(inbox)),
inbox_notify,
#[cfg(not(target_arch = "wasm32"))]
config,
#[cfg(target_arch = "wasm32")]
name: name.to_string(),
#[cfg(target_arch = "wasm32")]
inproc_namespace: namespace.clone(),
#[cfg(not(target_arch = "wasm32"))]
listener_handles: Vec::new(),
#[cfg(not(target_arch = "wasm32"))]
listeners_started: false,
keypair: Arc::new(keypair),
require_peer_auth: true,
dismiss_flag: AtomicBool::new(false),
subscriber_registry: crate::event_injector::new_subscriber_registry(),
interaction_stream_registry: Arc::new(Mutex::new(HashMap::new())),
silent_intents,
actionable_notify,
};
InprocRegistry::global().register_with_meta_in_namespace(
namespace.as_deref().unwrap_or(""),
name,
runtime.public_key,
inbox_sender,
crate::PeerMeta::default(),
);
Ok(runtime)
}
pub fn participant_name(&self) -> &str {
#[cfg(not(target_arch = "wasm32"))]
{
&self.config.name
}
#[cfg(target_arch = "wasm32")]
{
&self.name
}
}
pub fn inproc_namespace(&self) -> Option<&str> {
#[cfg(not(target_arch = "wasm32"))]
{
self.config.inproc_namespace.as_deref()
}
#[cfg(target_arch = "wasm32")]
{
self.inproc_namespace.as_deref()
}
}
#[cfg(not(target_arch = "wasm32"))]
pub fn advertised_address(&self) -> String {
if let Some(ref uds) = self.config.listen_uds {
return format!("uds://{}", uds.display());
}
if let Some(ref tcp) = self.config.listen_tcp {
return format!("tcp://{tcp}");
}
format!("inproc://{}", self.config.name)
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn start_listeners(&mut self) -> Result<(), CommsRuntimeError> {
if self.listeners_started {
return Err(CommsRuntimeError::AlreadyStarted);
}
let inbox_sender = self.router.inbox_sender().clone();
let max_line_length = self.config.comms_config.max_message_bytes as usize;
#[cfg(unix)]
if let Some(ref path) = self.config.listen_uds {
let handle = spawn_uds_listener(
path,
self.keypair.clone(),
self.trusted_peers.clone(),
self.require_peer_auth,
inbox_sender.clone(),
)
.await?;
self.listener_handles.push(handle);
}
if let Some(ref addr) = self.config.listen_tcp {
let handle = spawn_tcp_listener(
&addr.to_string(),
self.keypair.clone(),
self.trusted_peers.clone(),
self.require_peer_auth,
inbox_sender.clone(),
)
.await?;
self.listener_handles.push(handle);
}
if self.config.auth == meerkat_core::CommsAuthMode::Open {
if let Some(addr) = &self.config.event_listen_tcp
&& !addr.ip().is_loopback()
&& !self.config.allow_external_unauthenticated
{
return Err(CommsRuntimeError::UnsafeBinding(
"Plain event listener on non-loopback address is a prompt injection \
vector; set allow_external_unauthenticated=true to override"
.to_string(),
));
}
if let Some(ref addr) = self.config.event_listen_tcp {
let handle = spawn_plain_tcp_listener(
&addr.to_string(),
inbox_sender.clone(),
max_line_length,
)
.await?;
self.listener_handles.push(handle);
}
#[cfg(unix)]
if let Some(ref path) = self.config.event_listen_uds {
let handle =
spawn_plain_uds_listener(path, inbox_sender.clone(), max_line_length).await?;
self.listener_handles.push(handle);
}
}
self.listeners_started = true;
Ok(())
}
pub fn public_key(&self) -> PubKey {
self.public_key
}
pub fn router(&self) -> &Router {
&self.router
}
pub fn upsert_trusted_peer(&self, peer: TrustedPeer) {
self.router.add_trusted_peer(peer);
}
pub fn set_peer_meta(&self, meta: crate::PeerMeta) {
InprocRegistry::global().register_with_meta_in_namespace(
self.inproc_namespace().unwrap_or(""),
self.participant_name(),
self.public_key,
self.router.inbox_sender().clone(),
meta,
);
}
async fn resolve_peer_directory(&self) -> Vec<PeerDirectoryEntry> {
let sendable_kinds = vec![
"peer_message".to_string(),
"peer_request".to_string(),
"peer_response".to_string(),
];
let mut peers = Vec::new();
self.for_each_resolved_peer(|peer| {
peers.push(PeerDirectoryEntry {
name: peer.name,
peer_id: peer.peer_id,
address: peer.address,
source: peer.source,
sendable_kinds: sendable_kinds.clone(),
capabilities: serde_json::json!({}),
meta: peer.meta,
});
})
.await;
peers
}
async fn resolve_peer_count(&self) -> usize {
self.for_each_resolved_peer(|_| {}).await
}
async fn for_each_resolved_peer<F>(&self, mut on_peer: F) -> usize
where
F: FnMut(ResolvedPeer),
{
let inproc_peers =
InprocRegistry::global().peers_in_namespace(self.inproc_namespace().unwrap_or(""));
let inproc_by_name: std::collections::HashMap<String, crate::identity::PubKey> =
inproc_peers
.iter()
.map(|p| (p.name.clone(), p.pubkey))
.collect();
let participant_name = self.participant_name().to_string();
let mut trusted_names = HashSet::new();
let mut trusted_pubkeys = HashSet::new();
let mut emitted = 0usize;
{
let trusted = self.trusted_peers.read();
for peer in &trusted.peers {
if peer.name == participant_name || peer.pubkey == self.public_key {
continue;
}
trusted_names.insert(peer.name.clone());
trusted_pubkeys.insert(peer.pubkey);
let source = if inproc_by_name
.get(&peer.name)
.is_some_and(|key| *key == peer.pubkey)
{
PeerDirectorySource::TrustedAndInproc
} else {
PeerDirectorySource::Trusted
};
let name = match PeerName::new(peer.name.clone()) {
Ok(name) => name,
Err(_) => {
tracing::warn!(
peer_name = %peer.name,
peer_id = %peer.pubkey.to_peer_id(),
"skipping invalid peer name in trusted peers"
);
continue;
}
};
on_peer(ResolvedPeer {
name,
peer_id: peer.pubkey.to_peer_id(),
address: peer.addr.clone(),
source,
meta: peer.meta.clone(),
});
emitted += 1;
}
}
if self.require_peer_auth {
return emitted;
}
for inproc in &inproc_peers {
let name = match PeerName::new(inproc.name.clone()) {
Ok(name) => name,
Err(_) => {
tracing::warn!(
peer_name = %inproc.name,
peer_id = %inproc.pubkey.to_peer_id(),
"skipping invalid inproc peer name"
);
continue;
}
};
let peer_name_str = name.as_string();
if trusted_names.contains(name.as_str()) || trusted_pubkeys.contains(&inproc.pubkey) {
continue;
}
if peer_name_str == participant_name || inproc.pubkey == self.public_key {
continue;
}
on_peer(ResolvedPeer {
name,
peer_id: inproc.pubkey.to_peer_id(),
address: format!("inproc://{peer_name_str}"),
source: PeerDirectorySource::Inproc,
meta: inproc.meta.clone(),
});
emitted += 1;
}
emitted
}
async fn send_peer_command(
&self,
peer_name: &str,
kind: crate::types::MessageKind,
) -> Result<Uuid, SendError> {
self.router
.send(peer_name, kind)
.await
.map_err(map_router_send_error)
}
pub fn mark_interaction_complete(&self, interaction_id: Uuid) {
let mut registry = self.interaction_stream_registry.lock();
let mut should_remove = false;
if let Some(entry) = registry.get_mut(&interaction_id)
&& entry.transition(ReservationState::Attached, ReservationState::Completed)
{
tracing::debug!(
interaction_id = %interaction_id,
"interaction stream completed by terminal event"
);
should_remove = true;
}
if should_remove {
self.subscriber_registry.lock().remove(&interaction_id);
registry.remove(&interaction_id);
}
}
pub fn reap_expired_reservations(&self) {
let mut registry = self.interaction_stream_registry.lock();
registry.retain(|id, entry| {
if entry.state == ReservationState::Reserved
&& entry.created_at.elapsed() > RESERVATION_TTL
{
tracing::debug!(interaction_id = %id, "reservation expired (TTL)");
entry.state = ReservationState::Expired;
false
} else {
true
}
});
}
pub fn router_arc(&self) -> Arc<Router> {
self.router.clone()
}
pub fn trusted_peers_shared(&self) -> Arc<parking_lot::RwLock<TrustedPeers>> {
self.trusted_peers.clone()
}
pub fn inbox_notify(&self) -> Arc<tokio::sync::Notify> {
self.inbox_notify.clone()
}
pub fn event_injector(&self) -> Arc<dyn meerkat_core::EventInjector> {
Arc::new(crate::CommsEventInjector::new(
self.router.inbox_sender().clone(),
self.subscriber_registry.clone(),
))
}
#[doc(hidden)]
pub fn interaction_event_injector(
&self,
) -> Arc<dyn meerkat_core::event_injector::SubscribableInjector> {
Arc::new(crate::CommsEventInjector::new(
self.router.inbox_sender().clone(),
self.subscriber_registry.clone(),
))
}
fn register_interaction_stream(
&self,
interaction_id: Uuid,
body: String,
source: meerkat_core::InputSource,
) -> Result<(), SendError> {
let (sender, receiver) = mpsc::channel::<meerkat_core::AgentEvent>(4096);
self.subscriber_registry
.lock()
.insert(interaction_id, sender.clone());
self.interaction_stream_registry.lock().insert(
interaction_id,
StreamRegistryEntry::reserved(sender, receiver),
);
if let Err(error) =
self.router
.inbox_sender()
.send_classified(crate::types::InboxItem::PlainEvent {
body,
source: PlainEventSource::from(source),
interaction_id: Some(interaction_id),
})
{
self.interaction_stream_registry
.lock()
.remove(&interaction_id);
self.subscriber_registry.lock().remove(&interaction_id);
return Err(match error {
crate::inbox::InboxError::Full => SendError::Validation("input queue full".into()),
crate::inbox::InboxError::Closed => SendError::InputClosed,
});
}
Ok(())
}
pub async fn drain_messages(&self) -> Vec<CommsMessage> {
let mut inbox = self.inbox.lock().await;
let entries = inbox.try_drain_classified();
entries
.into_iter()
.filter_map(|entry| CommsMessage::from_classified_entry(&entry))
.collect()
}
pub async fn recv_message(&self) -> Option<CommsMessage> {
loop {
{
let mut inbox = self.inbox.lock().await;
while let Some(entry) = inbox.try_recv_one_classified() {
if let Some(msg) = CommsMessage::from_classified_entry(&entry) {
return Some(msg);
}
}
}
self.inbox_notify.notified().await;
}
}
pub fn shutdown(&mut self) {
#[cfg(not(target_arch = "wasm32"))]
{
for handle in self.listener_handles.drain(..) {
handle.abort();
}
self.listeners_started = false;
}
}
}
impl Drop for CommsRuntime {
fn drop(&mut self) {
self.shutdown();
InprocRegistry::global()
.unregister_in_namespace(self.inproc_namespace().unwrap_or(""), &self.public_key);
}
}
#[cfg(not(target_arch = "wasm32"))]
pub struct ListenerHandle {
handle: JoinHandle<()>,
}
#[cfg(not(target_arch = "wasm32"))]
impl ListenerHandle {
pub fn abort(&self) {
self.handle.abort();
}
}
#[cfg(unix)]
async fn spawn_uds_listener(
path: &Path,
keypair: Arc<Keypair>,
trusted: Arc<parking_lot::RwLock<TrustedPeers>>,
require_peer_auth: bool,
inbox_sender: InboxSender,
) -> Result<ListenerHandle, std::io::Error> {
use tokio::net::UnixListener;
let path = path.to_path_buf();
if let Err(err) = tokio::fs::remove_file(&path).await
&& err.kind() != std::io::ErrorKind::NotFound
{
return Err(err);
}
if let Some(parent) = path.parent().filter(|p| !p.as_os_str().is_empty()) {
tokio::fs::create_dir_all(parent).await?;
}
let listener = UnixListener::bind(&path)?;
let handle = tokio::spawn(async move {
while let Ok((stream, _)) = listener.accept().await {
let (keypair, trusted, inbox_sender) =
(keypair.clone(), trusted.clone(), inbox_sender.clone());
tokio::spawn(async move {
let trusted_snapshot = trusted.read().clone();
let _ = handle_connection(
stream,
require_peer_auth,
&keypair,
&trusted_snapshot,
&inbox_sender,
)
.await;
});
}
});
Ok(ListenerHandle { handle })
}
#[cfg(not(target_arch = "wasm32"))]
async fn spawn_tcp_listener(
addr: &str,
keypair: Arc<Keypair>,
trusted: Arc<parking_lot::RwLock<TrustedPeers>>,
require_peer_auth: bool,
inbox_sender: InboxSender,
) -> Result<ListenerHandle, std::io::Error> {
let listener = TcpListener::bind(addr).await?;
let handle = tokio::spawn(async move {
while let Ok((stream, _)) = listener.accept().await {
let (keypair, trusted, inbox_sender) =
(keypair.clone(), trusted.clone(), inbox_sender.clone());
tokio::spawn(async move {
let trusted_snapshot = trusted.read().clone();
let _ = handle_connection(
stream,
require_peer_auth,
&keypair,
&trusted_snapshot,
&inbox_sender,
)
.await;
});
}
});
Ok(ListenerHandle { handle })
}
#[cfg(not(target_arch = "wasm32"))]
const PLAIN_LISTENER_MAX_CONCURRENT: usize = 64;
#[cfg(not(target_arch = "wasm32"))]
async fn spawn_plain_tcp_listener(
addr: &str,
inbox_sender: InboxSender,
max_line_length: usize,
) -> Result<ListenerHandle, std::io::Error> {
let listener = TcpListener::bind(addr).await?;
let semaphore = Arc::new(tokio::sync::Semaphore::new(PLAIN_LISTENER_MAX_CONCURRENT));
let handle = tokio::spawn(async move {
while let Ok((stream, _peer)) = listener.accept().await {
let sender = inbox_sender.clone();
let sem = semaphore.clone();
tokio::spawn(async move {
let _permit = match sem.acquire().await {
Ok(p) => p,
Err(_) => return, };
crate::plain_listener::handle_plain_connection(
stream,
sender,
max_line_length,
meerkat_core::PlainEventSource::Tcp,
)
.await;
});
}
});
Ok(ListenerHandle { handle })
}
#[cfg(unix)]
async fn spawn_plain_uds_listener(
path: &std::path::Path,
inbox_sender: InboxSender,
max_line_length: usize,
) -> Result<ListenerHandle, std::io::Error> {
use tokio::net::UnixListener;
let path = path.to_path_buf();
if let Err(err) = tokio::fs::remove_file(&path).await
&& err.kind() != std::io::ErrorKind::NotFound
{
return Err(err);
}
if let Some(parent) = path.parent().filter(|p| !p.as_os_str().is_empty()) {
tokio::fs::create_dir_all(parent).await?;
}
let listener = UnixListener::bind(&path)?;
let semaphore = Arc::new(tokio::sync::Semaphore::new(PLAIN_LISTENER_MAX_CONCURRENT));
let handle = tokio::spawn(async move {
while let Ok((stream, _)) = listener.accept().await {
let sender = inbox_sender.clone();
let sem = semaphore.clone();
tokio::spawn(async move {
let _permit = match sem.acquire().await {
Ok(p) => p,
Err(_) => return,
};
crate::plain_listener::handle_plain_connection(
stream,
sender,
max_line_length,
meerkat_core::PlainEventSource::Uds,
)
.await;
});
}
});
Ok(ListenerHandle { handle })
}
#[cfg(all(test, not(target_arch = "wasm32")))]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
use crate::event_injector::CommsEventInjector;
use crate::identity::Signature;
use crate::types::{Envelope, InboxItem, MessageKind, Status};
use futures::StreamExt;
use meerkat_core::event_injector::SubscribableInjector;
use meerkat_core::{
SendError,
comms::{
InputSource, InputStreamMode, PeerDirectorySource, PeerName, StreamError, StreamScope,
},
interaction::InteractionId,
types::SessionId,
};
use tokio::time::{Duration, timeout};
use uuid::Uuid;
fn test_runtime_config(name: &str, tmp: &tempfile::TempDir) -> ResolvedCommsConfig {
ResolvedCommsConfig {
enabled: true,
name: name.to_string(),
inproc_namespace: None,
listen_uds: None,
listen_tcp: None,
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
identity_dir: tmp.path().join("identity"),
trusted_peers_path: tmp.path().join("trusted_peers.json"),
comms_config: crate::CommsConfig::default(),
auth: meerkat_core::CommsAuthMode::Open,
require_peer_auth: true,
allow_external_unauthenticated: false,
}
}
fn signed_envelope(from: &Keypair, to: PubKey, kind: MessageKind) -> Envelope {
let mut envelope = Envelope {
id: Uuid::new_v4(),
from: from.public_key(),
to,
kind,
sig: Signature::new([0u8; 64]),
};
envelope.sign(from);
envelope
}
#[tokio::test]
async fn test_auth_open_loads_keypair_and_peers() {
let tmp = tempfile::TempDir::new().unwrap();
let config = ResolvedCommsConfig {
enabled: true,
name: "test-agent".to_string(),
inproc_namespace: None,
listen_uds: None,
listen_tcp: None,
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
identity_dir: tmp.path().join("identity"),
trusted_peers_path: tmp.path().join("trusted_peers.json"),
comms_config: crate::CommsConfig::default(),
auth: meerkat_core::CommsAuthMode::Open,
require_peer_auth: true,
allow_external_unauthenticated: false,
};
let runtime = CommsRuntime::new(config).await.unwrap();
assert_ne!(runtime.public_key().as_bytes(), &[0u8; 32]);
}
#[tokio::test]
#[ignore] async fn test_signed_listener_starts_in_open_mode() {
let tmp = tempfile::TempDir::new().unwrap();
let config = ResolvedCommsConfig {
enabled: true,
name: "test-signed".to_string(),
inproc_namespace: None,
listen_uds: None,
listen_tcp: Some("127.0.0.1:0".parse().unwrap()),
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
identity_dir: tmp.path().join("identity"),
trusted_peers_path: tmp.path().join("trusted_peers.json"),
comms_config: crate::CommsConfig::default(),
auth: meerkat_core::CommsAuthMode::Open,
require_peer_auth: true,
allow_external_unauthenticated: false,
};
let mut runtime = CommsRuntime::new(config).await.unwrap();
runtime.start_listeners().await.unwrap();
assert!(runtime.listeners_started);
runtime.shutdown();
}
#[tokio::test]
async fn test_plain_listener_rejects_non_loopback() {
let tmp = tempfile::TempDir::new().unwrap();
let config = ResolvedCommsConfig {
enabled: true,
name: "test-reject".to_string(),
inproc_namespace: None,
listen_uds: None,
listen_tcp: None,
event_listen_tcp: Some("0.0.0.0:4201".parse().unwrap()),
#[cfg(unix)]
event_listen_uds: None,
identity_dir: tmp.path().join("identity"),
trusted_peers_path: tmp.path().join("trusted_peers.json"),
comms_config: crate::CommsConfig::default(),
auth: meerkat_core::CommsAuthMode::Open,
require_peer_auth: true,
allow_external_unauthenticated: false,
};
let mut runtime = CommsRuntime::new(config).await.unwrap();
let result = runtime.start_listeners().await;
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("prompt injection"));
}
#[tokio::test]
async fn test_drain_inbox_interactions_converts_all_authenticated_content_variants() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("variants", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
let sender = Keypair::generate();
runtime.router.add_trusted_peer(crate::TrustedPeer {
name: "sender".to_string(),
pubkey: sender.public_key(),
addr: "tcp://127.0.0.1:4200".to_string(),
meta: crate::PeerMeta::default(),
});
let request_id = Uuid::new_v4();
let reply_to = Uuid::new_v4();
let msg = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Message {
body: "hello".to_string(),
},
);
let req = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
intent: "review".to_string(),
params: serde_json::json!({"pr": 19}),
},
);
let mut req = req;
req.id = request_id;
req.sign(&sender);
let resp = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Response {
in_reply_to: reply_to,
status: Status::Completed,
result: serde_json::json!({"ok": true}),
},
);
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope: msg })
.unwrap();
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope: req })
.unwrap();
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope: resp })
.unwrap();
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 3);
assert!(interactions.iter().any(|i| {
matches!(
&i.content,
meerkat_core::InteractionContent::Message { body } if body == "hello"
)
}));
assert!(interactions.iter().any(|i| {
matches!(
&i.content,
meerkat_core::InteractionContent::Request { intent, params }
if intent == "review" && params["pr"] == 19
)
}));
assert!(interactions.iter().any(|i| {
matches!(
&i.content,
meerkat_core::InteractionContent::Response { in_reply_to, status, result }
if in_reply_to.0 == reply_to
&& *status == meerkat_core::ResponseStatus::Completed
&& result["ok"] == true
)
}));
}
#[tokio::test]
async fn test_subscription_correlation_e2e_one_shot() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("subscription", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
let injector = CommsEventInjector::new(
runtime.router.inbox_sender().clone(),
runtime.subscriber_registry.clone(),
);
let sub = injector
.inject_with_subscription("tracked".to_string(), meerkat_core::PlainEventSource::Rpc)
.unwrap();
let tracked_id = sub.id;
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 1);
assert_eq!(interactions[0].id, tracked_id);
let first = CoreCommsRuntime::interaction_subscriber(&runtime, &tracked_id);
assert!(first.is_some(), "subscriber should be found");
let second = CoreCommsRuntime::interaction_subscriber(&runtime, &tracked_id);
assert!(second.is_none(), "subscriber should be one-shot");
}
#[tokio::test]
async fn test_plain_event_interaction_id_is_preserved_in_drain_inbox_interactions() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("plain-id", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
let interaction_id = Uuid::new_v4();
runtime
.router
.inbox_sender()
.send_classified(InboxItem::PlainEvent {
body: "evt".to_string(),
source: meerkat_core::PlainEventSource::Tcp,
interaction_id: Some(interaction_id),
})
.unwrap();
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 1);
assert_eq!(interactions[0].id.0, interaction_id);
}
#[test]
fn test_interaction_subscriber_correlation_miss_returns_none() {
let runtime = CommsRuntime::inproc_only("corr-miss").unwrap();
let random = meerkat_core::InteractionId(Uuid::new_v4());
let sender = CoreCommsRuntime::interaction_subscriber(&runtime, &random);
assert!(sender.is_none());
}
#[tokio::test]
async fn test_core_send_input_no_reservation() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("input-nores-{suffix}")).unwrap();
let cmd = CommsCommand::Input {
session_id: SessionId::new(),
body: "standalone test input".to_string(),
source: InputSource::Rpc,
stream: InputStreamMode::None,
allow_self_session: true,
};
let receipt = CoreCommsRuntime::send(&runtime, cmd).await;
assert!(receipt.is_ok(), "send should succeed: {receipt:?}");
match receipt.unwrap() {
SendReceipt::InputAccepted {
interaction_id: _,
stream_reserved,
} => assert!(!stream_reserved),
other => panic!("Expected InputAccepted, got: {other:?}"),
}
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 1);
assert!(matches!(
&interactions[0].content,
meerkat_core::InteractionContent::Message { body } if body == "standalone test input"
));
}
#[tokio::test]
async fn test_core_send_input_reserves_stream() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("input-res-{suffix}")).unwrap();
let cmd = CommsCommand::Input {
session_id: SessionId::new(),
body: "streaming input".to_string(),
source: InputSource::Rpc,
stream: InputStreamMode::ReserveInteraction,
allow_self_session: true,
};
let receipt = CoreCommsRuntime::send(&runtime, cmd).await.unwrap();
let (interaction_id, reserved) = match receipt {
SendReceipt::InputAccepted {
interaction_id,
stream_reserved,
} => (interaction_id, stream_reserved),
other => panic!("Expected InputAccepted, got: {other:?}"),
};
assert!(reserved);
let sender = runtime
.take_interaction_stream_sender(&interaction_id)
.expect("reserved interaction should be registered");
drop(sender);
assert!(
runtime
.take_interaction_stream_sender(&interaction_id)
.is_none(),
"interaction registration must be one-shot"
);
}
#[tokio::test]
async fn test_core_stream_attachment_duplicate_attach_fails() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("dup-attach-{suffix}")).unwrap();
let cmd = CommsCommand::Input {
session_id: SessionId::new(),
body: "duplicate stream test".to_string(),
source: InputSource::Rpc,
stream: InputStreamMode::ReserveInteraction,
allow_self_session: true,
};
let interaction_id = match CoreCommsRuntime::send(&runtime, cmd).await.unwrap() {
SendReceipt::InputAccepted {
interaction_id,
stream_reserved,
} => {
assert!(stream_reserved);
interaction_id
}
other => panic!("Expected InputAccepted, got: {other:?}"),
};
let stream = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(interaction_id))
.expect("first attach should succeed");
let dup = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(interaction_id));
assert!(matches!(dup, Err(StreamError::AlreadyAttached(_))));
drop(stream);
let after_drop =
CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(interaction_id));
assert!(matches!(after_drop, Err(StreamError::NotReserved(_))));
}
#[tokio::test]
async fn test_core_stream_not_reserved_before_send() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("pre-send-{suffix}")).unwrap();
let random = InteractionId(Uuid::new_v4());
let missing = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(random));
assert!(matches!(missing, Err(StreamError::NotReserved(_))));
let cmd = CommsCommand::Input {
session_id: SessionId::new(),
body: "send before stream attach".to_string(),
source: InputSource::Rpc,
stream: InputStreamMode::ReserveInteraction,
allow_self_session: true,
};
let interaction_id = match CoreCommsRuntime::send(&runtime, cmd).await.unwrap() {
SendReceipt::InputAccepted {
interaction_id,
stream_reserved,
} => {
assert!(stream_reserved);
interaction_id
}
other => panic!("Expected InputAccepted, got: {other:?}"),
};
let _stream = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(interaction_id))
.expect("should attach after send");
}
#[tokio::test]
async fn test_core_send_and_stream_input_returns_stream_and_receipt() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("sas-{suffix}")).unwrap();
let cmd = CommsCommand::Input {
session_id: SessionId::new(),
body: "stream-first".to_string(),
source: InputSource::Rpc,
stream: InputStreamMode::ReserveInteraction,
allow_self_session: true,
};
let (receipt, _stream) = CoreCommsRuntime::send_and_stream(&runtime, cmd)
.await
.unwrap();
let interaction_id = match receipt {
SendReceipt::InputAccepted {
interaction_id,
stream_reserved,
} => {
assert!(stream_reserved);
interaction_id
}
other => panic!("Expected InputAccepted, got: {other:?}"),
};
let dup = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(interaction_id));
assert!(matches!(dup, Err(StreamError::AlreadyAttached(_))));
assert!(
CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(interaction_id)).is_err()
);
}
#[tokio::test]
async fn test_core_send_unknown_peer_fails() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("sender-{suffix}")).unwrap();
let cmd = CommsCommand::PeerMessage {
to: PeerName::new("missing-peer".to_string())
.expect("missing-peer is a valid peer name"),
body: "hello".to_string(),
};
let result = CoreCommsRuntime::send(&runtime, cmd).await;
assert!(matches!(
result,
Err(SendError::PeerNotFound(peer)) if peer == "missing-peer"
));
}
#[tokio::test]
async fn test_core_send_success_and_drain() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("sender-{suffix}");
let receiver_name = format!("receiver-{suffix}");
let sender = CommsRuntime::inproc_only(&sender_name).unwrap();
let receiver = CommsRuntime::inproc_only(&receiver_name).unwrap();
sender.router.add_trusted_peer(crate::TrustedPeer {
name: receiver_name.clone(),
pubkey: receiver.public_key(),
addr: format!("inproc://{receiver_name}"),
meta: crate::PeerMeta::default(),
});
receiver.router.add_trusted_peer(crate::TrustedPeer {
name: sender_name.clone(),
pubkey: sender.public_key(),
addr: format!("inproc://{sender_name}"),
meta: crate::PeerMeta::default(),
});
let cmd = CommsCommand::PeerMessage {
to: PeerName::new(receiver_name).expect("receiver_name is a valid peer name"),
body: "greeting".to_string(),
};
let receipt = CoreCommsRuntime::send(&sender, cmd).await;
match receipt {
Ok(SendReceipt::PeerMessageSent {
envelope_id: _,
acked,
}) => assert!(!acked),
other => panic!("Expected peer message receipt, got: {other:?}"),
}
let interactions = CoreCommsRuntime::drain_inbox_interactions(&receiver).await;
assert_eq!(interactions.len(), 1);
assert!(matches!(
&interactions[0].content,
meerkat_core::InteractionContent::Message { body } if body == "greeting"
));
}
#[tokio::test]
async fn test_core_send_fails_without_trusted_entry() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("sender-ipc-{suffix}");
let receiver_name = format!("receiver-ipc-{suffix}");
let sender = CommsRuntime::inproc_only(&sender_name).unwrap();
let receiver = CommsRuntime::inproc_only(&receiver_name).unwrap();
let cmd = CommsCommand::PeerMessage {
to: PeerName::new(receiver_name.clone()).expect("receiver_name is a valid peer name"),
body: "inproc-only hello".to_string(),
};
let receipt = CoreCommsRuntime::send(&sender, cmd).await;
match receipt {
Err(SendError::PeerNotFound(peer)) => assert_eq!(peer, receiver_name),
other => panic!("Expected peer-not-found for missing trust, got: {other:?}"),
}
let interactions = CoreCommsRuntime::drain_inbox_interactions(&receiver).await;
assert_eq!(interactions.len(), 0);
}
#[tokio::test]
async fn test_core_send_succeeds_without_trusted_entry_when_auth_disabled() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("sender-ipc-no-auth-{suffix}");
let receiver_name = format!("receiver-ipc-no-auth-{suffix}");
let sender_tmp = tempfile::tempdir().unwrap();
let receiver_tmp = tempfile::tempdir().unwrap();
let sender_config = ResolvedCommsConfig {
enabled: true,
name: sender_name.clone(),
inproc_namespace: None,
identity_dir: sender_tmp.path().join("identity"),
trusted_peers_path: sender_tmp.path().join("trusted_peers.json"),
comms_config: crate::CommsConfig::default(),
auth: meerkat_core::CommsAuthMode::Open,
require_peer_auth: false,
listen_uds: None,
listen_tcp: None,
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
allow_external_unauthenticated: false,
};
let receiver_config = ResolvedCommsConfig {
enabled: true,
name: receiver_name.clone(),
inproc_namespace: None,
identity_dir: receiver_tmp.path().join("identity"),
trusted_peers_path: receiver_tmp.path().join("trusted_peers.json"),
comms_config: crate::CommsConfig::default(),
auth: meerkat_core::CommsAuthMode::Open,
require_peer_auth: false,
listen_uds: None,
listen_tcp: None,
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
allow_external_unauthenticated: false,
};
let sender = CommsRuntime::new(sender_config).await.unwrap();
let receiver = CommsRuntime::new(receiver_config).await.unwrap();
let sender_pubkey = sender.public_key();
let receiver_pubkey = receiver.public_key();
let cmd = CommsCommand::PeerMessage {
to: PeerName::new(receiver_name.clone()).expect("receiver_name is a valid peer name"),
body: "hello without trusted".to_string(),
};
let receipt = CoreCommsRuntime::send(&sender, cmd).await;
match receipt {
Ok(SendReceipt::PeerMessageSent {
envelope_id: _,
acked,
}) => assert!(!acked),
other => panic!("Expected peer message receipt, got: {other:?}"),
}
let interactions = CoreCommsRuntime::drain_inbox_interactions(&receiver).await;
assert_eq!(interactions.len(), 1);
assert_eq!(interactions[0].from, sender_name);
assert!(matches!(
&interactions[0].content,
meerkat_core::InteractionContent::Message { body } if body == "hello without trusted"
));
assert!(crate::InprocRegistry::global().unregister(&sender_pubkey));
assert!(crate::InprocRegistry::global().unregister(&receiver_pubkey));
}
#[tokio::test]
async fn test_add_trusted_peer_updates_peers_and_enables_send() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("trust-sender-{suffix}");
let receiver_name = format!("trust-receiver-{suffix}");
let sender = CommsRuntime::inproc_only(&sender_name).unwrap();
let receiver = CommsRuntime::inproc_only(&receiver_name).unwrap();
let peer_spec = meerkat_core::comms::TrustedPeerSpec::new(
&receiver_name,
receiver.public_key().to_peer_id(),
format!("inproc://{receiver_name}"),
)
.expect("valid peer spec");
CoreCommsRuntime::add_trusted_peer(&sender, peer_spec)
.await
.expect("trusted peer add should succeed");
let reverse_spec = meerkat_core::comms::TrustedPeerSpec::new(
&sender_name,
sender.public_key().to_peer_id(),
format!("inproc://{sender_name}"),
)
.expect("valid reverse peer spec");
CoreCommsRuntime::add_trusted_peer(&receiver, reverse_spec)
.await
.expect("reverse trusted peer add should succeed");
let peers = CoreCommsRuntime::peers(&sender).await;
let receiver_entries: Vec<_> = peers
.iter()
.filter(|entry| entry.name.as_str() == receiver_name)
.collect();
assert_eq!(receiver_entries.len(), 1, "peer should appear in peers()");
let send_cmd = CommsCommand::PeerMessage {
to: PeerName::new(receiver_name.clone()).expect("receiver_name is a valid peer name"),
body: "hello trusted peer".to_string(),
};
let receipt = CoreCommsRuntime::send(&sender, send_cmd).await;
assert!(matches!(receipt, Ok(SendReceipt::PeerMessageSent { .. })));
let interactions = CoreCommsRuntime::drain_inbox_interactions(&receiver).await;
assert_eq!(interactions.len(), 1);
assert!(matches!(
&interactions[0].content,
meerkat_core::InteractionContent::Message { body } if body == "hello trusted peer"
));
}
#[tokio::test]
async fn test_add_trusted_peer_invalid_peer_id_is_rejected() {
let sender = CommsRuntime::inproc_only("trust-invalid-sender").unwrap();
let invalid_peer_spec = meerkat_core::comms::TrustedPeerSpec {
name: "invalid".to_string(),
peer_id: "bad-peer-id".to_string(),
address: "inproc://invalid".to_string(),
};
let result = CoreCommsRuntime::add_trusted_peer(&sender, invalid_peer_spec).await;
assert!(matches!(result, Err(SendError::Validation(_))));
}
#[test]
fn test_core_runtime_public_key_is_exposed_via_trait() {
let runtime = CommsRuntime::inproc_only("pub-key-trait").unwrap();
let public_key = <CommsRuntime as CoreCommsRuntime>::public_key(&runtime);
assert!(
public_key
.as_ref()
.is_some_and(|id| id.starts_with("ed25519:")),
"public_key should be available and formatted as ed25519 peer id"
);
}
#[tokio::test]
async fn test_core_peers_includes_trusted_without_self() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("trusted-only-{suffix}");
let runtime_name = format!("runtime-mixed-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = CommsRuntime::inproc_only(&runtime_name).unwrap();
{
let mut trusted = runtime.trusted_peers.write();
trusted.upsert(crate::TrustedPeer {
name: peer_name.clone(),
pubkey: peer.public_key(),
addr: format!("inproc://{peer_name}"),
meta: crate::PeerMeta::default(),
});
}
let peers = CoreCommsRuntime::peers(&runtime).await;
let names: Vec<_> = peers.iter().map(|entry| entry.name.as_string()).collect();
assert!(names.iter().any(|name| name == &peer_name));
assert!(!names.iter().any(|name| name == &runtime_name));
let matched: Vec<_> = peers
.iter()
.filter(|entry| entry.name.as_str() == peer_name)
.collect();
assert_eq!(matched.len(), 1);
let peer = matched[0];
assert!(
matches!(
peer.source,
PeerDirectorySource::Trusted | PeerDirectorySource::TrustedAndInproc
),
"expected Trusted or TrustedAndInproc, got {:?}",
peer.source
);
assert_eq!(peer.address, format!("inproc://{peer_name}"));
assert_eq!(peer.sendable_kinds.len(), 3);
assert!(peer.sendable_kinds.contains(&"peer_message".to_string()));
assert!(peer.sendable_kinds.contains(&"peer_request".to_string()));
assert!(peer.sendable_kinds.contains(&"peer_response".to_string()));
}
#[tokio::test]
async fn test_m3_truthfulness_invariant_all_advertised_peers_are_sendable() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("truth-peer-{suffix}");
let runtime_name = format!("truth-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = CommsRuntime::inproc_only(&runtime_name).unwrap();
{
let mut trusted = runtime.trusted_peers.write();
trusted.upsert(crate::TrustedPeer {
name: peer_name.clone(),
pubkey: peer.public_key(),
addr: format!("inproc://{peer_name}"),
meta: crate::PeerMeta::default(),
});
}
let all_peers = CoreCommsRuntime::peers(&runtime).await;
let peers: Vec<_> = all_peers
.iter()
.filter(|e| e.name.as_str() == peer_name)
.collect();
assert_eq!(
peers.len(),
1,
"peer should be visible when explicitly trusted"
);
for entry in peers {
for kind in &entry.sendable_kinds {
let cmd = match kind.as_str() {
"peer_message" => CommsCommand::PeerMessage {
to: entry.name.clone(),
body: "truthfulness test".to_string(),
},
"peer_request" => CommsCommand::PeerRequest {
to: entry.name.clone(),
intent: "test".to_string(),
params: serde_json::json!({}),
stream: InputStreamMode::None,
},
"peer_response" => CommsCommand::PeerResponse {
to: entry.name.clone(),
in_reply_to: meerkat_core::InteractionId(Uuid::new_v4()),
status: meerkat_core::ResponseStatus::Completed,
result: serde_json::json!({}),
},
_ => continue,
};
let result = CoreCommsRuntime::send(&runtime, cmd).await;
assert!(
!matches!(result, Err(SendError::PeerNotFound(_))),
"peer '{}' advertised kind '{}' but send failed with PeerNotFound",
entry.name.as_str(),
kind
);
}
}
}
#[tokio::test]
async fn test_m4_duplicate_close_is_safe() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("dup-close-{suffix}")).unwrap();
let session_id = meerkat_core::SessionId::new();
let cmd = CommsCommand::Input {
session_id: session_id.clone(),
body: "hello".to_string(),
source: meerkat_core::InputSource::Rpc,
stream: InputStreamMode::ReserveInteraction,
allow_self_session: true,
};
let receipt = CoreCommsRuntime::send(&runtime, cmd).await.unwrap();
let iid = match receipt {
SendReceipt::InputAccepted { interaction_id, .. } => interaction_id,
_ => panic!("expected InputAccepted"),
};
let stream = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(iid)).unwrap();
drop(stream);
let result = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(iid));
assert!(matches!(result, Err(StreamError::NotReserved(_))));
}
#[tokio::test]
async fn test_m4_mark_interaction_complete_cleans_up() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("complete-{suffix}")).unwrap();
let session_id = meerkat_core::SessionId::new();
let cmd = CommsCommand::Input {
session_id: session_id.clone(),
body: "hello".to_string(),
source: meerkat_core::InputSource::Rpc,
stream: InputStreamMode::ReserveInteraction,
allow_self_session: true,
};
let receipt = CoreCommsRuntime::send(&runtime, cmd).await.unwrap();
let iid = match receipt {
SendReceipt::InputAccepted { interaction_id, .. } => interaction_id,
_ => panic!("expected InputAccepted"),
};
let mut stream = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(iid)).unwrap();
runtime.mark_interaction_complete(iid.0);
let terminal = timeout(Duration::from_millis(100), stream.next())
.await
.expect("stream should terminate promptly after completion");
assert!(terminal.is_none());
let registry = runtime.interaction_stream_registry.lock();
assert!(
!registry.contains_key(&iid.0),
"completed entry should be cleaned from registry"
);
}
#[tokio::test]
async fn test_m4_reap_expired_reservations() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("reap-{suffix}")).unwrap();
let id = Uuid::new_v4();
let (sender, receiver) = mpsc::channel::<meerkat_core::AgentEvent>(16);
{
let mut entry = StreamRegistryEntry::reserved(sender, receiver);
entry.created_at = Instant::now()
.checked_sub(meerkat_core::time_compat::Duration::from_secs(60))
.unwrap();
runtime.interaction_stream_registry.lock().insert(id, entry);
}
runtime.reap_expired_reservations();
let registry = runtime.interaction_stream_registry.lock();
assert!(
!registry.contains_key(&id),
"expired reservation should have been reaped"
);
}
#[tokio::test]
async fn test_m5_send_and_stream_non_streamable_returns_stream_attach() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("m5-peer-{suffix}");
let runtime_name = format!("m5-runtime-{suffix}");
let _peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = CommsRuntime::inproc_only(&runtime_name).unwrap();
{
let mut trusted = runtime.trusted_peers.write();
trusted.upsert(crate::TrustedPeer {
name: peer_name.clone(),
pubkey: _peer.public_key(),
addr: format!("inproc://{peer_name}"),
meta: crate::PeerMeta::default(),
});
}
let cmd = CommsCommand::PeerMessage {
to: PeerName::new(peer_name).expect("peer_name is a valid peer name"),
body: "not streamable".to_string(),
};
let result = CoreCommsRuntime::send_and_stream(&runtime, cmd).await;
match result {
Err(SendAndStreamError::StreamAttach { receipt, error }) => {
assert!(
matches!(receipt, SendReceipt::PeerMessageSent { .. }),
"receipt should be PeerMessageSent"
);
assert!(
matches!(error, StreamError::NotFound(_)),
"error should be NotFound for non-streamable command"
);
}
Err(e) => panic!("expected StreamAttach error, got send error: {e}"),
Ok(_) => panic!("expected StreamAttach error, got Ok"),
}
}
#[tokio::test]
async fn test_m5_send_and_stream_input_none_is_not_streamable() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("m5-none-{suffix}")).unwrap();
let session_id = meerkat_core::SessionId::new();
let cmd = CommsCommand::Input {
session_id: session_id.clone(),
body: "no stream".to_string(),
source: meerkat_core::InputSource::Rpc,
stream: InputStreamMode::None,
allow_self_session: true,
};
let result = CoreCommsRuntime::send_and_stream(&runtime, cmd).await;
match result {
Err(SendAndStreamError::StreamAttach { receipt, error }) => {
assert!(
matches!(
receipt,
SendReceipt::InputAccepted {
stream_reserved: false,
..
}
),
"receipt should indicate no stream reserved"
);
assert!(matches!(error, StreamError::NotFound(_)));
}
Err(e) => panic!("expected StreamAttach error, got send error: {e}"),
Ok(_) => panic!("expected StreamAttach error, got Ok"),
}
}
#[tokio::test]
async fn test_m4_attach_after_completed_returns_not_reserved() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("post-comp-{suffix}")).unwrap();
let session_id = meerkat_core::SessionId::new();
let cmd = CommsCommand::Input {
session_id: session_id.clone(),
body: "hello".to_string(),
source: meerkat_core::InputSource::Rpc,
stream: InputStreamMode::ReserveInteraction,
allow_self_session: true,
};
let receipt = CoreCommsRuntime::send(&runtime, cmd).await.unwrap();
let iid = match receipt {
SendReceipt::InputAccepted { interaction_id, .. } => interaction_id,
_ => panic!("expected InputAccepted"),
};
let _stream = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(iid)).unwrap();
runtime.mark_interaction_complete(iid.0);
let result = CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(iid));
assert!(
matches!(result, Err(StreamError::NotReserved(_))),
"attach after completed should return NotReserved"
);
}
#[tokio::test]
async fn test_allow_self_session_guard_rejects_default() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("self-guard-{suffix}")).unwrap();
let cmd = CommsCommand::Input {
session_id: meerkat_core::SessionId::new(),
body: "blocked".to_string(),
source: meerkat_core::InputSource::Rpc,
stream: InputStreamMode::None,
allow_self_session: false,
};
let result = CoreCommsRuntime::send(&runtime, cmd).await;
assert!(
matches!(result, Err(SendError::Validation(_))),
"allow_self_session=false should reject self-input"
);
}
#[tokio::test]
async fn test_remove_trusted_peer_removes_from_peers_and_revokes_send() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("rm-sender-{suffix}");
let receiver_name = format!("rm-receiver-{suffix}");
let sender = CommsRuntime::inproc_only(&sender_name).unwrap();
let receiver = CommsRuntime::inproc_only(&receiver_name).unwrap();
let peer_spec = meerkat_core::comms::TrustedPeerSpec::new(
&receiver_name,
receiver.public_key().to_peer_id(),
format!("inproc://{receiver_name}"),
)
.expect("valid peer spec");
CoreCommsRuntime::add_trusted_peer(&sender, peer_spec)
.await
.expect("add_trusted_peer should succeed");
let peers = CoreCommsRuntime::peers(&sender).await;
let receiver_entries: Vec<_> = peers
.iter()
.filter(|entry| entry.name.as_str() == receiver_name)
.collect();
assert_eq!(
receiver_entries.len(),
1,
"receiver should appear in sender's peers after add"
);
let peer_id = receiver.public_key().to_peer_id();
let removed = CoreCommsRuntime::remove_trusted_peer(&sender, &peer_id)
.await
.expect("remove_trusted_peer should succeed");
assert!(removed, "remove should return true for existing peer");
let peers_after = CoreCommsRuntime::peers(&sender).await;
let receiver_entries_after: Vec<_> = peers_after
.iter()
.filter(|entry| entry.name.as_str() == receiver_name)
.collect();
assert!(
receiver_entries_after.is_empty(),
"receiver should not appear in sender's peers after removal"
);
let cmd = CommsCommand::PeerMessage {
to: PeerName::new(receiver_name.clone()).expect("valid peer name"),
body: "should fail".to_string(),
};
let result = CoreCommsRuntime::send(&sender, cmd).await;
assert!(
matches!(result, Err(SendError::PeerNotFound(_))),
"send to removed peer should fail with PeerNotFound, got: {result:?}"
);
}
#[tokio::test]
async fn test_remove_trusted_peer_returns_false_for_absent_peer() {
let suffix = Uuid::new_v4().simple().to_string();
let sender = CommsRuntime::inproc_only(&format!("rm-absent-{suffix}")).unwrap();
let other = CommsRuntime::inproc_only(&format!("rm-absent-other-{suffix}")).unwrap();
let peer_id = other.public_key().to_peer_id();
let removed = CoreCommsRuntime::remove_trusted_peer(&sender, &peer_id)
.await
.expect("remove_trusted_peer should succeed even for absent peer");
assert!(
!removed,
"remove should return false for peer that was never trusted"
);
}
}