#[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;
use crate::peer_directory_reachability_authority::{
PeerDirectoryReachabilityAuthority, ReachabilityKey,
};
#[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,
PeerReachabilityReason, SendAndStreamError, SendError, SendReceipt, StreamError, StreamScope,
TrustedPeerSpec,
};
use meerkat_core::config::PlainEventSource;
use meerkat_core::hydrate_content_blocks;
use meerkat_core::time_compat::Instant;
use meerkat_core::{BlobStore, MissingBlobBehavior};
use parking_lot::Mutex;
use std::collections::{HashMap, HashSet};
#[cfg(unix)]
use std::path::Path;
use std::pin::Pin;
use std::sync::Arc;
#[cfg(not(target_arch = "wasm32"))]
use std::sync::LazyLock;
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>>>;
#[cfg(not(target_arch = "wasm32"))]
static SESSION_IDENTITY_CLAIMS: LazyLock<Mutex<HashSet<String>>> =
LazyLock::new(|| Mutex::new(HashSet::new()));
#[cfg(not(target_arch = "wasm32"))]
struct SessionIdentityClaim {
session_id: String,
}
#[cfg(not(target_arch = "wasm32"))]
impl SessionIdentityClaim {
fn acquire(session_id: &meerkat_core::SessionId) -> Result<Self, CommsRuntimeError> {
let session_id = session_id.to_string();
let mut claims = SESSION_IDENTITY_CLAIMS.lock();
if !claims.insert(session_id.clone()) {
return Err(CommsRuntimeError::SessionIdentityInUse(session_id));
}
Ok(Self { session_id })
}
}
#[cfg(not(target_arch = "wasm32"))]
impl Drop for SessionIdentityClaim {
fn drop(&mut self) {
SESSION_IDENTITY_CLAIMS.lock().remove(&self.session_id);
}
}
#[cfg(not(target_arch = "wasm32"))]
pub fn release_session_claim(session_id: &meerkat_core::SessionId) -> bool {
SESSION_IDENTITY_CLAIMS
.lock()
.remove(&session_id.to_string())
}
#[cfg(not(target_arch = "wasm32"))]
pub fn clear_all_session_claims() {
SESSION_IDENTITY_CLAIMS.lock().clear();
}
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 ResolvedPeer {
fn reachability_key(&self) -> ReachabilityKey {
ReachabilityKey::new(self.name.as_str(), self.peer_id.as_str())
}
}
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: _,
blocks,
handling_mode,
} => {
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())
})?;
let content = match blocks {
Some(blocks) => meerkat_core::types::ContentInput::Blocks(blocks),
None => meerkat_core::types::ContentInput::Text(body),
};
injector
.inject(content, PlainEventSource::from(source), handling_mode, None)
.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,
blocks,
source,
handling_mode,
)?;
Ok(SendReceipt::InputAccepted {
interaction_id: meerkat_core::InteractionId(interaction_id),
stream_reserved: true,
})
}
}
}
CommsCommand::PeerMessage { to, body, blocks } => self
.send_peer_command(
to.as_str(),
crate::types::MessageKind::Message { body, blocks },
)
.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: _,
blocks,
handling_mode,
} => {
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,
blocks,
source,
handling_mode,
)?;
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 { blocks: None, 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, blocks } => {
let rendered = format!(
"[COMMS MESSAGE from {from_peer}]\n{body}"
);
(
meerkat_core::InteractionContent::Message { body, blocks },
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,
handling_mode: meerkat_core::types::HandlingMode::Queue,
render_metadata: None,
},
class: entry.class,
lifecycle_peer: entry.lifecycle_peer,
})
}
crate::types::InboxItem::PlainEvent {
body,
source,
handling_mode,
interaction_id,
blocks,
render_metadata,
..
} => {
let rendered = meerkat_core::interaction::format_external_event_projection(
&source.to_string(),
Some(&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, blocks },
rendered_text: rendered,
handling_mode,
render_metadata,
},
class: entry.class,
lifecycle_peer: entry.lifecycle_peer,
})
}
}
})
.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,
}
}
#[derive(Debug, Error)]
pub enum CommsRuntimeError {
#[error("Identity error: {0}")]
IdentityError(String),
#[error("Session identity already active: {0}")]
SessionIdentityInUse(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),
}
pub struct CommsToolMaterial {
router: Arc<Router>,
trusted_peers: Arc<parking_lot::RwLock<TrustedPeers>>,
self_pubkey: crate::identity::PubKey,
runtime: Arc<dyn meerkat_core::agent::CommsRuntime>,
}
impl CommsToolMaterial {
pub fn router(&self) -> &Arc<Router> {
&self.router
}
pub fn trusted_peers(&self) -> &Arc<parking_lot::RwLock<TrustedPeers>> {
&self.trusted_peers
}
pub fn trusted_peers_shared(&self) -> Arc<parking_lot::RwLock<TrustedPeers>> {
self.trusted_peers.clone()
}
pub fn self_pubkey(&self) -> crate::identity::PubKey {
self.self_pubkey
}
pub fn into_runtime(self) -> Arc<dyn meerkat_core::agent::CommsRuntime> {
self.runtime
}
}
#[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,
#[cfg(not(target_arch = "wasm32"))]
_session_identity_claim: Option<SessionIdentityClaim>,
keypair: Arc<Keypair>,
require_peer_auth: bool,
dismiss_flag: AtomicBool,
subscriber_registry: crate::event_injector::SubscriberRegistry,
interaction_stream_registry: InteractionStreamRegistry,
peer_directory_reachability: Arc<Mutex<PeerDirectoryReachabilityAuthority>>,
#[allow(dead_code)] silent_intents: Arc<HashSet<String>>,
actionable_notify: Option<Arc<tokio::sync::Notify>>,
blob_store: Option<Arc<dyn BlobStore>>,
}
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,
_session_identity_claim: None,
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())),
peer_directory_reachability: Arc::new(Mutex::new(
PeerDirectoryReachabilityAuthority::new(),
)),
silent_intents,
actionable_notify,
blob_store: None,
};
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,
#[cfg(not(target_arch = "wasm32"))]
_session_identity_claim: None,
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())),
peer_directory_reachability: Arc::new(Mutex::new(
PeerDirectoryReachabilityAuthority::new(),
)),
silent_intents,
actionable_notify,
blob_store: None,
};
InprocRegistry::global().register_with_meta_in_namespace(
namespace.as_deref().unwrap_or(""),
name,
runtime.public_key,
inbox_sender,
crate::PeerMeta::default(),
);
Ok(runtime)
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn inproc_only_session_scoped_with_silent_intents(
name: &str,
namespace: Option<String>,
identity_root: std::path::PathBuf,
session_id: &meerkat_core::SessionId,
silent_intents: Arc<HashSet<String>>,
) -> Result<Self, CommsRuntimeError> {
let claim = SessionIdentityClaim::acquire(session_id)?;
let identity_dir = identity_root.join(session_id.to_string());
let keypair = Keypair::load_or_generate(&identity_dir)
.await
.map_err(|e| CommsRuntimeError::IdentityError(e.to_string()))?;
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();
let config = ResolvedCommsConfig {
enabled: true,
name: name.to_string(),
inproc_namespace: namespace.clone(),
identity_dir: identity_dir.clone(),
trusted_peers_path: identity_dir.join("trusted_peers.json"),
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,
config,
listener_handles: Vec::new(),
listeners_started: false,
_session_identity_claim: Some(claim),
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())),
peer_directory_reachability: Arc::new(Mutex::new(
PeerDirectoryReachabilityAuthority::new(),
)),
silent_intents,
actionable_notify,
blob_store: None,
};
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 set_blob_store(&mut self, blob_store: Arc<dyn BlobStore>) {
self.blob_store = Some(blob_store);
}
async fn hydrate_message_kind_for_transport(
&self,
kind: crate::types::MessageKind,
) -> Result<crate::types::MessageKind, SendError> {
match kind {
crate::types::MessageKind::Message {
body,
blocks: Some(mut blocks),
} => {
if blocks.iter().any(|block| {
matches!(
block,
meerkat_core::types::ContentBlock::Image {
data: meerkat_core::types::ImageData::Blob { .. },
..
}
)
}) {
let blob_store = self.blob_store.as_ref().ok_or_else(|| {
SendError::Internal(
"blob-backed comms message requires blob store".to_string(),
)
})?;
hydrate_content_blocks(
blob_store.as_ref(),
&mut blocks,
MissingBlobBehavior::Error,
)
.await
.map_err(|err| SendError::Internal(err.to_string()))?;
}
Ok(crate::types::MessageKind::Message {
body,
blocks: Some(blocks),
})
}
other => Ok(other),
}
}
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 resolved = self.resolved_peers_snapshot().await;
self.reconcile_peer_directory(&resolved);
let sendable_kinds = vec![
"peer_message".to_string(),
"peer_request".to_string(),
"peer_response".to_string(),
];
let mut peers = Vec::new();
let reachability = self.peer_directory_reachability.lock();
for peer in resolved {
let snapshot = reachability.snapshot_for(&peer.reachability_key());
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!({}),
reachability: snapshot.reachability,
last_unreachable_reason: snapshot.last_unreachable_reason,
meta: peer.meta,
});
}
peers
}
async fn resolve_peer_count(&self) -> usize {
self.resolved_peers_snapshot().await.len()
}
async fn resolved_peers_snapshot(&self) -> Vec<ResolvedPeer> {
let mut peers = Vec::new();
self.for_each_resolved_peer(|peer| peers.push(peer)).await;
peers
}
fn reconcile_peer_directory(&self, peers: &[ResolvedPeer]) {
self.peer_directory_reachability
.lock()
.reconcile_resolved_directory(peers.iter().map(ResolvedPeer::reachability_key));
}
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> {
let resolved = self.resolved_peers_snapshot().await;
self.reconcile_peer_directory(&resolved);
let resolved_peer = resolved
.into_iter()
.find(|peer| peer.name.as_str() == peer_name);
let kind = self.hydrate_message_kind_for_transport(kind).await?;
let result = self.router.send(peer_name, kind).await;
match result {
Ok(envelope_id) => {
if let Some(peer) = resolved_peer.as_ref() {
self.peer_directory_reachability
.lock()
.record_send_succeeded(&peer.reachability_key());
}
Ok(envelope_id)
}
Err(crate::router::SendError::PeerNotFound(peer)) => {
if let Some(resolved_peer) = resolved_peer.as_ref() {
self.peer_directory_reachability.lock().record_send_failed(
&resolved_peer.reachability_key(),
PeerReachabilityReason::OfflineOrNoAck,
);
Err(SendError::PeerNotFound(peer))
} else {
Err(SendError::PeerNotFound(peer))
}
}
Err(crate::router::SendError::PeerOffline) => {
if let Some(peer) = resolved_peer.as_ref() {
self.peer_directory_reachability.lock().record_send_failed(
&peer.reachability_key(),
PeerReachabilityReason::OfflineOrNoAck,
);
}
Err(SendError::PeerOffline)
}
Err(
error @ (crate::router::SendError::Transport(_) | crate::router::SendError::Io(_)),
) => {
if let Some(peer) = resolved_peer.as_ref() {
self.peer_directory_reachability.lock().record_send_failed(
&peer.reachability_key(),
PeerReachabilityReason::TransportError,
);
}
Err(SendError::Internal(error.to_string()))
}
}
}
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 tool_material(self: &Arc<Self>) -> CommsToolMaterial {
CommsToolMaterial {
router: self.router_arc(),
trusted_peers: self.trusted_peers_shared(),
self_pubkey: self.public_key,
runtime: Arc::clone(self) as Arc<dyn meerkat_core::agent::CommsRuntime>,
}
}
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,
blocks: Option<Vec<meerkat_core::types::ContentBlock>>,
source: meerkat_core::InputSource,
handling_mode: meerkat_core::types::HandlingMode,
) -> 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),
handling_mode,
interaction_id: Some(interaction_id),
blocks,
render_metadata: None,
})
{
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 notified = self.inbox_notify.notified();
{
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);
}
}
}
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 async_trait::async_trait;
use futures::StreamExt;
use meerkat_core::event_injector::SubscribableInjector;
use meerkat_core::{
BlobId, BlobPayload, BlobRef, BlobStore, BlobStoreError, SendError,
comms::{
InputSource, InputStreamMode, PeerDirectorySource, PeerName, PeerReachability,
PeerReachabilityReason, StreamError, StreamScope, TrustedPeerSpec,
},
interaction::InteractionId,
types::{ContentBlock, ImageData, SessionId},
};
use parking_lot::Mutex;
use std::{collections::HashMap, sync::Arc};
use tokio::time::{Duration, timeout};
use uuid::Uuid;
#[derive(Default)]
struct TestBlobStore {
blobs: Mutex<HashMap<BlobId, BlobPayload>>,
}
#[async_trait]
impl BlobStore for TestBlobStore {
async fn put_image(&self, media_type: &str, data: &str) -> Result<BlobRef, BlobStoreError> {
let blob_id = BlobId::from(format!("sha256:test-{media_type}-{data}"));
let payload = BlobPayload {
blob_id: blob_id.clone(),
media_type: media_type.to_string(),
data: data.to_string(),
};
self.blobs.lock().insert(blob_id.clone(), payload);
Ok(BlobRef {
blob_id,
media_type: media_type.to_string(),
})
}
async fn get(&self, blob_id: &BlobId) -> Result<BlobPayload, BlobStoreError> {
self.blobs
.lock()
.get(blob_id)
.cloned()
.ok_or_else(|| BlobStoreError::NotFound(blob_id.clone()))
}
async fn delete(&self, blob_id: &BlobId) -> Result<(), BlobStoreError> {
self.blobs.lock().remove(blob_id);
Ok(())
}
fn is_persistent(&self) -> bool {
false
}
}
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 {
blocks: None,
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_drain_inbox_interactions_multimodal_message_keeps_body_as_projection() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("multimodal-body-projection", &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 blocks = vec![
meerkat_core::types::ContentBlock::Text {
text: "caption text".to_string(),
},
meerkat_core::types::ContentBlock::Image {
media_type: "image/png".to_string(),
data: "abc123".into(),
},
];
let msg = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Message {
body: "please inspect this".to_string(),
blocks: Some(blocks.clone()),
},
);
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope: msg })
.unwrap();
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 1);
let interaction = &interactions[0];
assert_eq!(
interaction.rendered_text,
"[COMMS MESSAGE from sender]\nplease inspect this"
);
match &interaction.content {
meerkat_core::InteractionContent::Message {
body,
blocks: got_blocks,
} => {
assert_eq!(body, "please inspect this");
assert_eq!(got_blocks.as_ref(), Some(&blocks));
}
other => panic!("expected message content, got {other:?}"),
}
}
#[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().into(),
meerkat_core::PlainEventSource::Rpc,
meerkat_core::types::HandlingMode::Queue,
None,
)
.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_subscription_correlation_preserves_inline_blocks() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("subscription-blocks", &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(
meerkat_core::types::ContentInput::Blocks(vec![
meerkat_core::types::ContentBlock::Text {
text: "tracked".to_string(),
},
meerkat_core::types::ContentBlock::Image {
media_type: "image/png".to_string(),
data: "aGVsbG8=".into(),
},
]),
meerkat_core::PlainEventSource::Rpc,
meerkat_core::types::HandlingMode::Queue,
None,
)
.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);
match &interactions[0].content {
meerkat_core::InteractionContent::Message { blocks, .. } => {
let blocks = blocks.as_ref().expect("blocks preserved");
assert!(
blocks.iter().any(|block| matches!(
block,
meerkat_core::types::ContentBlock::Image { .. }
)),
"expected inline image block to survive into drained interaction"
);
}
other => panic!("expected message interaction, got {other:?}"),
}
}
#[tokio::test]
#[ignore = "Phase 1 red-ok comms bridge + parent wait suite"]
async fn runtime_bridge_red_ok_send_and_stream_reserves_one_interaction_channel() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("phase1-bridge-{suffix}")).unwrap();
let cmd = CommsCommand::Input {
blocks: None,
session_id: SessionId::new(),
body: "phase 1 bridge input".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
source: InputSource::Rpc,
stream: InputStreamMode::ReserveInteraction,
allow_self_session: true,
};
let (receipt, _stream) = CoreCommsRuntime::send_and_stream(&runtime, cmd)
.await
.expect("send_and_stream should reserve interaction scope");
let interaction_id = match receipt {
SendReceipt::InputAccepted {
interaction_id,
stream_reserved,
} => {
assert!(
stream_reserved,
"bridge receipt should advertise reservation"
);
interaction_id
}
other => panic!("expected InputAccepted, got {other:?}"),
};
let duplicate =
CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(interaction_id));
assert!(
matches!(duplicate, Err(StreamError::AlreadyAttached(_))),
"reserved interaction streams should reject a second attachment"
);
}
#[tokio::test]
#[ignore = "Phase 1 red-ok comms bridge + parent wait suite"]
async fn runtime_bridge_red_ok_completed_interaction_terminates_reserved_stream() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("phase1-complete-{suffix}")).unwrap();
let receipt = CoreCommsRuntime::send(
&runtime,
CommsCommand::Input {
blocks: None,
session_id: SessionId::new(),
body: "complete reserved interaction".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
source: InputSource::Rpc,
stream: InputStreamMode::ReserveInteraction,
allow_self_session: true,
},
)
.await
.expect("send should succeed");
let interaction_id = match receipt {
SendReceipt::InputAccepted { interaction_id, .. } => interaction_id,
other => panic!("expected InputAccepted, got {other:?}"),
};
let mut stream =
CoreCommsRuntime::stream(&runtime, StreamScope::Interaction(interaction_id))
.expect("reserved stream should attach");
runtime.mark_interaction_complete(interaction_id.0);
let terminal = timeout(Duration::from_millis(100), stream.next())
.await
.expect("stream should terminate promptly after completion");
assert!(
terminal.is_none(),
"interaction completion should close the reserved stream for parent waiters"
);
}
#[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 {
blocks: None,
body: "evt".to_string(),
source: meerkat_core::PlainEventSource::Tcp,
handling_mode: meerkat_core::types::HandlingMode::Queue,
interaction_id: Some(interaction_id),
render_metadata: None,
})
.unwrap();
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 1);
assert_eq!(interactions[0].id.0, interaction_id);
}
#[tokio::test]
async fn test_plain_event_preserves_handling_mode_and_render_metadata_in_classified_drain() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("plain-hints", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
let render_metadata = meerkat_core::types::RenderMetadata {
class: meerkat_core::types::RenderClass::ExternalEvent,
salience: meerkat_core::types::RenderSalience::Urgent,
};
runtime
.router
.inbox_sender()
.send_classified(InboxItem::PlainEvent {
blocks: None,
body: "evt".to_string(),
source: meerkat_core::PlainEventSource::Tcp,
handling_mode: meerkat_core::types::HandlingMode::Steer,
interaction_id: None,
render_metadata: Some(render_metadata.clone()),
})
.unwrap();
let interactions = runtime.drain_classified_inbox_interactions().await.unwrap();
assert_eq!(interactions.len(), 1);
let interaction = &interactions[0].interaction;
assert_eq!(
interaction.handling_mode,
meerkat_core::types::HandlingMode::Steer
);
assert_eq!(interaction.render_metadata, Some(render_metadata));
}
#[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 {
blocks: None,
session_id: SessionId::new(),
body: "standalone test input".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
session_id: SessionId::new(),
body: "streaming input".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
session_id: SessionId::new(),
body: "duplicate stream test".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
session_id: SessionId::new(),
body: "send before stream attach".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
session_id: SessionId::new(),
body: "stream-first".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
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 {
blocks: None,
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_hydrates_blob_refs_before_transport() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("sender-blob-{suffix}");
let receiver_name = format!("receiver-blob-{suffix}");
let mut sender = CommsRuntime::inproc_only(&sender_name).unwrap();
let receiver = CommsRuntime::inproc_only(&receiver_name).unwrap();
let blob_store: Arc<dyn BlobStore> = Arc::new(TestBlobStore::default());
let blob_ref = blob_store
.put_image("image/png", "aGVsbG8=")
.await
.expect("blob stored");
sender.set_blob_store(blob_store);
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 {
blocks: Some(vec![ContentBlock::Image {
media_type: "image/png".to_string(),
data: ImageData::Blob {
blob_id: blob_ref.blob_id,
},
}]),
to: PeerName::new(receiver_name).expect("receiver_name is a valid peer name"),
body: "blob-backed image".to_string(),
};
let receipt = CoreCommsRuntime::send(&sender, cmd).await;
assert!(matches!(receipt, Ok(SendReceipt::PeerMessageSent { .. })));
let interactions = CoreCommsRuntime::drain_inbox_interactions(&receiver).await;
assert_eq!(interactions.len(), 1);
match &interactions[0].content {
meerkat_core::InteractionContent::Message { blocks, .. } => {
let blocks = blocks.as_ref().expect("received blocks");
assert!(matches!(
&blocks[0],
ContentBlock::Image {
data: ImageData::Inline { data },
..
} if data == "aGVsbG8="
));
}
other => panic!("expected message interaction, got {other:?}"),
}
}
#[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 {
blocks: None,
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 {
blocks: None,
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 {
blocks: None,
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_session_scoped_inproc_identity_preserves_peer_id_for_same_session() {
let tmp = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::new();
let runtime = CommsRuntime::inproc_only_session_scoped_with_silent_intents(
"session-scoped-peer",
None,
tmp.path().join("session-identities"),
&session_id,
Arc::new(HashSet::new()),
)
.await
.expect("create first session-scoped runtime");
let first_peer_id = runtime.public_key().to_peer_id();
drop(runtime);
let rebuilt = CommsRuntime::inproc_only_session_scoped_with_silent_intents(
"session-scoped-peer",
None,
tmp.path().join("session-identities"),
&session_id,
Arc::new(HashSet::new()),
)
.await
.expect("recreate session-scoped runtime");
let rebuilt_peer_id = rebuilt.public_key().to_peer_id();
assert_eq!(
rebuilt_peer_id, first_peer_id,
"same session id should preserve the inproc peer_id"
);
}
#[tokio::test]
async fn test_session_scoped_inproc_identity_rejects_second_live_activation() {
let tmp = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::new();
let _runtime = CommsRuntime::inproc_only_session_scoped_with_silent_intents(
"session-scoped-lock",
None,
tmp.path().join("session-identities"),
&session_id,
Arc::new(HashSet::new()),
)
.await
.expect("create first session-scoped runtime");
let error = match CommsRuntime::inproc_only_session_scoped_with_silent_intents(
"session-scoped-lock",
None,
tmp.path().join("session-identities"),
&session_id,
Arc::new(HashSet::new()),
)
.await
{
Ok(_) => panic!("second live activation should fail"),
Err(error) => error,
};
assert!(
matches!(error, CommsRuntimeError::SessionIdentityInUse(_)),
"unexpected error: {error:?}"
);
}
#[tokio::test]
async fn test_session_scoped_inproc_identity_does_not_restore_trust_from_disk() {
let tmp = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::new();
let sender = CommsRuntime::inproc_only_session_scoped_with_silent_intents(
"session-scoped-trust",
None,
tmp.path().join("session-identities"),
&session_id,
Arc::new(HashSet::new()),
)
.await
.expect("create session-scoped runtime");
let peer = CommsRuntime::inproc_only("session-scoped-trust-peer").unwrap();
CoreCommsRuntime::add_trusted_peer(
&sender,
TrustedPeerSpec::new(
"session-scoped-trust-peer",
peer.public_key().to_peer_id(),
format!("inproc://{}", peer.participant_name()),
)
.expect("trusted peer spec"),
)
.await
.expect("add trusted peer");
assert_eq!(
sender.peers().await.len(),
1,
"trust should be visible while live"
);
drop(sender);
let rebuilt = CommsRuntime::inproc_only_session_scoped_with_silent_intents(
"session-scoped-trust",
None,
tmp.path().join("session-identities"),
&session_id,
Arc::new(HashSet::new()),
)
.await
.expect("recreate session-scoped runtime");
assert!(
rebuilt.peers().await.is_empty(),
"session-scoped identity should preserve the keypair without replaying trusted peers"
);
}
#[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()));
assert_eq!(peer.reachability, PeerReachability::Unknown);
assert_eq!(peer.last_unreachable_reason, None);
}
#[tokio::test]
async fn test_core_peers_send_success_marks_peer_reachable() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("reach-sender-{suffix}");
let receiver_name = format!("reach-receiver-{suffix}");
let sender = CommsRuntime::inproc_only(&sender_name).unwrap();
let receiver = CommsRuntime::inproc_only(&receiver_name).unwrap();
CoreCommsRuntime::add_trusted_peer(
&sender,
TrustedPeerSpec::new(
&receiver_name,
receiver.public_key().to_peer_id(),
format!("inproc://{receiver_name}"),
)
.expect("valid trusted peer"),
)
.await
.expect("add trusted peer");
let peers_before = CoreCommsRuntime::peers(&sender).await;
let before = peers_before
.iter()
.find(|entry| entry.name.as_str() == receiver_name)
.expect("receiver should be listed before send");
assert_eq!(before.reachability, PeerReachability::Unknown);
CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
blocks: None,
to: PeerName::new(receiver_name.clone()).expect("valid peer name"),
body: "hello".to_string(),
},
)
.await
.expect("send should succeed");
let peers_after = CoreCommsRuntime::peers(&sender).await;
let after = peers_after
.iter()
.find(|entry| entry.name.as_str() == receiver_name)
.expect("receiver should still be listed");
assert_eq!(after.reachability, PeerReachability::Reachable);
assert_eq!(after.last_unreachable_reason, None);
}
#[tokio::test]
async fn test_core_peers_transport_failure_marks_peer_unreachable() {
let suffix = Uuid::new_v4().simple().to_string();
let sender = CommsRuntime::inproc_only(&format!("transport-sender-{suffix}")).unwrap();
let peer_name = format!("transport-peer-{suffix}");
let peer_key = Keypair::generate().public_key().to_peer_id();
CoreCommsRuntime::add_trusted_peer(
&sender,
TrustedPeerSpec::new(&peer_name, peer_key, "tcp://127.0.0.1:9")
.expect("valid trusted peer"),
)
.await
.expect("add trusted peer");
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
blocks: None,
to: PeerName::new(peer_name.clone()).expect("valid peer name"),
body: "hello".to_string(),
},
)
.await;
assert!(
matches!(result, Err(SendError::Internal(_))),
"transport failure should surface as internal send error, got: {result:?}"
);
let peers = CoreCommsRuntime::peers(&sender).await;
let entry = peers
.iter()
.find(|listed| listed.name.as_str() == peer_name)
.expect("trusted peer should remain listed after transport failure");
assert_eq!(entry.reachability, PeerReachability::Unreachable);
assert_eq!(
entry.last_unreachable_reason,
Some(PeerReachabilityReason::TransportError)
);
}
#[tokio::test]
async fn test_unknown_send_target_does_not_create_directory_entry() {
let suffix = Uuid::new_v4().simple().to_string();
let sender = CommsRuntime::inproc_only(&format!("unknown-target-{suffix}")).unwrap();
let missing_name = format!("missing-{suffix}");
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
blocks: None,
to: PeerName::new(missing_name.clone()).expect("valid peer name"),
body: "hello".to_string(),
},
)
.await;
assert!(
matches!(result, Err(SendError::PeerNotFound(ref peer)) if peer == &missing_name),
"missing peer should fail with PeerNotFound, got: {result:?}"
);
let peers = CoreCommsRuntime::peers(&sender).await;
assert!(
peers
.iter()
.all(|entry| entry.name.as_str() != missing_name),
"unknown attempted target must not become a directory entry"
);
}
#[tokio::test]
async fn test_core_peers_resolved_peer_not_found_marks_peer_unreachable() {
let suffix = Uuid::new_v4().simple().to_string();
let sender =
CommsRuntime::inproc_only(&format!("resolved-missing-sender-{suffix}")).unwrap();
let missing_name = format!("resolved-missing-peer-{suffix}");
let missing_key = Keypair::generate().public_key().to_peer_id();
CoreCommsRuntime::add_trusted_peer(
&sender,
TrustedPeerSpec::new(
&missing_name,
missing_key,
format!("inproc://{missing_name}"),
)
.expect("valid trusted peer"),
)
.await
.expect("add trusted peer");
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
blocks: None,
to: PeerName::new(missing_name.clone()).expect("valid peer name"),
body: "hello".to_string(),
},
)
.await;
assert!(
matches!(result, Err(SendError::PeerNotFound(ref peer)) if peer == &missing_name),
"resolved missing peer should preserve PeerNotFound, got: {result:?}"
);
let peers = CoreCommsRuntime::peers(&sender).await;
let entry = peers
.iter()
.find(|listed| listed.name.as_str() == missing_name)
.expect("trusted peer should remain listed after failed send");
assert_eq!(entry.reachability, PeerReachability::Unreachable);
assert_eq!(
entry.last_unreachable_reason,
Some(PeerReachabilityReason::OfflineOrNoAck)
);
}
#[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 {
blocks: None,
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 {
blocks: None,
session_id: session_id.clone(),
body: "hello".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
session_id: session_id.clone(),
body: "hello".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
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 {
blocks: None,
session_id: session_id.clone(),
body: "no stream".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
session_id: session_id.clone(),
body: "hello".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
session_id: meerkat_core::SessionId::new(),
body: "blocked".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
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 {
blocks: None,
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"
);
}
#[test]
fn test_release_session_claim_allows_reacquire() {
let sid = meerkat_core::SessionId::new();
let claim = SessionIdentityClaim::acquire(&sid).unwrap();
std::mem::forget(claim);
assert!(
SessionIdentityClaim::acquire(&sid).is_err(),
"second acquire should fail while claim is leaked"
);
assert!(release_session_claim(&sid));
let claim2 = SessionIdentityClaim::acquire(&sid)
.expect("acquire should succeed after force-release");
drop(claim2);
}
#[test]
fn test_clear_all_session_claims_allows_reacquire() {
let sid = meerkat_core::SessionId::new();
let claim = SessionIdentityClaim::acquire(&sid).unwrap();
std::mem::forget(claim);
clear_all_session_claims();
let claim2 =
SessionIdentityClaim::acquire(&sid).expect("acquire should succeed after clear_all");
drop(claim2);
}
}