#[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;
#[cfg(not(target_arch = "wasm32"))]
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD};
use futures::Stream;
use futures::task::{Context, Poll};
use meerkat_core::agent::CommsRuntime as CoreCommsRuntime;
use meerkat_core::comms::{
CommsCommand, EventStream, InputStreamMode, PeerAddress, PeerCapabilitySet, PeerDirectoryEntry,
PeerDirectorySource, PeerId, PeerName, PeerReachabilityReason, PeerRoute, PeerSendability,
SendAndStreamError, SendError, SendReceipt, StreamError, StreamScope, TrustedPeerDescriptor,
};
use meerkat_core::config::PlainEventSource;
#[cfg(not(target_arch = "wasm32"))]
use meerkat_core::handles::{SessionClaim, SessionClaimError, SessionClaimHandle};
use meerkat_core::hydrate_content_blocks;
use meerkat_core::time_compat::Instant;
use meerkat_core::{BlobStore, MissingBlobBehavior};
use parking_lot::Mutex;
use sha2::{Digest, Sha256};
use std::collections::{HashMap, HashSet};
#[cfg(not(target_arch = "wasm32"))]
use std::path::Path;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use thiserror::Error;
#[cfg(not(target_arch = "wasm32"))]
use tokio::net::TcpListener;
use tokio::sync::Mutex as AsyncMutex;
use tokio::sync::mpsc;
use tokio::sync::mpsc::Receiver;
#[cfg(not(target_arch = "wasm32"))]
use tokio::task::JoinHandle;
#[cfg(not(target_arch = "wasm32"))]
use tokio::{
io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
net::TcpStream,
};
use uuid::Uuid;
const RESERVATION_TTL: meerkat_core::time_compat::Duration =
meerkat_core::time_compat::Duration::from_secs(30);
#[derive(Debug)]
struct StreamRegistryEntry {
_sender: mpsc::Sender<meerkat_core::AgentEvent>,
receiver: Option<Receiver<meerkat_core::AgentEvent>>,
created_at: Instant,
lifecycle_authority: InteractionStreamLifecycleAuthority,
}
impl StreamRegistryEntry {
fn new(
sender: mpsc::Sender<meerkat_core::AgentEvent>,
receiver: Receiver<meerkat_core::AgentEvent>,
lifecycle_authority: InteractionStreamLifecycleAuthority,
) -> Self {
Self {
_sender: sender,
receiver: Some(receiver),
created_at: Instant::now(),
lifecycle_authority,
}
}
}
type InteractionStreamRegistry = Arc<Mutex<HashMap<Uuid, StreamRegistryEntry>>>;
type PeerRequestResponseAuthorityHandles = (
Arc<dyn meerkat_core::handles::PeerInteractionHandle>,
Arc<dyn meerkat_core::handles::InteractionStreamHandle>,
);
#[derive(Clone)]
pub struct PeerRequestResponseAuthority {
peer_interaction: Arc<dyn meerkat_core::handles::PeerInteractionHandle>,
interaction_stream: Arc<dyn meerkat_core::handles::InteractionStreamHandle>,
}
impl PeerRequestResponseAuthority {
pub fn new(
peer_interaction: Arc<dyn meerkat_core::handles::PeerInteractionHandle>,
interaction_stream: Arc<dyn meerkat_core::handles::InteractionStreamHandle>,
) -> Self {
Self {
peer_interaction,
interaction_stream,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum InteractionStreamLifecycleAuthority {
Machine,
LocalTransportOnly,
}
enum InteractionStreamReservationAuthority {
Machine(Arc<dyn meerkat_core::handles::InteractionStreamHandle>),
LocalTransportOnly,
}
struct InteractionStream {
id: Uuid,
receiver: Option<Receiver<meerkat_core::AgentEvent>>,
stream_handle: Option<Arc<dyn meerkat_core::handles::InteractionStreamHandle>>,
registry: InteractionStreamRegistry,
source: meerkat_core::EventSourceIdentity,
seq: u64,
}
struct ResolvedPeer {
name: PeerName,
peer_id: meerkat_core::comms::PeerId,
address: meerkat_core::comms::PeerAddress,
source: PeerDirectorySource,
meta: crate::PeerMeta,
}
impl ResolvedPeer {
fn reachability_key(&self) -> ReachabilityKey {
ReachabilityKey::new(self.name.as_str(), self.peer_id.as_str())
}
}
#[cfg(not(target_arch = "wasm32"))]
const PAIRING_VERSION: u32 = 1;
#[cfg(not(target_arch = "wasm32"))]
const PAIRING_HELLO_KIND: &str = "meerkat_pairing_hello";
#[cfg(not(target_arch = "wasm32"))]
const PAIRING_CHALLENGE_KIND: &str = "meerkat_pairing_challenge";
#[cfg(not(target_arch = "wasm32"))]
const PAIRING_PROOF_KIND: &str = "meerkat_pairing_proof";
#[cfg(not(target_arch = "wasm32"))]
const PAIRING_COMPLETE_KIND: &str = "meerkat_pairing_complete";
#[cfg(not(target_arch = "wasm32"))]
const PAIRING_CLASSIFY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, serde::Deserialize)]
struct PairingPeerIdentity {
kind: String,
public_key: String,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, serde::Deserialize)]
struct PairingPeer {
name: String,
address: String,
identity: PairingPeerIdentity,
}
#[cfg(not(target_arch = "wasm32"))]
fn comms_pairing_password_proof(
password: &str,
challenge: &str,
caller_public_key: &str,
caller_address: &str,
) -> String {
let mut hasher = Sha256::new();
hasher.update(b"meerkat-comms-pairing-v1");
hasher.update([0]);
hasher.update(password.as_bytes());
hasher.update([0]);
hasher.update(challenge.as_bytes());
hasher.update([0]);
hasher.update(caller_public_key.as_bytes());
hasher.update([0]);
hasher.update(caller_address.as_bytes());
BASE64_STANDARD.encode(hasher.finalize())
}
#[cfg(not(target_arch = "wasm32"))]
fn constant_time_str_eq(left: &str, right: &str) -> bool {
let left = left.as_bytes();
let right = right.as_bytes();
let mut diff = left.len() ^ right.len();
let max = left.len().max(right.len());
for idx in 0..max {
let left_byte = left.get(idx).copied().unwrap_or_default();
let right_byte = right.get(idx).copied().unwrap_or_default();
diff |= usize::from(left_byte ^ right_byte);
}
diff == 0
}
#[cfg(not(target_arch = "wasm32"))]
fn validate_pairing_secret(secret: &str) -> Result<(), std::io::Error> {
if secret.len() < 32 {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"comms pairing secret must be at least 32 bytes; use a random one-time bootstrap token",
));
}
Ok(())
}
fn peer_id_from_pubkey(pubkey: &crate::identity::PubKey) -> meerkat_core::comms::PeerId {
crate::router::peer_id_from_pubkey(pubkey)
}
fn parse_peer_address(raw: &str) -> Result<meerkat_core::comms::PeerAddress, String> {
meerkat_core::comms::PeerAddress::parse(raw).map_err(|err| err.to_string())
}
fn descriptor_for_inproc_peer(name: &str, pubkey: PubKey) -> Result<TrustedPeerDescriptor, String> {
let peer_name = PeerName::new(name.to_string())?;
Ok(TrustedPeerDescriptor {
peer_id: crate::router::peer_id_from_pubkey(&pubkey),
name: peer_name,
address: PeerAddress::new(meerkat_core::comms::PeerTransport::Inproc, name),
pubkey: *pubkey.as_bytes(),
})
}
fn descriptor_to_trusted_peer(descriptor: TrustedPeerDescriptor) -> Result<TrustedPeer, SendError> {
let pubkey = PubKey::new(descriptor.pubkey);
if descriptor.has_zero_pubkey() {
return Err(SendError::Validation(
"TrustedPeerDescriptor.pubkey must be non-zero for trust registration".to_string(),
));
}
let derived = pubkey.to_peer_id();
if derived != descriptor.peer_id {
return Err(SendError::Validation(format!(
"TrustedPeerDescriptor.peer_id {} does not match pubkey-derived id {}",
descriptor.peer_id, derived
)));
}
Ok(TrustedPeer {
name: descriptor.name.as_string(),
pubkey,
addr: descriptor.address.to_string(),
meta: crate::PeerMeta::default(),
})
}
impl InteractionStream {
fn finish(&mut self) {
if let Some(mut receiver) = self.receiver.take() {
receiver.close();
}
match self.stream_handle.as_ref() {
Some(handle) => {
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(self.id);
if let Err(err) = handle.closed_early(corr_id) {
tracing::trace!(
interaction_id = %self.id,
error = %err,
"InteractionStreamHandle::closed_early rejected (likely terminal won race)"
);
}
}
None => {
self.registry.lock().remove(&self.id);
}
}
}
}
fn map_peer_ingress_phase(
phase: crate::peer_types::PeerIngressState,
) -> meerkat_core::PeerIngressAuthorityPhase {
match phase {
crate::peer_types::PeerIngressState::Absent => {
meerkat_core::PeerIngressAuthorityPhase::Absent
}
crate::peer_types::PeerIngressState::Received => {
meerkat_core::PeerIngressAuthorityPhase::Received
}
crate::peer_types::PeerIngressState::Dropped => {
meerkat_core::PeerIngressAuthorityPhase::Dropped
}
crate::peer_types::PeerIngressState::Delivered => {
meerkat_core::PeerIngressAuthorityPhase::Delivered
}
}
}
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_with_source(
this.source.clone(),
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 peer_id(&self) -> Option<PeerId> {
Some(self.public_key.to_peer_id())
}
fn public_key(&self) -> Option<String> {
Some(self.public_key.to_pubkey_string())
}
fn public_key_bytes(&self) -> Option<[u8; 32]> {
Some(*self.public_key.as_bytes())
}
fn comms_name(&self) -> Option<String> {
Some(self.participant_name().to_string())
}
fn advertised_address(&self) -> Option<String> {
#[cfg(not(target_arch = "wasm32"))]
{
Some(crate::runtime::comms_runtime::CommsRuntime::advertised_address(self))
}
#[cfg(target_arch = "wasm32")]
{
None
}
}
fn bridge_bootstrap_token(&self) -> Option<String> {
Some(self.bridge_bootstrap_token.clone())
}
async fn add_trusted_peer(&self, peer: TrustedPeerDescriptor) -> Result<(), SendError> {
self.register_trusted_peer_descriptor(peer).await
}
async fn remove_trusted_peer(&self, peer_id: &str) -> Result<bool, SendError> {
self.unregister_trusted_peer(peer_id).await
}
async fn add_private_trusted_peer(&self, peer: TrustedPeerDescriptor) -> Result<(), SendError> {
self.register_private_trusted_peer_descriptor(peer).await
}
async fn remove_private_trusted_peer(&self, peer_id: &str) -> Result<bool, SendError> {
self.unregister_private_trusted_peer(peer_id).await
}
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 corr_id = meerkat_core::PeerCorrelationId::from_uuid(id);
let (expired, lifecycle_authority) = {
let registry = self.interaction_stream_registry.lock();
match registry.get(&id) {
Some(entry) => (
entry.created_at.elapsed() > RESERVATION_TTL,
entry.lifecycle_authority,
),
None => return Err(StreamError::NotReserved(interaction_id)),
}
};
if expired {
match lifecycle_authority {
InteractionStreamLifecycleAuthority::Machine => {
if let Some(handle) = self.interaction_stream_handle() {
let _ = handle.expired(corr_id);
}
}
InteractionStreamLifecycleAuthority::LocalTransportOnly => {
self.interaction_stream_registry.lock().remove(&id);
}
}
if let Some(handle) = self.peer_interaction_handle()
&& handle.outbound_state(corr_id).is_some()
{
let _ = handle.request_timed_out(corr_id);
}
return Err(StreamError::Timeout(format!(
"reservation expired for interaction {}",
interaction_id.0
)));
}
let stream_handle = match lifecycle_authority {
InteractionStreamLifecycleAuthority::Machine => {
Some(self.interaction_stream_handle().ok_or_else(|| {
StreamError::Internal(
"machine interaction stream authority missing for reserved stream"
.to_string(),
)
})?)
}
InteractionStreamLifecycleAuthority::LocalTransportOnly => None,
};
if let Some(handle) = stream_handle.as_ref() {
handle.attached(corr_id).map_err(|err| {
match handle.state(corr_id) {
Some(meerkat_core::InteractionStreamState::Attached) => {
StreamError::AlreadyAttached(interaction_id)
}
None => StreamError::NotReserved(interaction_id),
_ => StreamError::Internal(format!(
"unexpected interaction stream state for {}: {err}",
interaction_id.0
)),
}
})?;
}
let mut registry = self.interaction_stream_registry.lock();
let entry = registry
.get_mut(&id)
.ok_or(StreamError::NotReserved(interaction_id))?;
let receiver = entry.receiver.take().ok_or_else(|| {
if lifecycle_authority == InteractionStreamLifecycleAuthority::Machine {
StreamError::Internal("interaction stream receiver missing".to_string())
} else {
StreamError::AlreadyAttached(interaction_id)
}
})?;
drop(registry);
Ok(Box::pin(InteractionStream {
id,
receiver: Some(receiver),
stream_handle,
registry: self.interaction_stream_registry.clone(),
source: meerkat_core::EventSourceIdentity::interaction(
meerkat_core::InteractionId(id),
),
seq: 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,
handling_mode,
} => self
.send_peer_command(
&to,
crate::types::MessageKind::Message {
body,
blocks,
handling_mode: Some(handling_mode),
},
)
.await
.map(|envelope_id| SendReceipt::PeerMessageSent {
envelope_id,
acked: false,
}),
CommsCommand::PeerLifecycle { to, kind, params } => self
.send_peer_command(&to, crate::types::MessageKind::Lifecycle { kind, params })
.await
.map(|envelope_id| SendReceipt::PeerLifecycleSent { envelope_id }),
CommsCommand::PeerRequest {
to,
intent,
params,
blocks,
handling_mode,
stream,
} => {
let interaction_id = Uuid::new_v4();
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(interaction_id);
let stream_reserved = stream == InputStreamMode::ReserveInteraction;
let (peer_handle, stream_authority) =
self.require_peer_request_response_authority("PeerRequest")?;
let stream_handle = if stream_reserved {
Some(stream_authority)
} else {
None
};
if let Err(err) = peer_handle.request_sent(corr_id, to.peer_id.to_string()) {
tracing::warn!(
error = %err,
corr_id = %corr_id,
"PeerInteractionHandle::request_sent rejected — refusing to send"
);
return Err(SendError::Validation(format!(
"DSL rejected PeerRequestSent for corr_id {corr_id}: {err}"
)));
}
if let Some(stream_handle) = stream_handle
&& let Err(err) = self.reserve_interaction_stream_channels(
corr_id,
interaction_id,
InteractionStreamReservationAuthority::Machine(stream_handle),
)
{
let _ = peer_handle.request_timed_out(corr_id);
return Err(err);
}
let envelope_id = match self
.send_peer_command_with_id(
&to,
interaction_id,
crate::types::MessageKind::Request {
intent,
params,
blocks,
handling_mode: Some(handling_mode),
},
)
.await
{
Ok(id) => id,
Err(e) => {
if stream_reserved {
self.expire_interaction_stream_on_send_failure(corr_id, interaction_id);
}
let _ = peer_handle.request_timed_out(corr_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,
blocks,
handling_mode,
} => {
let (peer_handle, _stream_authority) =
self.require_peer_request_response_authority("PeerResponse")?;
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(in_reply_to.0);
let core_status = status;
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,
};
if peer_handle.inbound_state(corr_id).is_none() {
return Err(SendError::Validation(format!(
"PeerResponse requires machine inbound peer request state for corr_id {corr_id}"
)));
}
let is_terminal_reply = matches!(
core_status,
meerkat_core::ResponseStatus::Completed | meerkat_core::ResponseStatus::Failed
);
let effective_handling_mode = handling_mode.or_else(|| {
is_terminal_reply
.then(|| {
self.inbound_request_handling_modes
.lock()
.get(&corr_id)
.copied()
})
.flatten()
});
let envelope_id = self
.send_peer_command(
&to,
crate::types::MessageKind::Response {
in_reply_to: in_reply_to.0,
status,
result,
blocks,
handling_mode: effective_handling_mode,
},
)
.await?;
if is_terminal_reply && let Err(err) = peer_handle.response_replied(corr_id) {
tracing::warn!(
error = %err,
corr_id = %corr_id,
"PeerInteractionHandle::response_replied rejected"
);
return Err(SendError::Validation(format!(
"DSL rejected PeerResponseReplied for corr_id {corr_id}: {err}"
)));
}
if is_terminal_reply {
self.inbound_request_handling_modes.lock().remove(&corr_id);
}
Ok(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,
blocks,
handling_mode,
stream: InputStreamMode::ReserveInteraction,
} => {
let interaction_id = Uuid::new_v4();
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(interaction_id);
let (peer_handle, stream_handle) =
self.require_peer_request_response_authority("PeerRequest")?;
if let Err(err) = peer_handle.request_sent(corr_id, to.peer_id.to_string()) {
tracing::warn!(
error = %err,
corr_id = %corr_id,
"PeerInteractionHandle::request_sent rejected — refusing to send"
);
return Err(SendAndStreamError::Send(SendError::Validation(format!(
"DSL rejected PeerRequestSent for corr_id {corr_id}: {err}"
))));
}
if let Err(err) = self.reserve_interaction_stream_channels(
corr_id,
interaction_id,
InteractionStreamReservationAuthority::Machine(stream_handle),
) {
let _ = peer_handle.request_timed_out(corr_id);
return Err(SendAndStreamError::Send(err));
}
let envelope_id = match self
.send_peer_command_with_id(
&to,
interaction_id,
crate::types::MessageKind::Request {
intent,
params,
blocks,
handling_mode: Some(handling_mode),
},
)
.await
{
Ok(id) => id,
Err(e) => {
self.expire_interaction_stream_on_send_failure(corr_id, interaction_id);
let _ = peer_handle.request_timed_out(corr_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 corr_id = meerkat_core::PeerCorrelationId::from_uuid(id.0);
let lifecycle_authority = self
.interaction_stream_registry
.lock()
.get(&id.0)
.map(|entry| entry.lifecycle_authority);
match lifecycle_authority {
Some(InteractionStreamLifecycleAuthority::Machine) => {
if let Some(handle) = self.interaction_stream_handle()
&& handle.state(corr_id) == Some(meerkat_core::InteractionStreamState::Reserved)
{
let _ = handle.expired(corr_id);
}
}
Some(InteractionStreamLifecycleAuthority::LocalTransportOnly) => {
let mut registry = self.interaction_stream_registry.lock();
if let Some(entry) = registry.get(&id.0)
&& entry.receiver.is_some()
{
registry.remove(&id.0);
}
}
None => {}
}
sender
}
fn mark_interaction_complete(&self, id: &meerkat_core::InteractionId) {
self.mark_interaction_complete(id.0);
}
fn peer_interaction_handle(
&self,
) -> Option<Arc<dyn meerkat_core::handles::PeerInteractionHandle>> {
self.peer_interaction_handle()
}
fn peer_request_response_authority_handle(
&self,
) -> Option<Arc<dyn meerkat_core::handles::PeerInteractionHandle>> {
let peer_handle = self.peer_interaction_handle();
if peer_handle.is_some() && self.interaction_stream_handle().is_some() {
peer_handle
} else {
None
}
}
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 ingress = entry.ingress_fact;
let lifecycle_peer = entry.lifecycle_peer;
let from_peer = entry.from_peer.unwrap_or_else(|| "unknown".to_string());
let from_peer_id = ingress.route.as_ref().map(|route| route.peer_id);
let rendered_text = entry.text_projection.clone();
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 envelope_handling_mode = match &envelope.kind {
MessageKind::Message {
handling_mode: Some(mode),
..
}
| MessageKind::Request {
handling_mode: Some(mode),
..
}
| MessageKind::Response {
handling_mode: Some(mode),
..
} => *mode,
MessageKind::Response {
handling_mode: None,
..
}
| MessageKind::Message {
handling_mode: None,
..
}
| MessageKind::Request {
handling_mode: None,
..
}
| MessageKind::Lifecycle { .. }
| MessageKind::Ack { .. } => meerkat_core::types::HandlingMode::Queue,
};
let is_peer_request_envelope =
matches!(&envelope.kind, MessageKind::Request { .. });
let content = match envelope.kind {
MessageKind::Message {
body,
blocks,
handling_mode: _,
} => meerkat_core::InteractionContent::Message { body, blocks },
MessageKind::Request {
intent,
params,
blocks,
handling_mode: _,
} => {
let typed_intent = MessageIntent::from(intent.as_str());
meerkat_core::InteractionContent::Request {
intent: typed_intent.to_string(),
params,
blocks,
}
}
MessageKind::Lifecycle { kind, params } => {
meerkat_core::InteractionContent::Request {
intent: kind.to_string(),
params,
blocks: None,
}
}
MessageKind::Response {
in_reply_to,
status,
result,
blocks,
handling_mode: _,
} => {
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
}
};
meerkat_core::InteractionContent::Response {
in_reply_to: meerkat_core::InteractionId(in_reply_to),
status: core_status,
result,
blocks,
}
}
MessageKind::Ack { .. } => {
return None;
}
};
if is_peer_request_envelope {
self.inbound_request_handling_modes.lock().insert(
meerkat_core::PeerCorrelationId::from_uuid(envelope.id),
envelope_handling_mode,
);
}
Some(meerkat_core::ClassifiedInboxInteraction {
interaction: meerkat_core::InboxInteraction {
id: meerkat_core::InteractionId(envelope.id),
from_route: from_peer_id,
from: from_peer,
content,
rendered_text,
handling_mode: envelope_handling_mode,
render_metadata: None,
},
ingress,
lifecycle_peer,
response_terminality: entry.response_terminality,
})
}
crate::types::InboxItem::PlainEvent {
body,
source,
handling_mode,
interaction_id,
blocks,
render_metadata,
..
} => Some(meerkat_core::ClassifiedInboxInteraction {
interaction: meerkat_core::InboxInteraction {
id: meerkat_core::InteractionId(
interaction_id.unwrap_or_else(uuid::Uuid::new_v4),
),
from_route: None,
from: format!("event:{source}"),
content: meerkat_core::InteractionContent::Message { body, blocks },
rendered_text,
handling_mode,
render_metadata,
},
ingress,
lifecycle_peer,
response_terminality: entry.response_terminality,
}),
}
})
.collect())
}
async fn peer_ingress_queue_snapshot(
&self,
) -> Result<meerkat_core::PeerIngressQueueSnapshot, meerkat_core::CommsCapabilityError> {
let inbox = self.inbox.lock().await;
inbox.classified_snapshot().ok_or_else(|| {
meerkat_core::CommsCapabilityError::Unsupported(
"peer_ingress_queue_snapshot: classified inbox not initialized".to_string(),
)
})
}
async fn peer_ingress_runtime_snapshot(
&self,
) -> Result<meerkat_core::PeerIngressRuntimeSnapshot, meerkat_core::CommsCapabilityError> {
let (queue, authority_phase, submission_queue_len) = {
let inbox = self.inbox.lock().await;
inbox.peer_runtime_snapshot().ok_or_else(|| {
meerkat_core::CommsCapabilityError::Unsupported(
"peer_ingress_runtime_snapshot: classified inbox not initialized".to_string(),
)
})?
};
let mut trusted_peers: Vec<TrustedPeerDescriptor> = self
.trusted_peers
.read()
.peers
.iter()
.filter_map(|peer| {
if peer.pubkey.is_zero() {
tracing::warn!(
peer_name = %peer.name,
"skipping zero-pubkey trusted peer in ingress runtime snapshot"
);
return None;
}
let name = match PeerName::new(peer.name.clone()) {
Ok(name) => name,
Err(err) => {
tracing::warn!(
peer_name = %peer.name,
error = %err,
"skipping trusted peer with invalid name in snapshot"
);
return None;
}
};
let address = match parse_peer_address(&peer.addr) {
Ok(address) => address,
Err(err) => {
tracing::warn!(
peer_name = %peer.name,
peer_id = %peer.pubkey.to_peer_id(),
address = %peer.addr,
error = %err,
"skipping trusted peer with invalid address in snapshot"
);
return None;
}
};
Some(TrustedPeerDescriptor {
peer_id: self.router.peer_id_for_pubkey(&peer.pubkey),
name,
address,
pubkey: *peer.pubkey.as_bytes(),
})
})
.collect();
trusted_peers.sort_by(|left, right| {
left.name
.as_str()
.cmp(right.name.as_str())
.then_with(|| left.peer_id.cmp(&right.peer_id))
.then_with(|| left.address.to_string().cmp(&right.address.to_string()))
});
Ok(meerkat_core::PeerIngressRuntimeSnapshot {
self_peer_id: peer_id_from_pubkey(&self.public_key),
auth_required: self.require_peer_auth,
authority_phase: map_peer_ingress_phase(authority_phase),
trusted_peers,
submission_queue_len,
queue,
})
}
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),
}
#[cfg(not(target_arch = "wasm32"))]
impl From<SessionClaimError> for CommsRuntimeError {
fn from(err: SessionClaimError) -> Self {
match err {
SessionClaimError::SessionIdentityInUse(sid) => {
Self::SessionIdentityInUse(sid.to_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: Mutex<Vec<ListenerHandle>>,
#[cfg(not(target_arch = "wasm32"))]
listeners_started: bool,
#[cfg(not(target_arch = "wasm32"))]
_session_identity_claim: Option<SessionClaim>,
keypair: Arc<Keypair>,
bridge_bootstrap_token: String,
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)] ingress_policy: Arc<meerkat_core::PeerIngressMachinePolicy>,
peer_comms_handle: crate::classify::PeerCommsHandleSlot,
require_peer_comms_machine_authority: Arc<AtomicBool>,
actionable_notify: Option<Arc<tokio::sync::Notify>>,
blob_store: Option<Arc<dyn BlobStore>>,
peer_interaction_handle:
parking_lot::RwLock<Option<Arc<dyn meerkat_core::handles::PeerInteractionHandle>>>,
interaction_stream_handle:
parking_lot::RwLock<Option<Arc<dyn meerkat_core::handles::InteractionStreamHandle>>>,
inbound_request_handling_modes:
Arc<Mutex<HashMap<meerkat_core::PeerCorrelationId, meerkat_core::types::HandlingMode>>>,
}
impl CommsRuntime {
fn ingress_policy_from_silent_intents(
silent_intents: &Arc<HashSet<String>>,
) -> Arc<meerkat_core::PeerIngressMachinePolicy> {
Arc::new(meerkat_core::PeerIngressMachinePolicy::from_silent_intents(
silent_intents.iter().cloned(),
))
}
fn derive_bridge_bootstrap_token(keypair: &Keypair) -> String {
let mut digest = Sha256::new();
digest.update(b"meerkat.supervisor-bridge.bootstrap-token.v1");
digest.update(keypair.secret_bytes());
base64::Engine::encode(
&base64::engine::general_purpose::URL_SAFE_NO_PAD,
digest.finalize(),
)
}
#[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> {
Self::new_with_silent_intents_and_machine_authority_requirement(
config,
silent_intents,
false,
)
.await
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn new_machine_authority_required_with_silent_intents(
config: ResolvedCommsConfig,
silent_intents: Arc<HashSet<String>>,
) -> Result<Self, CommsRuntimeError> {
Self::new_with_silent_intents_and_machine_authority_requirement(
config,
silent_intents,
true,
)
.await
}
#[cfg(not(target_arch = "wasm32"))]
async fn new_with_silent_intents_and_machine_authority_requirement(
config: ResolvedCommsConfig,
silent_intents: Arc<HashSet<String>>,
require_machine_authority: bool,
) -> 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 ingress_policy = Self::ingress_policy_from_silent_intents(&silent_intents);
let peer_comms_handle = Arc::new(parking_lot::RwLock::new(None));
let require_peer_comms_machine_authority =
Arc::new(AtomicBool::new(require_machine_authority));
let classification_context = Arc::new(crate::classify::IngressClassificationContext {
require_peer_auth: config.require_peer_auth,
trusted_peers: trusted_peers.clone(),
ingress_policy: ingress_policy.clone(),
peer_comms_handle: peer_comms_handle.clone(),
require_machine_authority: require_peer_comms_machine_authority.clone(),
inproc_namespace: config.inproc_namespace.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: Mutex::new(Vec::new()),
listeners_started: false,
_session_identity_claim: None,
bridge_bootstrap_token: Self::derive_bridge_bootstrap_token(&keypair),
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())),
inbound_request_handling_modes: Arc::new(Mutex::new(HashMap::new())),
peer_directory_reachability: Arc::new(Mutex::new(
PeerDirectoryReachabilityAuthority::new(),
)),
ingress_policy,
peer_comms_handle,
require_peer_comms_machine_authority,
actionable_notify,
blob_store: None,
peer_interaction_handle: parking_lot::RwLock::new(None),
interaction_stream_handle: parking_lot::RwLock::new(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();
Self::inproc_only_with_keypair_and_silent_intents(name, namespace, keypair, silent_intents)
}
pub fn inproc_only_with_keypair_and_silent_intents(
name: &str,
namespace: Option<String>,
keypair: Keypair,
silent_intents: Arc<HashSet<String>>,
) -> Result<Self, CommsRuntimeError> {
let public_key = keypair.public_key();
let trusted_peers = Arc::new(parking_lot::RwLock::new(TrustedPeers::new()));
let ingress_policy = Self::ingress_policy_from_silent_intents(&silent_intents);
let peer_comms_handle = Arc::new(parking_lot::RwLock::new(None));
let require_peer_comms_machine_authority = Arc::new(AtomicBool::new(false));
let classification_context = Arc::new(crate::classify::IngressClassificationContext {
require_peer_auth: true,
trusted_peers: trusted_peers.clone(),
ingress_policy: ingress_policy.clone(),
peer_comms_handle: peer_comms_handle.clone(),
require_machine_authority: require_peer_comms_machine_authority.clone(),
inproc_namespace: namespace.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,
advertise_address: 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,
pairing_password: None,
};
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: Mutex::new(Vec::new()),
#[cfg(not(target_arch = "wasm32"))]
listeners_started: false,
#[cfg(not(target_arch = "wasm32"))]
_session_identity_claim: None,
bridge_bootstrap_token: Self::derive_bridge_bootstrap_token(&keypair),
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())),
inbound_request_handling_modes: Arc::new(Mutex::new(HashMap::new())),
peer_directory_reachability: Arc::new(Mutex::new(
PeerDirectoryReachabilityAuthority::new(),
)),
ingress_policy,
peer_comms_handle,
require_peer_comms_machine_authority,
actionable_notify,
blob_store: None,
peer_interaction_handle: parking_lot::RwLock::new(None),
interaction_stream_handle: parking_lot::RwLock::new(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 fn set_listen_tcp_for_unstarted_runtime(
&mut self,
listen_tcp: std::net::SocketAddr,
) -> Result<(), CommsRuntimeError> {
if self.listeners_started {
return Err(CommsRuntimeError::AlreadyStarted);
}
self.config.listen_tcp = Some(listen_tcp);
Ok(())
}
#[cfg(not(target_arch = "wasm32"))]
pub fn set_advertise_address_for_unstarted_runtime(
&mut self,
advertise_address: String,
) -> Result<(), CommsRuntimeError> {
if self.listeners_started {
return Err(CommsRuntimeError::AlreadyStarted);
}
self.config.advertise_address = Some(advertise_address);
Ok(())
}
pub async fn inproc_pair_with_mutual_trust(
name_a: &str,
name_b: &str,
) -> Result<(Arc<Self>, Arc<Self>), CommsRuntimeError> {
let a = Arc::new(Self::inproc_only(name_a)?);
let b = Arc::new(Self::inproc_only(name_b)?);
let descriptor_for_a = descriptor_for_inproc_peer(name_b, b.public_key())
.map_err(CommsRuntimeError::TrustLoadError)?;
let descriptor_for_b = descriptor_for_inproc_peer(name_a, a.public_key())
.map_err(CommsRuntimeError::TrustLoadError)?;
meerkat_core::agent::CommsRuntime::add_trusted_peer(a.as_ref(), descriptor_for_a)
.await
.map_err(|err| CommsRuntimeError::TrustLoadError(err.to_string()))?;
meerkat_core::agent::CommsRuntime::add_trusted_peer(b.as_ref(), descriptor_for_b)
.await
.map_err(|err| CommsRuntimeError::TrustLoadError(err.to_string()))?;
Ok((a, b))
}
#[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>>,
session_claim_handle: Arc<dyn SessionClaimHandle>,
) -> Result<Self, CommsRuntimeError> {
let claim = session_claim_handle
.try_acquire(session_id)
.map_err(CommsRuntimeError::from)?;
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 ingress_policy = Self::ingress_policy_from_silent_intents(&silent_intents);
let peer_comms_handle = Arc::new(parking_lot::RwLock::new(None));
let require_peer_comms_machine_authority = Arc::new(AtomicBool::new(true));
let classification_context = Arc::new(crate::classify::IngressClassificationContext {
require_peer_auth: true,
trusted_peers: trusted_peers.clone(),
ingress_policy: ingress_policy.clone(),
peer_comms_handle: peer_comms_handle.clone(),
require_machine_authority: require_peer_comms_machine_authority.clone(),
inproc_namespace: namespace.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,
advertise_address: 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,
pairing_password: None,
};
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: Mutex::new(Vec::new()),
listeners_started: false,
_session_identity_claim: Some(claim),
bridge_bootstrap_token: Self::derive_bridge_bootstrap_token(&keypair),
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())),
inbound_request_handling_modes: Arc::new(Mutex::new(HashMap::new())),
peer_directory_reachability: Arc::new(Mutex::new(
PeerDirectoryReachabilityAuthority::new(),
)),
ingress_policy,
peer_comms_handle,
require_peer_comms_machine_authority,
actionable_notify,
blob_store: None,
peer_interaction_handle: parking_lot::RwLock::new(None),
interaction_stream_handle: parking_lot::RwLock::new(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);
}
pub fn install_peer_comms_handle(
&self,
handle: Arc<dyn meerkat_core::handles::PeerCommsHandle>,
) {
*self.peer_comms_handle.write() = Some(handle);
}
pub fn require_peer_comms_machine_authority(&self) {
self.require_peer_comms_machine_authority
.store(true, Ordering::SeqCst);
}
pub fn peer_comms_machine_authority_required(&self) -> bool {
self.require_peer_comms_machine_authority
.load(Ordering::SeqCst)
}
pub fn peer_comms_handle(&self) -> Option<Arc<dyn meerkat_core::handles::PeerCommsHandle>> {
self.peer_comms_handle.read().clone()
}
pub fn install_peer_request_response_authority(
self: &Arc<Self>,
authority: PeerRequestResponseAuthority,
) {
self.install_peer_interaction_handle(authority.peer_interaction);
self.install_interaction_stream_handle(authority.interaction_stream);
}
pub fn install_peer_interaction_handle(
self: &Arc<Self>,
handle: Arc<dyn meerkat_core::handles::PeerInteractionHandle>,
) {
let observer: Arc<dyn meerkat_core::handles::PeerInteractionCleanupObserver> =
Arc::clone(self) as Arc<dyn meerkat_core::handles::PeerInteractionCleanupObserver>;
handle.install_cleanup_observer(observer);
*self.peer_interaction_handle.write() = Some(handle);
}
pub fn peer_interaction_handle(
&self,
) -> Option<Arc<dyn meerkat_core::handles::PeerInteractionHandle>> {
self.peer_interaction_handle.read().clone()
}
pub fn install_interaction_stream_handle(
self: &Arc<Self>,
handle: Arc<dyn meerkat_core::handles::InteractionStreamHandle>,
) {
let observer: Arc<dyn meerkat_core::handles::InteractionStreamCleanupObserver> =
Arc::clone(self) as Arc<dyn meerkat_core::handles::InteractionStreamCleanupObserver>;
handle.install_cleanup_observer(observer);
*self.interaction_stream_handle.write() = Some(handle);
}
pub fn interaction_stream_handle(
&self,
) -> Option<Arc<dyn meerkat_core::handles::InteractionStreamHandle>> {
self.interaction_stream_handle.read().clone()
}
fn require_peer_interaction_authority(
&self,
command: &'static str,
) -> Result<Arc<dyn meerkat_core::handles::PeerInteractionHandle>, SendError> {
self.peer_interaction_handle().ok_or_else(|| {
SendError::Validation(format!(
"{command} requires machine peer interaction authority; this CommsRuntime is transport-only for semantic peer request/response"
))
})
}
fn require_interaction_stream_authority(
&self,
command: &'static str,
) -> Result<Arc<dyn meerkat_core::handles::InteractionStreamHandle>, SendError> {
self.interaction_stream_handle().ok_or_else(|| {
SendError::Validation(format!(
"{command} requires machine interaction stream authority; this CommsRuntime is transport-only for semantic peer request/response streams"
))
})
}
fn require_peer_request_response_authority(
&self,
command: &'static str,
) -> Result<PeerRequestResponseAuthorityHandles, SendError> {
let peer_handle = self.require_peer_interaction_authority(command)?;
let stream_handle = self.require_interaction_stream_authority(command)?;
Ok((peer_handle, stream_handle))
}
fn peer_request_response_sendable_kinds(&self) -> Vec<PeerSendability> {
let mut kinds = vec![PeerSendability::PeerMessage];
if self.peer_interaction_handle().is_some() && self.interaction_stream_handle().is_some() {
kinds.push(PeerSendability::PeerRequest);
kinds.push(PeerSendability::PeerResponse);
}
kinds
}
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),
handling_mode,
} => {
self.hydrate_blocks_for_transport(&mut blocks, "comms message")
.await?;
Ok(crate::types::MessageKind::Message {
body,
blocks: Some(blocks),
handling_mode,
})
}
crate::types::MessageKind::Request {
intent,
params,
blocks: Some(mut blocks),
handling_mode,
} => {
self.hydrate_blocks_for_transport(&mut blocks, "comms request")
.await?;
Ok(crate::types::MessageKind::Request {
intent,
params,
blocks: Some(blocks),
handling_mode,
})
}
crate::types::MessageKind::Response {
in_reply_to,
status,
result,
blocks: Some(mut blocks),
handling_mode,
} => {
self.hydrate_blocks_for_transport(&mut blocks, "comms response")
.await?;
Ok(crate::types::MessageKind::Response {
in_reply_to,
status,
result,
blocks: Some(blocks),
handling_mode,
})
}
other => Ok(other),
}
}
async fn hydrate_blocks_for_transport(
&self,
blocks: &mut [meerkat_core::types::ContentBlock],
label: &str,
) -> Result<(), SendError> {
if !blocks.iter().any(|block| {
matches!(
block,
meerkat_core::types::ContentBlock::Image {
data: meerkat_core::types::ImageData::Blob { .. },
..
}
)
}) {
return Ok(());
}
let blob_store = self.blob_store.as_ref().ok_or_else(|| {
SendError::Internal(format!("blob-backed {label} requires blob store"))
})?;
hydrate_content_blocks(blob_store.as_ref(), blocks, MissingBlobBehavior::Error)
.await
.map_err(|err| SendError::Internal(err.to_string()))
}
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 advertised) = self.config.advertise_address {
return advertised.clone();
}
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 router = self.router.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.lock().push(handle);
}
if let Some(ref addr) = self.config.listen_tcp {
let handle = spawn_tcp_listener(TcpListenerConfig {
addr: addr.to_string(),
keypair: self.keypair.clone(),
router,
trusted: self.trusted_peers.clone(),
require_peer_auth: self.require_peer_auth,
inbox_sender: inbox_sender.clone(),
pairing_password: self.config.pairing_password.clone(),
advertise_address: self.config.advertise_address.clone(),
trusted_peers_path: self.config.trusted_peers_path.clone(),
participant_name: self.config.name.clone(),
bridge_bootstrap_token: self.bridge_bootstrap_token.clone(),
})
.await?;
if let Some(local_addr) = handle.local_addr {
self.config.listen_tcp = Some(local_addr);
}
self.listener_handles.lock().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.lock().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.lock().push(handle);
}
}
self.listeners_started = true;
Ok(())
}
pub fn public_key(&self) -> PubKey {
self.public_key
}
pub fn bridge_bootstrap_token(&self) -> &str {
&self.bridge_bootstrap_token
}
pub fn router(&self) -> &Router {
&self.router
}
pub async fn register_trusted_peer(&self, mut peer: TrustedPeer) -> Result<(), SendError> {
peer.addr = parse_peer_address(&peer.addr)
.map_err(SendError::Validation)?
.to_string();
self.router
.add_trusted_peer(peer)
.map_err(|err| SendError::Validation(err.to_string()))?;
Ok(())
}
async fn register_trusted_peer_descriptor(
&self,
descriptor: TrustedPeerDescriptor,
) -> Result<(), SendError> {
let peer_id = descriptor.peer_id;
let peer = descriptor_to_trusted_peer(descriptor)?;
self.router
.add_trusted_peer_with_peer_id(peer_id, peer)
.map_err(|err| SendError::Validation(err.to_string()))?;
Ok(())
}
pub async fn unregister_trusted_peer(&self, peer_id: &str) -> Result<bool, SendError> {
let peer_id = meerkat_core::comms::PeerId::parse(peer_id)
.map_err(|err| SendError::Validation(err.to_string()))?;
Ok(self.router.remove_trusted_peer(&peer_id))
}
pub async fn register_private_trusted_peer(
&self,
mut peer: TrustedPeer,
) -> Result<(), SendError> {
peer.addr = parse_peer_address(&peer.addr)
.map_err(SendError::Validation)?
.to_string();
let pubkey = peer.pubkey;
self.router
.add_trusted_peer(peer)
.map_err(|err| SendError::Validation(err.to_string()))?;
self.router.mark_private(pubkey);
Ok(())
}
async fn register_private_trusted_peer_descriptor(
&self,
descriptor: TrustedPeerDescriptor,
) -> Result<(), SendError> {
let peer_id = descriptor.peer_id;
let peer = descriptor_to_trusted_peer(descriptor)?;
self.router
.add_trusted_peer_with_peer_id(peer_id, peer)
.map_err(|err| SendError::Validation(err.to_string()))?;
self.router.mark_private_peer_id(peer_id);
Ok(())
}
pub async fn unregister_private_trusted_peer(&self, peer_id: &str) -> Result<bool, SendError> {
let peer_id = meerkat_core::comms::PeerId::parse(peer_id)
.map_err(|err| SendError::Validation(err.to_string()))?;
let removed = self.router.remove_trusted_peer(&peer_id);
self.router.unmark_private(&peer_id);
Ok(removed)
}
pub async fn unregister_trusted_pubkey(&self, public_key: &PubKey) -> Result<bool, SendError> {
self.unregister_trusted_peer(&public_key.to_peer_id().to_string())
.await
}
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 = self.peer_request_response_sendable_kinds();
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: PeerCapabilitySet::default(),
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 registry = InprocRegistry::global();
let inproc_peers = registry.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 private_peer_ids = self.router.private_peer_ids();
let mut trusted_names = HashSet::new();
let mut trusted_pubkeys = HashSet::new();
let mut emitted = 0usize;
{
let trusted = self.trusted_peers.read();
let trusted_peer_id_counts: HashMap<PeerId, usize> = trusted
.peers
.iter()
.filter(|peer| peer.has_raw_sendable_identity())
.fold(HashMap::new(), |mut counts, peer| {
*counts
.entry(self.router.peer_id_for_pubkey(&peer.pubkey))
.or_default() += 1;
counts
});
for peer in &trusted.peers {
if !peer.has_raw_sendable_identity() {
tracing::warn!(
peer_name = %peer.name,
"skipping zero-pubkey trusted peer in peer directory"
);
continue;
}
if peer.name == participant_name || peer.pubkey == self.public_key {
continue;
}
let peer_id = self.router.peer_id_for_pubkey(&peer.pubkey);
trusted_names.insert(peer.name.clone());
trusted_pubkeys.insert(peer.pubkey);
if trusted_peer_id_counts.get(&peer_id).copied().unwrap_or(0) != 1 {
tracing::warn!(
peer_name = %peer.name,
peer_id = %peer_id,
"skipping duplicate trusted peer id in peer directory"
);
continue;
}
if private_peer_ids.contains(&peer_id) {
continue;
}
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;
}
};
let address = match parse_peer_address(&peer.addr) {
Ok(address) => address,
Err(err) => {
tracing::warn!(
peer_name = %peer.name,
peer_id = %peer_id,
address = %peer.addr,
error = %err,
"skipping trusted peer with invalid address in peer directory"
);
continue;
}
};
on_peer(ResolvedPeer {
name,
peer_id,
address,
source,
meta: peer.meta.clone(),
});
emitted += 1;
}
}
if self.require_peer_auth {
return emitted;
}
for inproc in &inproc_peers {
if inproc.pubkey.is_zero() {
continue;
}
let peer_id = peer_id_from_pubkey(&inproc.pubkey);
if registry.pubkey_registration_count_any_namespace(&inproc.pubkey) != 1 {
tracing::warn!(
peer_name = %inproc.name,
peer_id = %peer_id,
"skipping duplicate inproc peer id in auth-disabled peer directory"
);
continue;
}
if private_peer_ids.contains(&peer_id) {
continue;
}
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,
address: meerkat_core::comms::PeerAddress::new(
meerkat_core::comms::PeerTransport::Inproc,
peer_name_str.clone(),
),
source: PeerDirectorySource::Inproc,
meta: inproc.meta.clone(),
});
emitted += 1;
}
emitted
}
async fn send_peer_command(
&self,
route: &PeerRoute,
kind: crate::types::MessageKind,
) -> Result<Uuid, SendError> {
self.send_peer_command_with_id(route, Uuid::new_v4(), kind)
.await
}
async fn send_peer_command_with_id(
&self,
route: &PeerRoute,
envelope_id: Uuid,
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.peer_id == route.peer_id);
let dest_peer_id = route.peer_id;
let kind = self.hydrate_message_kind_for_transport(kind).await?;
let result = self
.router
.send_with_id(dest_peer_id, envelope_id, 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_id)) => {
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_id.to_string()))
}
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(crate::router::SendError::AdmissionDropped { reason }) => {
if let Some(peer) = resolved_peer.as_ref() {
self.peer_directory_reachability.lock().record_send_failed(
&peer.reachability_key(),
PeerReachabilityReason::AdmissionDropped,
);
}
Err(SendError::AdmissionDropped {
reason: reason.into(),
})
}
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()))
}
}
}
fn drop_peer_interaction_projection(&self, interaction_id: Uuid) {
self.inbound_request_handling_modes
.lock()
.remove(&meerkat_core::PeerCorrelationId::from_uuid(interaction_id));
let removed_sender = self.subscriber_registry.lock().remove(&interaction_id);
let removed_stream = self
.interaction_stream_registry
.lock()
.remove(&interaction_id);
if removed_sender.is_some() || removed_stream.is_some() {
tracing::debug!(
interaction_id = %interaction_id,
"peer interaction projection dropped via DSL cleanup effect"
);
}
}
pub fn mark_interaction_complete(&self, interaction_id: Uuid) {
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(interaction_id);
let lifecycle_authority = self
.interaction_stream_registry
.lock()
.get(&interaction_id)
.map(|entry| entry.lifecycle_authority);
match lifecycle_authority {
Some(InteractionStreamLifecycleAuthority::Machine) => {
if let Some(handle) = self.interaction_stream_handle()
&& let Err(err) = handle.completed(corr_id)
{
tracing::trace!(
interaction_id = %interaction_id,
error = %err,
"InteractionStreamHandle::completed rejected (likely closed-early won race)"
);
}
}
Some(InteractionStreamLifecycleAuthority::LocalTransportOnly) | None => {
self.subscriber_registry.lock().remove(&interaction_id);
self.interaction_stream_registry
.lock()
.remove(&interaction_id);
}
}
}
pub fn reap_expired_reservations(&self) {
let candidates: Vec<(Uuid, InteractionStreamLifecycleAuthority)> = {
let registry = self.interaction_stream_registry.lock();
registry
.iter()
.filter_map(|(id, entry)| {
(entry.created_at.elapsed() > RESERVATION_TTL)
.then_some((*id, entry.lifecycle_authority))
})
.collect()
};
if candidates.is_empty() {
return;
}
let stream_handle = self.interaction_stream_handle();
let peer_handle = self.peer_interaction_handle();
for (id, lifecycle_authority) in candidates {
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(id);
match lifecycle_authority {
InteractionStreamLifecycleAuthority::Machine => {
let Some(handle) = stream_handle.as_ref() else {
continue;
};
if handle.state(corr_id) == Some(meerkat_core::InteractionStreamState::Reserved)
{
tracing::debug!(interaction_id = %id, "reservation expired (TTL)");
let _ = handle.expired(corr_id);
}
}
InteractionStreamLifecycleAuthority::LocalTransportOnly => {
let mut registry = self.interaction_stream_registry.lock();
if let Some(entry) = registry.get(&id)
&& entry.receiver.is_some()
{
tracing::debug!(interaction_id = %id, "reservation expired (TTL)");
registry.remove(&id);
}
}
}
if let Some(handle) = peer_handle.as_ref()
&& handle.outbound_state(corr_id).is_some()
{
let _ = handle.request_timed_out(corr_id);
}
}
}
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 corr_id = meerkat_core::PeerCorrelationId::from_uuid(interaction_id);
self.reserve_interaction_stream_channels(
corr_id,
interaction_id,
InteractionStreamReservationAuthority::LocalTransportOnly,
)?;
if let crate::inbox::AdmissionOutcome::Dropped { reason } = 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.expire_interaction_stream_on_send_failure(corr_id, interaction_id);
return Err(match reason {
crate::inbox::DropReason::InboxFull => {
SendError::Validation("input queue full".into())
}
crate::inbox::DropReason::SessionClosed => SendError::InputClosed,
crate::inbox::DropReason::UntrustedSender
| crate::inbox::DropReason::ClassificationRejected => {
SendError::Validation(format!("input rejected at ingress: {reason:?}"))
}
});
}
Ok(())
}
fn reserve_interaction_stream_channels(
&self,
corr_id: meerkat_core::PeerCorrelationId,
interaction_id: Uuid,
authority: InteractionStreamReservationAuthority,
) -> Result<(), SendError> {
let lifecycle_authority = match &authority {
InteractionStreamReservationAuthority::Machine(_) => {
InteractionStreamLifecycleAuthority::Machine
}
InteractionStreamReservationAuthority::LocalTransportOnly => {
InteractionStreamLifecycleAuthority::LocalTransportOnly
}
};
match authority {
InteractionStreamReservationAuthority::Machine(handle) => {
if let Err(err) = handle.reserved(corr_id) {
return Err(SendError::Validation(format!(
"DSL rejected InteractionStreamReserved for corr_id {corr_id}: {err}"
)));
}
}
InteractionStreamReservationAuthority::LocalTransportOnly => {}
}
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::new(sender, receiver, lifecycle_authority),
);
Ok(())
}
fn expire_interaction_stream_on_send_failure(
&self,
corr_id: meerkat_core::PeerCorrelationId,
interaction_id: Uuid,
) {
let lifecycle_authority = self
.interaction_stream_registry
.lock()
.get(&interaction_id)
.map(|entry| entry.lifecycle_authority);
match lifecycle_authority {
Some(InteractionStreamLifecycleAuthority::Machine) => {
if let Some(handle) = self.interaction_stream_handle() {
let _ = handle.expired(corr_id);
}
}
Some(InteractionStreamLifecycleAuthority::LocalTransportOnly) | None => {
self.interaction_stream_registry
.lock()
.remove(&interaction_id);
self.subscriber_registry.lock().remove(&interaction_id);
}
}
}
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.lock().drain(..) {
handle.abort();
}
self.listeners_started = false;
}
}
#[cfg(not(target_arch = "wasm32"))]
pub fn abort_listeners_for_rebind(&self) {
for handle in self.listener_handles.lock().drain(..) {
handle.abort();
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn stop_listeners_for_rebind(&self) {
let handles = self.listener_handles.lock().drain(..).collect::<Vec<_>>();
for handle in handles {
handle.abort_and_wait().await;
}
}
}
impl meerkat_core::handles::PeerInteractionCleanupObserver for CommsRuntime {
fn on_peer_interaction_cleanup(&self, corr_id: meerkat_core::PeerCorrelationId) {
self.drop_peer_interaction_projection(corr_id.as_uuid());
}
}
impl meerkat_core::handles::InteractionStreamCleanupObserver for CommsRuntime {
fn on_interaction_stream_cleanup(&self, corr_id: meerkat_core::PeerCorrelationId) {
let id = corr_id.as_uuid();
self.subscriber_registry.lock().remove(&id);
self.interaction_stream_registry.lock().remove(&id);
}
}
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<()>,
local_addr: Option<std::net::SocketAddr>,
}
#[cfg(not(target_arch = "wasm32"))]
impl ListenerHandle {
pub fn abort(&self) {
self.handle.abort();
}
async fn abort_and_wait(self) {
self.handle.abort();
let _ = self.handle.await;
}
}
#[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 _ =
handle_connection(stream, require_peer_auth, &keypair, &trusted, &inbox_sender)
.await;
});
}
});
Ok(ListenerHandle {
handle,
local_addr: None,
})
}
#[cfg(not(target_arch = "wasm32"))]
struct TcpListenerConfig {
addr: String,
keypair: Arc<Keypair>,
router: Arc<Router>,
trusted: Arc<parking_lot::RwLock<TrustedPeers>>,
require_peer_auth: bool,
inbox_sender: InboxSender,
pairing_password: Option<String>,
advertise_address: Option<String>,
trusted_peers_path: std::path::PathBuf,
participant_name: String,
bridge_bootstrap_token: String,
}
#[cfg(not(target_arch = "wasm32"))]
async fn spawn_tcp_listener(config: TcpListenerConfig) -> Result<ListenerHandle, std::io::Error> {
let TcpListenerConfig {
addr,
keypair,
router,
trusted,
require_peer_auth,
inbox_sender,
pairing_password,
advertise_address,
trusted_peers_path,
participant_name,
bridge_bootstrap_token,
} = config;
if let Some(secret) = pairing_password.as_deref() {
validate_pairing_secret(secret)?;
}
let listener = TcpListener::bind(&addr).await?;
let local_addr = listener.local_addr()?;
let target_address = if let Some(advertise_address) = advertise_address {
parse_peer_address(&advertise_address).map_err(|error| {
std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!("invalid comms advertised address '{advertise_address}': {error}"),
)
})?;
advertise_address
} else {
if local_addr.ip().is_unspecified() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"signed comms listener bound to a wildcard address requires an explicit advertised address",
));
}
format!("tcp://{local_addr}")
};
let handle = tokio::spawn(async move {
while let Ok((stream, _)) = listener.accept().await {
let (keypair, router, trusted, inbox_sender) = (
keypair.clone(),
router.clone(),
trusted.clone(),
inbox_sender.clone(),
);
let pairing_password = pairing_password.clone();
let trusted_peers_path = trusted_peers_path.clone();
let participant_name = participant_name.clone();
let bridge_bootstrap_token = bridge_bootstrap_token.clone();
let target_address = target_address.clone();
tokio::spawn(async move {
let _ = handle_tcp_connection_or_pairing(
stream,
require_peer_auth,
&keypair,
&router,
&trusted,
&inbox_sender,
pairing_password.as_deref(),
&trusted_peers_path,
&participant_name,
&bridge_bootstrap_token,
&target_address,
)
.await;
});
}
});
Ok(ListenerHandle {
handle,
local_addr: Some(local_addr),
})
}
#[cfg(not(target_arch = "wasm32"))]
#[allow(clippy::too_many_arguments)]
async fn handle_tcp_connection_or_pairing(
stream: TcpStream,
require_peer_auth: bool,
keypair: &Keypair,
router: &Router,
trusted: &Arc<parking_lot::RwLock<TrustedPeers>>,
inbox_sender: &InboxSender,
pairing_password: Option<&str>,
trusted_peers_path: &Path,
participant_name: &str,
bridge_bootstrap_token: &str,
target_address: &str,
) -> Result<(), std::io::Error> {
let mut first = [0_u8; 1];
let is_pairing = pairing_password.is_some()
&& matches!(
tokio::time::timeout(PAIRING_CLASSIFY_TIMEOUT, stream.peek(&mut first)).await,
Ok(Ok(1))
)
&& first[0] == b'{';
if is_pairing {
let Some(pairing_password) = pairing_password else {
return handle_connection(stream, require_peer_auth, keypair, trusted, inbox_sender)
.await
.map_err(|err| std::io::Error::other(err.to_string()));
};
return handle_pairing_connection(
stream,
keypair,
router,
trusted,
pairing_password,
trusted_peers_path,
participant_name,
bridge_bootstrap_token,
target_address,
)
.await;
}
handle_connection(stream, require_peer_auth, keypair, trusted, inbox_sender)
.await
.map_err(|err| std::io::Error::other(err.to_string()))
}
#[cfg(not(target_arch = "wasm32"))]
async fn read_pairing_line(
reader: &mut BufReader<TcpStream>,
buf: &mut String,
) -> Result<serde_json::Value, std::io::Error> {
buf.clear();
let bytes = tokio::time::timeout(std::time::Duration::from_secs(10), reader.read_line(buf))
.await
.map_err(|_| {
std::io::Error::new(std::io::ErrorKind::TimedOut, "pairing read timed out")
})??;
if bytes == 0 {
return Err(std::io::Error::new(
std::io::ErrorKind::UnexpectedEof,
"pairing peer closed connection",
));
}
serde_json::from_str(buf).map_err(|err| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("invalid pairing JSON: {err}"),
)
})
}
#[cfg(not(target_arch = "wasm32"))]
fn pairing_field<'a>(value: &'a serde_json::Value, field: &str) -> Result<&'a str, std::io::Error> {
value
.get(field)
.and_then(serde_json::Value::as_str)
.ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("pairing request missing string field '{field}'"),
)
})
}
#[cfg(not(target_arch = "wasm32"))]
async fn write_pairing_json(
writer: &mut TcpStream,
value: serde_json::Value,
) -> Result<(), std::io::Error> {
let bytes = serde_json::to_vec(&value).map_err(|err| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("pairing response encode failed: {err}"),
)
})?;
writer.write_all(&bytes).await?;
writer.write_all(b"\n").await?;
writer.flush().await
}
#[cfg(not(target_arch = "wasm32"))]
#[allow(clippy::too_many_arguments)]
async fn handle_pairing_connection(
stream: TcpStream,
keypair: &Keypair,
router: &Router,
trusted: &Arc<parking_lot::RwLock<TrustedPeers>>,
pairing_password: &str,
trusted_peers_path: &Path,
participant_name: &str,
bridge_bootstrap_token: &str,
target_address: &str,
) -> Result<(), std::io::Error> {
let mut reader = BufReader::new(stream);
let mut line = String::new();
let hello = read_pairing_line(&mut reader, &mut line).await?;
if pairing_field(&hello, "kind")? != PAIRING_HELLO_KIND {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"expected meerkat pairing hello",
));
}
if hello
.get("version")
.and_then(serde_json::Value::as_u64)
.unwrap_or_default()
!= u64::from(PAIRING_VERSION)
{
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"unsupported meerkat pairing version",
));
}
let challenge = Uuid::new_v4().to_string();
write_pairing_json(
reader.get_mut(),
serde_json::json!({
"kind": PAIRING_CHALLENGE_KIND,
"version": PAIRING_VERSION,
"challenge": challenge,
"target": {
"name": participant_name,
"address": target_address,
"identity": {
"kind": "ed25519_public_key",
"public_key": keypair.public_key().to_pubkey_string(),
}
}
}),
)
.await?;
let proof = read_pairing_line(&mut reader, &mut line).await?;
if pairing_field(&proof, "kind")? != PAIRING_PROOF_KIND {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"expected meerkat pairing proof",
));
}
let caller_value = proof.get("caller").cloned().ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"pairing proof missing caller",
)
})?;
let caller: PairingPeer = serde_json::from_value(caller_value).map_err(|err| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("invalid pairing caller: {err}"),
)
})?;
if caller.identity.kind != "ed25519_public_key" {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"pairing caller identity must be ed25519_public_key",
));
}
let supplied = pairing_field(&proof, "password_proof")?;
let expected = comms_pairing_password_proof(
pairing_password,
&challenge,
&caller.identity.public_key,
&caller.address,
);
if !constant_time_str_eq(supplied, &expected) {
return Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"invalid pairing password proof",
));
}
let pubkey = PubKey::from_pubkey_string(&caller.identity.public_key).map_err(|err| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("invalid pairing caller public key: {err}"),
)
})?;
router
.add_trusted_peer(TrustedPeer {
name: caller.name.clone(),
pubkey,
addr: caller.address.clone(),
meta: crate::PeerMeta::default(),
})
.map_err(|err| std::io::Error::new(std::io::ErrorKind::InvalidData, err.to_string()))?;
let peers_snapshot = trusted.read().clone();
peers_snapshot
.save(trusted_peers_path)
.await
.map_err(|err| std::io::Error::other(format!("persist paired trust: {err}")))?;
write_pairing_json(
reader.get_mut(),
serde_json::json!({
"kind": PAIRING_COMPLETE_KIND,
"version": PAIRING_VERSION,
"binding": {
"kind": "external",
"address": target_address,
"bootstrap_token": bridge_bootstrap_token,
"identity": {
"kind": "ed25519_public_key",
"public_key": keypair.public_key().to_pubkey_string(),
}
},
"comms_name": participant_name,
}),
)
.await
}
#[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,
local_addr: None,
})
}
#[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,
local_addr: None,
})
}
#[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, PeerId, PeerName, PeerReachability,
PeerReachabilityReason, PeerRoute, PeerSendability, StreamError, StreamScope,
TrustedPeerDescriptor,
},
interaction::InteractionId,
types::{ContentBlock, ImageData, SessionId},
};
fn trusted_descriptor(name: &str, pubkey: PubKey, address: &str) -> TrustedPeerDescriptor {
TrustedPeerDescriptor {
peer_id: crate::router::peer_id_from_pubkey(&pubkey),
name: PeerName::new(name.to_string()).expect("valid peer name"),
address: parse_peer_address(address).expect("valid peer address"),
pubkey: *pubkey.as_bytes(),
}
}
fn peer_route(name: &str, pubkey: PubKey) -> PeerRoute {
PeerRoute::with_display_name(
crate::router::peer_id_from_pubkey(&pubkey),
PeerName::new(name.to_string()).expect("valid peer name"),
)
}
fn missing_peer_route(name: &str) -> PeerRoute {
PeerRoute::with_display_name(
meerkat_core::comms::PeerId::new(),
PeerName::new(name.to_string()).expect("valid peer name"),
)
}
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
}
}
#[derive(Default)]
struct TestPeerInteractionHandle {
outbound:
Mutex<HashMap<meerkat_core::PeerCorrelationId, meerkat_core::OutboundPeerRequestState>>,
inbound:
Mutex<HashMap<meerkat_core::PeerCorrelationId, meerkat_core::InboundPeerRequestState>>,
cleanup_observer:
Mutex<Option<Arc<dyn meerkat_core::handles::PeerInteractionCleanupObserver>>>,
}
impl TestPeerInteractionHandle {
fn notify_cleanup(&self, corr_id: meerkat_core::PeerCorrelationId) {
if let Some(observer) = self.cleanup_observer.lock().clone() {
observer.on_peer_interaction_cleanup(corr_id);
}
}
fn guard_rejected(
context: &'static str,
corr_id: meerkat_core::PeerCorrelationId,
) -> meerkat_core::handles::DslTransitionError {
meerkat_core::handles::DslTransitionError::guard_rejected(
context,
format!("test authority rejected corr_id {corr_id}"),
)
}
}
impl meerkat_core::handles::PeerInteractionHandle for TestPeerInteractionHandle {
fn request_sent(
&self,
corr_id: meerkat_core::PeerCorrelationId,
_to: String,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
let mut outbound = self.outbound.lock();
if outbound.contains_key(&corr_id) {
return Err(Self::guard_rejected(
"PeerInteractionHandle::request_sent",
corr_id,
));
}
outbound.insert(corr_id, meerkat_core::OutboundPeerRequestState::Sent);
Ok(())
}
fn response_progress(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
let mut outbound = self.outbound.lock();
if !outbound.contains_key(&corr_id) {
return Err(Self::guard_rejected(
"PeerInteractionHandle::response_progress",
corr_id,
));
}
outbound.insert(
corr_id,
meerkat_core::OutboundPeerRequestState::AcceptedProgress,
);
Ok(())
}
fn response_terminal(
&self,
corr_id: meerkat_core::PeerCorrelationId,
_disposition: meerkat_core::handles::PeerTerminalDisposition,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
if self.outbound.lock().remove(&corr_id).is_none() {
return Err(Self::guard_rejected(
"PeerInteractionHandle::response_terminal",
corr_id,
));
}
self.notify_cleanup(corr_id);
Ok(())
}
fn request_timed_out(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
if self.outbound.lock().remove(&corr_id).is_none() {
return Err(Self::guard_rejected(
"PeerInteractionHandle::request_timed_out",
corr_id,
));
}
self.notify_cleanup(corr_id);
Ok(())
}
fn request_received(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
let mut inbound = self.inbound.lock();
if inbound.contains_key(&corr_id) {
return Err(Self::guard_rejected(
"PeerInteractionHandle::request_received",
corr_id,
));
}
inbound.insert(corr_id, meerkat_core::InboundPeerRequestState::Received);
Ok(())
}
fn response_replied(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
if self.inbound.lock().remove(&corr_id).is_none() {
return Err(Self::guard_rejected(
"PeerInteractionHandle::response_replied",
corr_id,
));
}
Ok(())
}
fn outbound_state(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Option<meerkat_core::OutboundPeerRequestState> {
self.outbound.lock().get(&corr_id).copied()
}
fn inbound_state(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Option<meerkat_core::InboundPeerRequestState> {
self.inbound.lock().get(&corr_id).copied()
}
fn install_cleanup_observer(
&self,
observer: Arc<dyn meerkat_core::handles::PeerInteractionCleanupObserver>,
) {
*self.cleanup_observer.lock() = Some(observer);
}
}
#[derive(Default)]
struct TestInteractionStreamHandle {
states:
Mutex<HashMap<meerkat_core::PeerCorrelationId, meerkat_core::InteractionStreamState>>,
cleanup_observer:
Mutex<Option<Arc<dyn meerkat_core::handles::InteractionStreamCleanupObserver>>>,
}
impl TestInteractionStreamHandle {
fn notify_cleanup(&self, corr_id: meerkat_core::PeerCorrelationId) {
if let Some(observer) = self.cleanup_observer.lock().clone() {
observer.on_interaction_stream_cleanup(corr_id);
}
}
fn guard_rejected(
context: &'static str,
corr_id: meerkat_core::PeerCorrelationId,
) -> meerkat_core::handles::DslTransitionError {
meerkat_core::handles::DslTransitionError::guard_rejected(
context,
format!("test stream authority rejected corr_id {corr_id}"),
)
}
fn terminal(
&self,
corr_id: meerkat_core::PeerCorrelationId,
context: &'static str,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
if self.states.lock().remove(&corr_id).is_none() {
return Err(Self::guard_rejected(context, corr_id));
}
self.notify_cleanup(corr_id);
Ok(())
}
}
impl meerkat_core::handles::InteractionStreamHandle for TestInteractionStreamHandle {
fn reserved(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
let mut states = self.states.lock();
if states.contains_key(&corr_id) {
return Err(Self::guard_rejected(
"InteractionStreamHandle::reserved",
corr_id,
));
}
states.insert(corr_id, meerkat_core::InteractionStreamState::Reserved);
Ok(())
}
fn attached(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
let mut states = self.states.lock();
match states.get_mut(&corr_id) {
Some(state) if *state == meerkat_core::InteractionStreamState::Reserved => {
*state = meerkat_core::InteractionStreamState::Attached;
Ok(())
}
_ => Err(Self::guard_rejected(
"InteractionStreamHandle::attached",
corr_id,
)),
}
}
fn completed(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
self.terminal(corr_id, "InteractionStreamHandle::completed")
}
fn expired(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
self.terminal(corr_id, "InteractionStreamHandle::expired")
}
fn closed_early(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Result<(), meerkat_core::handles::DslTransitionError> {
self.terminal(corr_id, "InteractionStreamHandle::closed_early")
}
fn state(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Option<meerkat_core::InteractionStreamState> {
self.states.lock().get(&corr_id).copied()
}
fn install_cleanup_observer(
&self,
observer: Arc<dyn meerkat_core::handles::InteractionStreamCleanupObserver>,
) {
*self.cleanup_observer.lock() = Some(observer);
}
}
fn install_test_peer_request_response_authority(runtime: &Arc<CommsRuntime>) {
runtime.install_peer_request_response_authority(PeerRequestResponseAuthority::new(
Arc::new(TestPeerInteractionHandle::default()),
Arc::new(TestInteractionStreamHandle::default()),
));
}
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,
advertise_address: 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,
pairing_password: None,
}
}
#[tokio::test]
async fn pairing_password_installs_trusted_peer_and_returns_binding() {
let tmp = tempfile::TempDir::new().unwrap();
let mut config = test_runtime_config("pairing-target", &tmp);
config.listen_tcp = Some("127.0.0.1:0".parse().unwrap());
config.pairing_password = Some("0123456789abcdef0123456789abcdef".to_string());
let mut runtime = CommsRuntime::new_machine_authority_required_with_silent_intents(
config,
Arc::new(HashSet::new()),
)
.await
.unwrap();
runtime.start_listeners().await.unwrap();
let addr = runtime.config.listen_tcp.unwrap();
let caller_keypair = Keypair::generate();
let caller_public_key = caller_keypair.public_key().to_pubkey_string();
let caller_address = "tcp://127.0.0.1:65530";
let stream = TcpStream::connect(addr).await.unwrap();
let mut reader = BufReader::new(stream);
reader
.get_mut()
.write_all(
serde_json::json!({
"kind": PAIRING_HELLO_KIND,
"version": PAIRING_VERSION,
})
.to_string()
.as_bytes(),
)
.await
.unwrap();
reader.get_mut().write_all(b"\n").await.unwrap();
let mut line = String::new();
reader.read_line(&mut line).await.unwrap();
let challenge: serde_json::Value = serde_json::from_str(&line).unwrap();
assert_eq!(challenge["kind"], PAIRING_CHALLENGE_KIND);
let challenge_token = challenge["challenge"].as_str().unwrap();
let proof = comms_pairing_password_proof(
"0123456789abcdef0123456789abcdef",
challenge_token,
&caller_public_key,
caller_address,
);
line.clear();
reader
.get_mut()
.write_all(
serde_json::json!({
"kind": PAIRING_PROOF_KIND,
"version": PAIRING_VERSION,
"password_proof": proof,
"caller": {
"name": "hive",
"address": caller_address,
"identity": {
"kind": "ed25519_public_key",
"public_key": caller_public_key,
}
}
})
.to_string()
.as_bytes(),
)
.await
.unwrap();
reader.get_mut().write_all(b"\n").await.unwrap();
reader.read_line(&mut line).await.unwrap();
let complete: serde_json::Value = serde_json::from_str(&line).unwrap();
assert_eq!(complete["kind"], PAIRING_COMPLETE_KIND);
assert_eq!(complete["binding"]["kind"], "external");
assert_eq!(
complete["binding"]["identity"]["public_key"],
runtime.public_key().to_pubkey_string()
);
let persisted = TrustedPeers::load(&runtime.config.trusted_peers_path)
.await
.unwrap();
assert!(
persisted
.peers
.iter()
.any(|peer| { peer.name == "hive" && peer.pubkey == caller_keypair.public_key() })
);
let live_peers = CoreCommsRuntime::peers(&runtime).await;
assert!(
live_peers
.iter()
.any(|peer| peer.peer_id == caller_keypair.public_key().to_peer_id()),
"pairing must update the live router peer-id index, not only trusted_peers.json"
);
}
#[tokio::test]
async fn pairing_classification_tolerates_delayed_hello() {
let tmp = tempfile::TempDir::new().unwrap();
let mut config = test_runtime_config("pairing-delayed-target", &tmp);
config.listen_tcp = Some("127.0.0.1:0".parse().unwrap());
config.pairing_password = Some("0123456789abcdef0123456789abcdef".to_string());
let mut runtime = CommsRuntime::new_machine_authority_required_with_silent_intents(
config,
Arc::new(HashSet::new()),
)
.await
.unwrap();
runtime.start_listeners().await.unwrap();
let addr = runtime.config.listen_tcp.unwrap();
let stream = TcpStream::connect(addr).await.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let mut reader = BufReader::new(stream);
reader
.get_mut()
.write_all(
serde_json::json!({
"kind": PAIRING_HELLO_KIND,
"version": PAIRING_VERSION,
})
.to_string()
.as_bytes(),
)
.await
.unwrap();
reader.get_mut().write_all(b"\n").await.unwrap();
let mut line = String::new();
reader.read_line(&mut line).await.unwrap();
let challenge: serde_json::Value = serde_json::from_str(&line).unwrap();
assert_eq!(challenge["kind"], PAIRING_CHALLENGE_KIND);
}
#[tokio::test]
async fn wildcard_signed_listener_requires_advertised_address() {
let tmp = tempfile::TempDir::new().unwrap();
let mut config = test_runtime_config("wildcard-target", &tmp);
config.listen_tcp = Some("0.0.0.0:0".parse().unwrap());
let mut runtime = CommsRuntime::new_machine_authority_required_with_silent_intents(
config,
Arc::new(HashSet::new()),
)
.await
.unwrap();
let error = runtime
.start_listeners()
.await
.expect_err("wildcard listener without advertise address should fail");
assert!(
error
.to_string()
.contains("requires an explicit advertised address"),
"unexpected error: {error}"
);
}
#[tokio::test]
async fn explicit_advertised_address_overrides_bound_listener_address() {
let tmp = tempfile::TempDir::new().unwrap();
let mut config = test_runtime_config("advertised-target", &tmp);
config.listen_tcp = Some("127.0.0.1:0".parse().unwrap());
config.advertise_address = Some("tcp://203.0.113.10:4200".to_string());
let mut runtime = CommsRuntime::new_machine_authority_required_with_silent_intents(
config,
Arc::new(HashSet::new()),
)
.await
.unwrap();
runtime.start_listeners().await.unwrap();
assert_eq!(runtime.advertised_address(), "tcp://203.0.113.10:4200");
}
#[tokio::test]
async fn machine_required_tcp_runtime_rejects_ingress_before_handle_or_listener_start() {
let tmp = tempfile::TempDir::new().unwrap();
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("pre-authority-sender-{suffix}");
let receiver_name = format!("pre-authority-receiver-{suffix}");
let sender = CommsRuntime::inproc_only(&sender_name).unwrap();
let mut config = test_runtime_config(&receiver_name, &tmp);
config.listen_tcp = Some("127.0.0.1:0".parse().unwrap());
let receiver = CommsRuntime::new_machine_authority_required_with_silent_intents(
config,
Arc::new(HashSet::new()),
)
.await
.unwrap();
assert!(receiver.peer_comms_machine_authority_required());
assert!(receiver.peer_comms_handle().is_none());
CoreCommsRuntime::add_trusted_peer(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await
.unwrap();
CoreCommsRuntime::add_trusted_peer(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await
.unwrap();
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
blocks: None,
to: peer_route(&receiver_name, receiver.public_key()),
body: "must not pass local classifier".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
},
)
.await;
assert!(matches!(
result,
Err(SendError::AdmissionDropped {
reason: meerkat_core::comms::AdmissionDropReason::ClassificationRejected
})
));
let interactions = CoreCommsRuntime::drain_inbox_interactions(&receiver).await;
assert!(
interactions.is_empty(),
"runtime-required ingress without a machine handle must not reach the inbox"
);
}
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
}
#[test]
fn bridge_bootstrap_token_is_stable_for_the_same_keypair() {
let keypair = Keypair::generate();
let first = CommsRuntime::derive_bridge_bootstrap_token(&keypair);
let second = CommsRuntime::derive_bridge_bootstrap_token(&keypair);
assert_eq!(
first, second,
"same keypair should derive same bootstrap token"
);
}
#[test]
fn bridge_bootstrap_token_changes_with_a_different_keypair() {
let first = CommsRuntime::derive_bridge_bootstrap_token(&Keypair::generate());
let second = CommsRuntime::derive_bridge_bootstrap_token(&Keypair::generate());
assert_ne!(
first, second,
"different keypairs should not derive the same bootstrap token"
);
}
#[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,
advertise_address: 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,
pairing_password: None,
};
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()),
advertise_address: 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,
pairing_password: None,
};
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,
advertise_address: 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,
pairing_password: None,
};
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();
CoreCommsRuntime::add_trusted_peer(
&runtime,
trusted_descriptor("sender", sender.public_key(), "tcp://127.0.0.1:4200"),
)
.await
.unwrap();
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(),
handling_mode: None,
},
);
let req = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
intent: "review".to_string(),
params: serde_json::json!({"pr": 19}),
blocks: None,
handling_mode: None,
},
);
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}),
blocks: None,
handling_mode: None,
},
);
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope: msg })
.into_result()
.unwrap();
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope: req })
.into_result()
.unwrap();
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope: resp })
.into_result()
.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 lifecycle_envelopes_do_not_seed_peer_response_handling_mode_cache() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("lifecycle-no-response-cache", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
let sender = Keypair::generate();
CoreCommsRuntime::add_trusted_peer(
&runtime,
trusted_descriptor("sender", sender.public_key(), "tcp://127.0.0.1:4200"),
)
.await
.unwrap();
let lifecycle_id = Uuid::new_v4();
let mut envelope = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Lifecycle {
kind: meerkat_core::comms::PeerLifecycleKind::PeerAdded,
params: serde_json::json!({"name": "sender"}),
},
);
envelope.id = lifecycle_id;
envelope.sign(&sender);
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope })
.into_result()
.unwrap();
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 1);
assert!(matches!(
&interactions[0].content,
meerkat_core::InteractionContent::Request { intent, .. }
if intent == meerkat_core::comms::PeerLifecycleKind::PeerAdded.as_str()
));
assert!(
!runtime
.inbound_request_handling_modes
.lock()
.contains_key(&meerkat_core::PeerCorrelationId::from_uuid(lifecycle_id)),
"lifecycle projections are request-shaped but are not responseable peer requests"
);
}
#[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();
CoreCommsRuntime::add_trusted_peer(
&runtime,
trusted_descriptor("sender", sender.public_key(), "tcp://127.0.0.1:4200"),
)
.await
.unwrap();
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()),
handling_mode: None,
},
);
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope: msg })
.into_result()
.unwrap();
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 1);
let interaction = &interactions[0];
assert_eq!(
interaction.rendered_text,
"Peer 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_peer_request_reserved_stream_correlates_via_request_envelope_id() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("corr-peer-{suffix}");
let runtime_name = format!("corr-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = Arc::new(CommsRuntime::inproc_only(&runtime_name).unwrap());
install_test_peer_request_response_authority(&runtime);
{
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(),
})
.expect("valid test peer should upsert");
}
CoreCommsRuntime::add_trusted_peer(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await
.expect("peer should accept trust entry for runtime");
let receipt = CoreCommsRuntime::send(
runtime.as_ref(),
CommsCommand::PeerRequest {
to: peer_route(&peer_name, peer.public_key()),
intent: "checksum".to_string(),
params: serde_json::json!({"path": "README.md"}),
blocks: None,
handling_mode: meerkat_core::types::HandlingMode::Queue,
stream: InputStreamMode::ReserveInteraction,
},
)
.await
.expect("peer request send should succeed");
let envelope_id = match receipt {
SendReceipt::PeerRequestSent { envelope_id, .. } => envelope_id,
other => panic!("expected PeerRequestSent, got {other:?}"),
};
let subscriber =
CoreCommsRuntime::interaction_subscriber(runtime.as_ref(), &InteractionId(envelope_id));
assert!(
subscriber.is_some(),
"reserved peer-request stream should correlate via the request envelope id carried in in_reply_to"
);
}
#[tokio::test]
async fn transport_only_runtime_rejects_peer_request_receipts_without_authority() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("transport-only-peer-{suffix}");
let runtime_name = format!("transport-only-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = CommsRuntime::inproc_only(&runtime_name).unwrap();
CoreCommsRuntime::add_trusted_peer(
&runtime,
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await
.expect("runtime should trust peer");
CoreCommsRuntime::add_trusted_peer(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await
.expect("peer should trust runtime");
let result = CoreCommsRuntime::send(
&runtime,
CommsCommand::PeerRequest {
to: peer_route(&peer_name, peer.public_key()),
intent: "must-have-authority".to_string(),
params: serde_json::json!({}),
blocks: None,
handling_mode: meerkat_core::types::HandlingMode::Queue,
stream: InputStreamMode::None,
},
)
.await;
assert!(
matches!(result, Err(SendError::Validation(ref message)) if message.contains("machine peer interaction authority")),
"transport-only runtime must fail before emitting PeerRequestSent, got {result:?}"
);
}
#[tokio::test]
async fn transport_only_runtime_rejects_peer_request_send_and_stream_without_authority() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("transport-only-stream-peer-{suffix}");
let runtime_name = format!("transport-only-stream-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = CommsRuntime::inproc_only(&runtime_name).unwrap();
CoreCommsRuntime::add_trusted_peer(
&runtime,
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await
.expect("runtime should trust peer");
CoreCommsRuntime::add_trusted_peer(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await
.expect("peer should trust runtime");
let result = CoreCommsRuntime::send_and_stream(
&runtime,
CommsCommand::PeerRequest {
to: peer_route(&peer_name, peer.public_key()),
intent: "must-have-stream-authority".to_string(),
params: serde_json::json!({}),
blocks: None,
handling_mode: meerkat_core::types::HandlingMode::Queue,
stream: InputStreamMode::ReserveInteraction,
},
)
.await;
assert!(
matches!(
result,
Err(SendAndStreamError::Send(SendError::Validation(ref message)))
if message.contains("machine peer interaction authority")
),
"transport-only runtime must fail before emitting streamed PeerRequestSent"
);
}
#[tokio::test]
async fn transport_only_runtime_rejects_peer_response_receipts_without_authority() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("transport-only-response-peer-{suffix}");
let runtime_name = format!("transport-only-response-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = CommsRuntime::inproc_only(&runtime_name).unwrap();
CoreCommsRuntime::add_trusted_peer(
&runtime,
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await
.expect("runtime should trust peer");
CoreCommsRuntime::add_trusted_peer(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await
.expect("peer should trust runtime");
let result = CoreCommsRuntime::send(
&runtime,
CommsCommand::PeerResponse {
to: peer_route(&peer_name, peer.public_key()),
in_reply_to: meerkat_core::InteractionId(Uuid::new_v4()),
status: meerkat_core::ResponseStatus::Completed,
result: serde_json::json!({"ok": true}),
blocks: None,
handling_mode: Some(meerkat_core::types::HandlingMode::Queue),
},
)
.await;
assert!(
matches!(result, Err(SendError::Validation(ref message)) if message.contains("machine peer interaction authority")),
"transport-only runtime must fail before emitting PeerResponseSent, got {result:?}"
);
}
#[tokio::test]
async fn terminal_peer_response_route_failure_keeps_inbound_request_retryable() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = Arc::new(
CommsRuntime::inproc_only(&format!("response-route-failure-{suffix}")).unwrap(),
);
let peer_handle = Arc::new(TestPeerInteractionHandle::default());
runtime.install_peer_request_response_authority(PeerRequestResponseAuthority::new(
peer_handle.clone(),
Arc::new(TestInteractionStreamHandle::default()),
));
let interaction_id = Uuid::new_v4();
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(interaction_id);
meerkat_core::handles::PeerInteractionHandle::request_received(
peer_handle.as_ref(),
corr_id,
)
.expect("test authority should seed inbound request state");
let result = CoreCommsRuntime::send(
runtime.as_ref(),
CommsCommand::PeerResponse {
to: missing_peer_route(&format!("missing-response-peer-{suffix}")),
in_reply_to: InteractionId(interaction_id),
status: meerkat_core::ResponseStatus::Completed,
result: serde_json::json!({"ok": true}),
blocks: None,
handling_mode: Some(meerkat_core::types::HandlingMode::Queue),
},
)
.await;
assert!(
matches!(result, Err(SendError::PeerNotFound(_))),
"test must fail at transport routing, got {result:?}"
);
assert_eq!(
meerkat_core::handles::PeerInteractionHandle::inbound_state(
peer_handle.as_ref(),
corr_id
),
Some(meerkat_core::InboundPeerRequestState::Received),
"failed terminal PeerResponse send must leave inbound request state retryable"
);
}
#[tokio::test]
async fn transport_only_peer_directory_does_not_advertise_semantic_request_response() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("directory-transport-peer-{suffix}");
let runtime_name = format!("directory-transport-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = CommsRuntime::inproc_only(&runtime_name).unwrap();
CoreCommsRuntime::add_trusted_peer(
&runtime,
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await
.expect("runtime should trust peer");
let peers = CoreCommsRuntime::peers(&runtime).await;
let entry = peers
.iter()
.find(|entry| entry.name.as_str() == peer_name)
.expect("trusted peer should be listed");
assert_eq!(entry.sendable_kinds, vec![PeerSendability::PeerMessage]);
}
#[tokio::test]
async fn partial_peer_authority_does_not_advertise_semantic_request_response() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("directory-partial-peer-{suffix}");
let runtime_name = format!("directory-partial-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = Arc::new(CommsRuntime::inproc_only(&runtime_name).unwrap());
runtime.install_peer_interaction_handle(Arc::new(TestPeerInteractionHandle::default()));
CoreCommsRuntime::add_trusted_peer(
runtime.as_ref(),
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await
.expect("runtime should trust peer");
let peers = CoreCommsRuntime::peers(runtime.as_ref()).await;
let entry = peers
.iter()
.find(|entry| entry.name.as_str() == peer_name)
.expect("trusted peer should be listed");
assert_eq!(entry.sendable_kinds, vec![PeerSendability::PeerMessage]);
}
#[tokio::test]
async fn partial_peer_authority_rejects_peer_request_receipts_without_stream_authority() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("partial-authority-peer-{suffix}");
let runtime_name = format!("partial-authority-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = Arc::new(CommsRuntime::inproc_only(&runtime_name).unwrap());
runtime.install_peer_interaction_handle(Arc::new(TestPeerInteractionHandle::default()));
CoreCommsRuntime::add_trusted_peer(
runtime.as_ref(),
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await
.expect("runtime should trust peer");
CoreCommsRuntime::add_trusted_peer(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await
.expect("peer should trust runtime");
let result = CoreCommsRuntime::send(
runtime.as_ref(),
CommsCommand::PeerRequest {
to: peer_route(&peer_name, peer.public_key()),
intent: "must-have-complete-authority".to_string(),
params: serde_json::json!({}),
blocks: None,
handling_mode: meerkat_core::types::HandlingMode::Queue,
stream: InputStreamMode::None,
},
)
.await;
assert!(
matches!(result, Err(SendError::Validation(ref message)) if message.contains("machine interaction stream authority")),
"partial peer authority must fail before emitting PeerRequestSent, got {result:?}"
);
}
#[tokio::test]
async fn partial_peer_authority_rejects_peer_response_receipts_without_stream_authority() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("partial-response-peer-{suffix}");
let runtime_name = format!("partial-response-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = Arc::new(CommsRuntime::inproc_only(&runtime_name).unwrap());
let peer_handle = Arc::new(TestPeerInteractionHandle::default());
runtime.install_peer_interaction_handle(peer_handle.clone());
CoreCommsRuntime::add_trusted_peer(
runtime.as_ref(),
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await
.expect("runtime should trust peer");
CoreCommsRuntime::add_trusted_peer(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await
.expect("peer should trust runtime");
let interaction_id = Uuid::new_v4();
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(interaction_id);
meerkat_core::handles::PeerInteractionHandle::request_received(
peer_handle.as_ref(),
corr_id,
)
.expect("seed inbound request state");
let result = CoreCommsRuntime::send(
runtime.as_ref(),
CommsCommand::PeerResponse {
to: peer_route(&peer_name, peer.public_key()),
in_reply_to: InteractionId(interaction_id),
status: meerkat_core::ResponseStatus::Completed,
result: serde_json::json!({"ok": true}),
blocks: None,
handling_mode: Some(meerkat_core::types::HandlingMode::Queue),
},
)
.await;
assert!(
matches!(result, Err(SendError::Validation(ref message)) if message.contains("machine interaction stream authority")),
"partial peer authority must fail before emitting PeerResponseSent, got {result:?}"
);
assert_eq!(
meerkat_core::handles::PeerInteractionHandle::inbound_state(
peer_handle.as_ref(),
corr_id
),
Some(meerkat_core::InboundPeerRequestState::Received),
"failed PeerResponse must leave inbound request state retryable"
);
}
#[tokio::test]
async fn peer_directory_advertises_semantic_request_response_with_authority() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("directory-semantic-peer-{suffix}");
let runtime_name = format!("directory-semantic-runtime-{suffix}");
let peer = CommsRuntime::inproc_only(&peer_name).unwrap();
let runtime = Arc::new(CommsRuntime::inproc_only(&runtime_name).unwrap());
install_test_peer_request_response_authority(&runtime);
CoreCommsRuntime::add_trusted_peer(
runtime.as_ref(),
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await
.expect("runtime should trust peer");
let peers = CoreCommsRuntime::peers(runtime.as_ref()).await;
let entry = peers
.iter()
.find(|entry| entry.name.as_str() == peer_name)
.expect("trusted peer should be listed");
assert_eq!(
entry.sendable_kinds,
vec![
PeerSendability::PeerMessage,
PeerSendability::PeerRequest,
PeerSendability::PeerResponse,
]
);
}
#[tokio::test]
async fn test_peer_lifecycle_send_stays_typed_and_non_rendered() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("life-sender-{suffix}");
let receiver_name = format!("life-receiver-{suffix}");
let sender = CommsRuntime::inproc_only(&sender_name).unwrap();
let receiver = CommsRuntime::inproc_only(&receiver_name).unwrap();
CoreCommsRuntime::add_trusted_peer(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await
.expect("sender should trust receiver");
CoreCommsRuntime::add_trusted_peer(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await
.expect("receiver should trust sender");
let receipt = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerLifecycle {
to: peer_route(&receiver_name, receiver.public_key()),
kind: meerkat_core::comms::PeerLifecycleKind::PeerAdded,
params: serde_json::json!({
"peer": "worker-1",
"role": "worker",
"peer_name": "mob/worker/worker-1",
}),
},
)
.await
.expect("peer lifecycle send should succeed");
match receipt {
SendReceipt::PeerLifecycleSent { .. } => {}
other => panic!("expected PeerLifecycleSent, got {other:?}"),
}
let interactions = receiver
.drain_classified_inbox_interactions()
.await
.expect("classified drain should succeed");
assert_eq!(interactions.len(), 1);
assert_eq!(
interactions[0].class(),
meerkat_core::PeerInputClass::PeerLifecycleAdded
);
assert_eq!(interactions[0].lifecycle_peer.as_deref(), Some("worker-1"));
assert!(
interactions[0].interaction.rendered_text.is_empty(),
"silent lifecycle notices must not synthesize prompt text"
);
}
#[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]
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]
async fn local_input_stream_remains_transport_only_with_machine_stream_authority() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = Arc::new(
CommsRuntime::inproc_only(&format!("local-transport-authority-{suffix}")).unwrap(),
);
let stream_handle = Arc::new(TestInteractionStreamHandle::default());
runtime.install_interaction_stream_handle(stream_handle.clone());
let cmd = CommsCommand::Input {
blocks: None,
session_id: SessionId::new(),
body: "local input stream".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.as_ref(), cmd)
.await
.expect("local input send_and_stream should attach without semantic stream authority");
let interaction_id = match receipt {
SendReceipt::InputAccepted {
interaction_id,
stream_reserved,
} => {
assert!(stream_reserved);
interaction_id
}
other => panic!("expected InputAccepted, got {other:?}"),
};
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(interaction_id.0);
assert_eq!(
meerkat_core::handles::InteractionStreamHandle::state(stream_handle.as_ref(), corr_id),
None,
"local input streams must not create semantic stream DSL state"
);
let duplicate =
CoreCommsRuntime::stream(runtime.as_ref(), StreamScope::Interaction(interaction_id));
assert!(
matches!(duplicate, Err(StreamError::AlreadyAttached(_))),
"local input streams should keep transport-only duplicate attach semantics"
);
}
#[tokio::test]
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,
})
.into_result()
.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()),
})
.into_result()
.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));
}
#[tokio::test]
async fn test_plain_event_without_interaction_id_gets_stable_id_between_snapshot_and_drain() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("plain-generated-id", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
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: None,
render_metadata: None,
})
.into_result()
.unwrap();
let snapshot = CoreCommsRuntime::peer_ingress_queue_snapshot(&runtime)
.await
.expect("classified queue snapshot should be available");
let ingress_id = snapshot
.queued_entries
.first()
.and_then(|entry| entry.interaction_id)
.expect("plain event should have a stable generated interaction id at ingress")
.0;
assert_eq!(
snapshot.queued_entries[0].raw_item_id,
meerkat_core::InteractionId(ingress_id),
"the queue snapshot should expose the same stable ingress id it uses as the raw item key"
);
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 1);
assert_eq!(interactions[0].id.0, ingress_id);
}
#[tokio::test]
async fn test_peer_ingress_queue_snapshot_reflects_classified_queue_without_draining() {
let tmp = tempfile::TempDir::new().unwrap();
let mut config = test_runtime_config("peer-snapshot", &tmp);
config.require_peer_auth = false;
let runtime = CommsRuntime::new(config).await.unwrap();
let sender = Keypair::generate();
let envelope = signed_envelope(
&sender,
runtime.public_key,
crate::types::MessageKind::Request {
intent: "review".to_string(),
params: serde_json::json!({"pr": 42}),
blocks: None,
handling_mode: None,
},
);
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope })
.into_result()
.unwrap();
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: None,
render_metadata: None,
})
.into_result()
.unwrap();
let snapshot = CoreCommsRuntime::peer_ingress_queue_snapshot(&runtime)
.await
.expect("classified queue snapshot should be available");
assert_eq!(snapshot.total_count, 2);
assert_eq!(snapshot.actionable_count, 2);
assert_eq!(snapshot.plain_event_count, 1);
assert_eq!(snapshot.response_count, 0);
assert_eq!(snapshot.lifecycle_count, 0);
assert_eq!(snapshot.queued_entries.len(), 2);
assert_eq!(
snapshot.queued_entries[0].kind,
meerkat_core::PeerIngressKind::Request
);
assert_eq!(
snapshot.queued_entries[1].kind,
meerkat_core::PeerIngressKind::PlainEvent
);
let interactions = runtime.drain_classified_inbox_interactions().await.unwrap();
assert_eq!(
interactions.len(),
2,
"snapshot must not consume classified peer ingress"
);
}
#[tokio::test]
async fn test_peer_ingress_runtime_snapshot_reflects_trusted_peers_and_queue() {
let tmp = tempfile::TempDir::new().unwrap();
let mut config = test_runtime_config("peer-runtime-snapshot", &tmp);
config.require_peer_auth = false;
let runtime = CommsRuntime::new(config).await.unwrap();
let peer_key = Keypair::generate();
let trusted_peer = trusted_descriptor("ally", peer_key.public_key(), "inproc://ally");
CoreCommsRuntime::add_trusted_peer(&runtime, trusted_peer.clone())
.await
.expect("add_trusted_peer should succeed");
runtime
.event_injector()
.inject(
"evt".to_string().into(),
meerkat_core::PlainEventSource::Tcp,
meerkat_core::types::HandlingMode::Queue,
None,
)
.expect("plain event injection should succeed");
let snapshot = CoreCommsRuntime::peer_ingress_runtime_snapshot(&runtime)
.await
.expect("peer runtime snapshot should be available");
assert_eq!(snapshot.self_peer_id, runtime.public_key().to_peer_id());
assert!(!snapshot.auth_required);
assert_eq!(
snapshot.authority_phase,
meerkat_core::PeerIngressAuthorityPhase::Absent
);
assert_eq!(snapshot.trusted_peers, vec![trusted_peer]);
assert_eq!(snapshot.submission_queue_len, 1);
assert_eq!(snapshot.queue.total_count, 1);
assert_eq!(snapshot.queue.plain_event_count, 1);
assert_eq!(snapshot.queue.queued_entries.len(), 1);
assert_eq!(
snapshot.queue.queued_entries[0].kind,
meerkat_core::PeerIngressKind::PlainEvent
);
assert_eq!(snapshot.queue.queued_entries[0].admission_diagnostic, None);
assert_eq!(
snapshot.queue.queued_entries[0].raw_item_id,
snapshot.queue.queued_entries[0]
.interaction_id
.expect("plain event should retain an ingress interaction id")
);
}
#[tokio::test]
async fn test_peer_ingress_runtime_snapshot_skips_direct_zero_pubkey_trust_entry() {
let tmp = tempfile::TempDir::new().unwrap();
let runtime = CommsRuntime::new(test_runtime_config("peer-runtime-zero-snapshot", &tmp))
.await
.unwrap();
let zero_pubkey = PubKey::new([0u8; 32]);
runtime
.trusted_peers_shared()
.write()
.peers
.push(crate::TrustedPeer {
name: "zero-peer".to_string(),
pubkey: zero_pubkey,
addr: "inproc://zero-peer".to_string(),
meta: crate::PeerMeta::default(),
});
let snapshot = CoreCommsRuntime::peer_ingress_runtime_snapshot(&runtime)
.await
.expect("peer runtime snapshot should be available");
assert!(
snapshot.trusted_peers.is_empty(),
"direct-injected zero-pubkey trust must not enter the machine snapshot"
);
assert!(
CoreCommsRuntime::peers(&runtime).await.is_empty(),
"direct-injected zero-pubkey trust must not enter the peer directory"
);
assert!(
!runtime
.trusted_peers_shared()
.read()
.is_trusted(&zero_pubkey),
"direct-injected zero-pubkey trust must not pass admission lookup"
);
}
#[tokio::test]
async fn test_peer_ingress_runtime_snapshot_reflects_dropped_peer_authority_state() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("peer-runtime-dropped", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
let sender = Keypair::generate();
let envelope = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
intent: "review".to_string(),
params: serde_json::json!({"scope": "peer"}),
blocks: None,
handling_mode: None,
},
);
let outcome = runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope });
assert_eq!(
outcome,
crate::inbox::AdmissionOutcome::Dropped {
reason: crate::inbox::DropReason::UntrustedSender
}
);
let snapshot = CoreCommsRuntime::peer_ingress_runtime_snapshot(&runtime)
.await
.expect("peer runtime snapshot should be available");
assert_eq!(
snapshot.authority_phase,
meerkat_core::PeerIngressAuthorityPhase::Dropped
);
assert_eq!(snapshot.submission_queue_len, 0);
assert_eq!(snapshot.queue.total_count, 0);
assert!(snapshot.trusted_peers.is_empty());
}
#[tokio::test]
async fn test_live_peer_authority_syncs_trust_receive_and_drain() {
use crate::peer_types::PeerIngressState;
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("peer-authority-sync", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
let sender = Keypair::generate();
let pubkey_for_authority_check = sender.public_key();
let envelope = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
intent: "review".to_string(),
params: serde_json::json!({"scope": "peer"}),
blocks: None,
handling_mode: None,
},
);
let outcome = runtime
.router
.inbox_sender()
.send_classified(InboxItem::External {
envelope: envelope.clone(),
});
assert_eq!(
outcome,
crate::inbox::AdmissionOutcome::Dropped {
reason: crate::inbox::DropReason::UntrustedSender
}
);
{
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox
.classified_snapshot()
.expect("classified snapshot should exist")
.total_count,
0
);
assert_eq!(
inbox.peer_authority_test_snapshot(),
Some((PeerIngressState::Dropped, 0))
);
assert_eq!(
inbox.peer_authority_trusts_peer_for_test(&pubkey_for_authority_check),
Some(false)
);
assert_eq!(inbox.dropped_count(), Some(1));
}
let trusted_peer = trusted_descriptor("sender", sender.public_key(), "inproc://sender");
CoreCommsRuntime::add_trusted_peer(&runtime, trusted_peer)
.await
.expect("add_trusted_peer should succeed");
{
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox.peer_authority_test_snapshot(),
Some((PeerIngressState::Dropped, 0))
);
assert_eq!(
inbox.peer_authority_trusts_peer_for_test(&pubkey_for_authority_check),
Some(true)
);
}
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope })
.into_result()
.unwrap();
let snapshot = CoreCommsRuntime::peer_ingress_runtime_snapshot(&runtime)
.await
.expect("peer runtime snapshot should be available");
assert_eq!(
snapshot.authority_phase,
meerkat_core::PeerIngressAuthorityPhase::Received
);
assert_eq!(snapshot.submission_queue_len, 1);
assert_eq!(snapshot.queue.total_count, 1);
assert_eq!(
snapshot.queue.queued_entries[0].admission_diagnostic,
Some(meerkat_core::PeerIngressAdmissionDiagnostic::TrustedAtAdmission)
);
{
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox
.classified_snapshot()
.expect("classified snapshot should exist")
.total_count,
1
);
assert_eq!(
inbox.peer_authority_test_snapshot(),
Some((PeerIngressState::Received, 1))
);
}
let interactions = runtime.drain_classified_inbox_interactions().await.unwrap();
assert_eq!(interactions.len(), 1);
{
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox.peer_authority_test_snapshot(),
Some((PeerIngressState::Delivered, 0))
);
}
let removed = CoreCommsRuntime::remove_trusted_peer(
&runtime,
&pubkey_for_authority_check.to_peer_id().to_string(),
)
.await
.expect("remove_trusted_peer should succeed");
assert!(removed);
{
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox.peer_authority_test_snapshot(),
Some((PeerIngressState::Delivered, 0))
);
assert_eq!(
inbox.peer_authority_trusts_peer_for_test(&pubkey_for_authority_check),
Some(false)
);
}
}
#[tokio::test]
async fn test_live_peer_authority_accepts_unknown_peer_when_auth_open() {
use crate::peer_types::PeerIngressState;
let tmp = tempfile::TempDir::new().unwrap();
let mut config = test_runtime_config("peer-authority-auth-open", &tmp);
config.require_peer_auth = false;
let runtime = CommsRuntime::new(config).await.unwrap();
let sender = Keypair::generate();
let envelope = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
intent: "status".to_string(),
params: serde_json::json!({}),
blocks: None,
handling_mode: None,
},
);
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope })
.into_result()
.unwrap();
{
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox
.classified_snapshot()
.expect("classified snapshot should exist")
.total_count,
1
);
assert_eq!(
inbox.peer_authority_test_snapshot(),
Some((PeerIngressState::Received, 1))
);
}
let interactions = runtime.drain_classified_inbox_interactions().await.unwrap();
assert_eq!(interactions.len(), 1);
{
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox.peer_authority_test_snapshot(),
Some((PeerIngressState::Delivered, 0))
);
}
}
#[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: missing_peer_route("missing-peer"),
body: "hello".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
};
let result = CoreCommsRuntime::send(&runtime, cmd).await;
assert!(matches!(result, Err(SendError::PeerNotFound(_))));
}
#[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();
CoreCommsRuntime::add_trusted_peer(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await
.unwrap();
CoreCommsRuntime::add_trusted_peer(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await
.unwrap();
let cmd = CommsCommand::PeerMessage {
blocks: None,
to: peer_route(&receiver_name, receiver.public_key()),
body: "greeting".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
};
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);
CoreCommsRuntime::add_trusted_peer(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await
.unwrap();
CoreCommsRuntime::add_trusted_peer(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await
.unwrap();
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: peer_route(&receiver_name, receiver.public_key()),
body: "blob-backed image".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
};
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_hydrates_blob_refs_on_peer_request_before_transport() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("sender-request-blob-{suffix}");
let receiver_name = format!("receiver-request-blob-{suffix}");
let mut sender_runtime = 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_runtime.set_blob_store(blob_store);
let sender = Arc::new(sender_runtime);
install_test_peer_request_response_authority(&sender);
CoreCommsRuntime::add_trusted_peer(
sender.as_ref(),
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await
.unwrap();
CoreCommsRuntime::add_trusted_peer(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await
.unwrap();
let cmd = CommsCommand::PeerRequest {
blocks: Some(vec![ContentBlock::Image {
media_type: "image/png".to_string(),
data: ImageData::Blob {
blob_id: blob_ref.blob_id,
},
}]),
to: peer_route(&receiver_name, receiver.public_key()),
intent: "checksum".to_string(),
params: serde_json::json!({"subject": "image"}),
handling_mode: meerkat_core::types::HandlingMode::Queue,
stream: InputStreamMode::ReserveInteraction,
};
let receipt = CoreCommsRuntime::send(sender.as_ref(), cmd).await;
assert!(matches!(receipt, Ok(SendReceipt::PeerRequestSent { .. })));
let interactions = CoreCommsRuntime::drain_inbox_interactions(&receiver).await;
assert_eq!(interactions.len(), 1);
match &interactions[0].content {
meerkat_core::InteractionContent::Request { 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 request interaction, got {other:?}"),
}
}
#[tokio::test]
async fn test_core_send_hydrates_blob_refs_on_peer_response_before_transport() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("sender-response-blob-{suffix}");
let receiver_name = format!("receiver-response-blob-{suffix}");
let mut sender_runtime = 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_runtime.set_blob_store(blob_store);
let sender = Arc::new(sender_runtime);
let peer_handle = Arc::new(TestPeerInteractionHandle::default());
sender.install_peer_request_response_authority(PeerRequestResponseAuthority::new(
peer_handle.clone(),
Arc::new(TestInteractionStreamHandle::default()),
));
CoreCommsRuntime::add_trusted_peer(
sender.as_ref(),
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await
.unwrap();
CoreCommsRuntime::add_trusted_peer(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await
.unwrap();
let interaction_id = Uuid::new_v4();
let corr_id = meerkat_core::PeerCorrelationId::from_uuid(interaction_id);
meerkat_core::handles::PeerInteractionHandle::request_received(
peer_handle.as_ref(),
corr_id,
)
.expect("seed inbound request state");
let cmd = CommsCommand::PeerResponse {
blocks: Some(vec![ContentBlock::Image {
media_type: "image/png".to_string(),
data: ImageData::Blob {
blob_id: blob_ref.blob_id,
},
}]),
to: peer_route(&receiver_name, receiver.public_key()),
in_reply_to: InteractionId(interaction_id),
status: meerkat_core::ResponseStatus::Completed,
result: serde_json::json!({"ok": true}),
handling_mode: Some(meerkat_core::types::HandlingMode::Queue),
};
let receipt = CoreCommsRuntime::send(sender.as_ref(), cmd).await;
assert!(matches!(receipt, Ok(SendReceipt::PeerResponseSent { .. })));
let interactions = CoreCommsRuntime::drain_inbox_interactions(&receiver).await;
assert_eq!(interactions.len(), 1);
match &interactions[0].content {
meerkat_core::InteractionContent::Response { 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 response 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: peer_route(&receiver_name, receiver.public_key()),
body: "inproc-only hello".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
};
let receipt = CoreCommsRuntime::send(&sender, cmd).await;
match receipt {
Err(SendError::PeerNotFound(_)) => {}
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,
advertise_address: None,
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
allow_external_unauthenticated: false,
pairing_password: None,
};
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,
advertise_address: None,
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
allow_external_unauthenticated: false,
pairing_password: None,
};
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: peer_route(&receiver_name, receiver.public_key()),
body: "hello without trusted".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
};
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 = trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
);
CoreCommsRuntime::add_trusted_peer(&sender, peer_spec)
.await
.expect("trusted peer add should succeed");
let reverse_spec = trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
);
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: peer_route(&receiver_name, receiver.public_key()),
body: "hello trusted peer".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
};
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_register_trusted_peer_preserves_meta_and_syncs_peer_authority() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("trust-register-{suffix}")).unwrap();
let peer_key = Keypair::generate().public_key();
let pubkey_for_authority_check = peer_key;
let meta = crate::PeerMeta::default()
.with_description("Coordinates delegated work")
.with_label("team", "platform");
runtime
.register_trusted_peer(crate::TrustedPeer {
name: "ally".to_string(),
pubkey: peer_key,
addr: "inproc://ally".to_string(),
meta: meta.clone(),
})
.await
.expect("trusted peer registration should succeed");
{
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox.peer_authority_trusts_peer_for_test(&pubkey_for_authority_check),
Some(true)
);
}
let peers = CoreCommsRuntime::peers(&runtime).await;
let ally = peers
.iter()
.find(|entry| entry.name.as_str() == "ally")
.expect("registered peer should be listed");
assert_eq!(ally.meta, meta);
assert!(
runtime
.unregister_trusted_pubkey(&peer_key)
.await
.expect("trusted peer removal should succeed")
);
{
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox.peer_authority_trusts_peer_for_test(&pubkey_for_authority_check),
Some(false)
);
}
assert!(
CoreCommsRuntime::peers(&runtime)
.await
.iter()
.all(|entry| entry.name.as_str() != "ally")
);
}
#[tokio::test]
async fn test_add_trusted_peer_invalid_peer_id_is_rejected() {
let sender = CommsRuntime::inproc_only("trust-invalid-sender").unwrap();
let bogus_peer_key = Keypair::generate();
let real_peer_key = Keypair::generate();
let invalid_peer_spec = meerkat_core::comms::TrustedPeerDescriptor {
name: meerkat_core::comms::PeerName::new("invalid").expect("valid peer name"),
peer_id: bogus_peer_key.public_key().to_peer_id(),
address: meerkat_core::comms::PeerAddress::new(
meerkat_core::comms::PeerTransport::Inproc,
"invalid",
),
pubkey: *real_peer_key.public_key().as_bytes(),
};
let result = CoreCommsRuntime::add_trusted_peer(&sender, invalid_peer_spec).await;
assert!(matches!(result, Err(SendError::Validation(_))));
}
#[tokio::test]
async fn test_add_trusted_peer_zero_pubkey_is_rejected_at_trust_registration() {
let runtime = CommsRuntime::inproc_only("trust-zero-pubkey-sender").unwrap();
let zero_pubkey_spec = TrustedPeerDescriptor::test_only_unsigned_typed(
"zero-pubkey-peer",
PeerId::new(),
"inproc://zero-pubkey-peer",
)
.expect("test-only descriptor should build before trust registration");
let result = CoreCommsRuntime::add_trusted_peer(&runtime, zero_pubkey_spec).await;
match result {
Err(SendError::Validation(message)) => {
assert!(
message.contains("pubkey") && message.contains("non-zero"),
"unexpected validation message: {message}"
);
}
other => panic!("zero-pubkey trust registration must fail, got {other:?}"),
}
assert!(
CoreCommsRuntime::peers(&runtime)
.await
.iter()
.all(|entry| entry.name.as_str() != "zero-pubkey-peer"),
"rejected zero-pubkey peer must not enter the peer directory"
);
}
#[tokio::test]
async fn test_register_trusted_peer_zero_pubkey_is_rejected_at_raw_boundary() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("raw-zero-pubkey-{suffix}")).unwrap();
let zero_pubkey = PubKey::new([0u8; 32]);
let result = runtime
.register_trusted_peer(crate::TrustedPeer {
name: format!("raw-zero-peer-{suffix}"),
pubkey: zero_pubkey,
addr: format!("inproc://raw-zero-peer-{suffix}"),
meta: crate::PeerMeta::default(),
})
.await;
match result {
Err(SendError::Validation(message)) => {
assert!(
message.contains("pubkey") && message.contains("non-zero"),
"unexpected validation message: {message}"
);
}
other => panic!("raw zero-pubkey trust registration must fail, got {other:?}"),
}
assert!(
!runtime
.trusted_peers_shared()
.read()
.is_trusted(&zero_pubkey),
"raw zero-pubkey peer must not enter the shared trust authority"
);
let inbox = runtime.inbox.lock().await;
assert_eq!(
inbox.peer_authority_trusts_peer_for_test(&zero_pubkey),
Some(false),
"classified inbox authority must not trust the rejected zero pubkey"
);
}
#[tokio::test]
async fn test_runtime_startup_rejects_persisted_zero_pubkey_trust() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("persisted-zero-pubkey", &tmp);
tokio::fs::write(
&config.trusted_peers_path,
serde_json::json!({
"peers": [{
"name": "persisted-zero-peer",
"pubkey": PubKey::new([0u8; 32]).to_pubkey_string(),
"addr": "inproc://persisted-zero-peer"
}]
})
.to_string(),
)
.await
.unwrap();
let error = match CommsRuntime::new(config).await {
Ok(_) => panic!("runtime startup must reject persisted zero-pubkey trust"),
Err(error) => error,
};
match error {
CommsRuntimeError::TrustLoadError(message) => {
assert!(
message.contains("pubkey") && message.contains("non-zero"),
"unexpected trust-load error: {message}"
);
}
other => panic!("unexpected runtime startup error: {other:?}"),
}
}
#[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 an ed25519 public key"
);
let public_key_bytes = <CommsRuntime as CoreCommsRuntime>::public_key_bytes(&runtime)
.expect("typed public key bytes should be available");
assert_ne!(public_key_bytes, [0u8; 32]);
assert_eq!(
<CommsRuntime as CoreCommsRuntime>::peer_id(&runtime),
Some(PeerId::from_ed25519_pubkey(&public_key_bytes))
);
}
fn test_claim_handle() -> Arc<dyn SessionClaimHandle> {
Arc::new(meerkat_core::handles::DefaultSessionClaimRegistry::new())
}
#[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 claim_handle = test_claim_handle();
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()),
Arc::clone(&claim_handle),
)
.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()),
claim_handle,
)
.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 claim_handle = test_claim_handle();
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()),
Arc::clone(&claim_handle),
)
.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()),
claim_handle,
)
.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 claim_handle = test_claim_handle();
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()),
Arc::clone(&claim_handle),
)
.await
.expect("create session-scoped runtime");
let peer = CommsRuntime::inproc_only("session-scoped-trust-peer").unwrap();
CoreCommsRuntime::add_trusted_peer(
&sender,
trusted_descriptor(
"session-scoped-trust-peer",
peer.public_key(),
&format!("inproc://{}", peer.participant_name()),
),
)
.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()),
claim_handle,
)
.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(),
})
.expect("valid test peer should upsert");
}
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.to_string(), format!("inproc://{peer_name}"));
assert_eq!(peer.sendable_kinds, vec![PeerSendability::PeerMessage]);
assert_eq!(peer.reachability, PeerReachability::Unknown);
assert_eq!(peer.last_unreachable_reason, None);
}
#[tokio::test]
async fn test_core_peers_rejects_unknown_and_schemeless_addresses() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime_name = format!("runtime-address-parse-{suffix}");
let unknown_name = format!("unknown-scheme-{suffix}");
let schemeless_name = format!("schemeless-{suffix}");
let runtime = CommsRuntime::inproc_only(&runtime_name).unwrap();
let unknown_pubkey = Keypair::generate().public_key();
let schemeless_pubkey = Keypair::generate().public_key();
{
let mut trusted = runtime.trusted_peers.write();
trusted
.upsert(crate::TrustedPeer {
name: unknown_name.clone(),
pubkey: unknown_pubkey,
addr: "http://127.0.0.1:4200".to_string(),
meta: crate::PeerMeta::default(),
})
.expect("valid test peer should upsert");
trusted
.upsert(crate::TrustedPeer {
name: schemeless_name.clone(),
pubkey: schemeless_pubkey,
addr: "127.0.0.1:4201".to_string(),
meta: crate::PeerMeta::default(),
})
.expect("valid test peer should upsert");
}
let peers = CoreCommsRuntime::peers(&runtime).await;
assert!(
peers
.iter()
.all(|entry| entry.name.as_str() != unknown_name),
"unknown-scheme trusted peer must not be advertised as TCP"
);
assert!(
peers
.iter()
.all(|entry| entry.name.as_str() != schemeless_name),
"schemeless trusted peer must not be advertised as TCP"
);
}
#[tokio::test]
async fn test_register_trusted_peer_rejects_unknown_and_schemeless_addresses() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime =
CommsRuntime::inproc_only(&format!("runtime-register-address-{suffix}")).unwrap();
let unknown_name = format!("register-unknown-{suffix}");
let schemeless_name = format!("register-schemeless-{suffix}");
let err = runtime
.register_trusted_peer(crate::TrustedPeer {
name: unknown_name.clone(),
pubkey: Keypair::generate().public_key(),
addr: "http://127.0.0.1:4202".to_string(),
meta: crate::PeerMeta::default(),
})
.await
.expect_err("unknown address scheme must be rejected at registration");
assert!(
matches!(err, SendError::Validation(ref message) if message.contains("unknown peer address transport")),
"unexpected error: {err:?}",
);
let err = runtime
.register_trusted_peer(crate::TrustedPeer {
name: schemeless_name.clone(),
pubkey: Keypair::generate().public_key(),
addr: "127.0.0.1:4203".to_string(),
meta: crate::PeerMeta::default(),
})
.await
.expect_err("schemeless TCP address must be rejected at registration");
assert!(
matches!(err, SendError::Validation(ref message) if message.contains("missing transport scheme")),
"unexpected error: {err:?}",
);
assert!(
runtime
.trusted_peers
.read()
.peers
.iter()
.all(|peer| peer.name != unknown_name && peer.name != schemeless_name),
"invalid peers must not be stored"
);
}
#[tokio::test]
async fn test_register_private_trusted_peer_rejects_unknown_and_schemeless_addresses() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime =
CommsRuntime::inproc_only(&format!("runtime-private-address-{suffix}")).unwrap();
let unknown_name = format!("private-unknown-{suffix}");
let schemeless_name = format!("private-schemeless-{suffix}");
let err = runtime
.register_private_trusted_peer(crate::TrustedPeer {
name: unknown_name.clone(),
pubkey: Keypair::generate().public_key(),
addr: "http://127.0.0.1:4204".to_string(),
meta: crate::PeerMeta::default(),
})
.await
.expect_err("unknown address scheme must be rejected for private peers");
assert!(
matches!(err, SendError::Validation(ref message) if message.contains("unknown peer address transport")),
"unexpected error: {err:?}",
);
let err = runtime
.register_private_trusted_peer(crate::TrustedPeer {
name: schemeless_name.clone(),
pubkey: Keypair::generate().public_key(),
addr: "127.0.0.1:4205".to_string(),
meta: crate::PeerMeta::default(),
})
.await
.expect_err("schemeless TCP address must be rejected for private peers");
assert!(
matches!(err, SendError::Validation(ref message) if message.contains("missing transport scheme")),
"unexpected error: {err:?}",
);
assert!(
runtime
.trusted_peers
.read()
.peers
.iter()
.all(|peer| peer.name != unknown_name && peer.name != schemeless_name),
"invalid private peers must not be stored"
);
}
#[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,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await
.expect("add trusted peer");
CoreCommsRuntime::add_trusted_peer(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.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: peer_route(&receiver_name, receiver.public_key()),
body: "hello".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
},
)
.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_pubkey = Keypair::generate().public_key();
CoreCommsRuntime::add_trusted_peer(
&sender,
trusted_descriptor(&peer_name, peer_pubkey, "tcp://127.0.0.1:9"),
)
.await
.expect("add trusted peer");
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
blocks: None,
to: peer_route(&peer_name, peer_pubkey),
body: "hello".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
},
)
.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_send_to_untrusted_peer_surfaces_admission_dropped_not_offline() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("admit-sender-{suffix}");
let receiver_name = format!("admit-receiver-{suffix}");
let sender = CommsRuntime::inproc_only(&sender_name).unwrap();
let receiver = CommsRuntime::inproc_only(&receiver_name).unwrap();
CoreCommsRuntime::add_trusted_peer(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await
.expect("sender must trust receiver");
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
blocks: None,
to: peer_route(&receiver_name, receiver.public_key()),
body: "policy-rejected".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
},
)
.await;
assert!(
matches!(
result,
Err(SendError::AdmissionDropped {
reason: meerkat_core::comms::AdmissionDropReason::UntrustedSender
})
),
"untrusted-sender ingress drop must surface as typed AdmissionDropped, \
got: {result:?}"
);
let peers = CoreCommsRuntime::peers(&sender).await;
let entry = peers
.iter()
.find(|listed| listed.name.as_str() == receiver_name)
.expect("trusted peer should remain listed after admission drop");
assert_eq!(
entry.last_unreachable_reason,
Some(PeerReachabilityReason::AdmissionDropped),
"policy rejection must record AdmissionDropped, not OfflineOrNoAck"
);
}
#[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: missing_peer_route(&missing_name),
body: "hello".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
},
)
.await;
assert!(
matches!(result, Err(SendError::PeerNotFound(_))),
"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_pubkey = Keypair::generate().public_key();
let missing_peer_id = missing_pubkey.to_peer_id().as_str();
CoreCommsRuntime::add_trusted_peer(
&sender,
trusted_descriptor(
&missing_name,
missing_pubkey,
&format!("inproc://{missing_name}"),
),
)
.await
.expect("add trusted peer");
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
blocks: None,
to: peer_route(&missing_name, missing_pubkey),
body: "hello".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
},
)
.await;
assert!(
matches!(result, Err(SendError::PeerNotFound(ref peer)) if peer == &missing_peer_id),
"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(),
})
.expect("valid test peer should upsert");
}
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 {
PeerSendability::PeerMessage => CommsCommand::PeerMessage {
blocks: None,
to: PeerRoute::with_display_name(entry.peer_id, entry.name.clone()),
body: "truthfulness test".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
},
PeerSendability::PeerRequest => CommsCommand::PeerRequest {
to: PeerRoute::with_display_name(entry.peer_id, entry.name.clone()),
intent: "test".to_string(),
params: serde_json::json!({}),
blocks: None,
handling_mode: meerkat_core::types::HandlingMode::Queue,
stream: InputStreamMode::None,
},
PeerSendability::PeerResponse => CommsCommand::PeerResponse {
to: PeerRoute::with_display_name(entry.peer_id, entry.name.clone()),
in_reply_to: meerkat_core::InteractionId(Uuid::new_v4()),
status: meerkat_core::ResponseStatus::Completed,
result: serde_json::json!({}),
blocks: None,
handling_mode: None,
},
};
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::new(
sender,
receiver,
InteractionStreamLifecycleAuthority::LocalTransportOnly,
);
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(),
})
.expect("valid test peer should upsert");
}
CoreCommsRuntime::add_trusted_peer(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await
.expect("peer should accept trust entry for runtime");
let cmd = CommsCommand::PeerMessage {
blocks: None,
to: peer_route(&peer_name, peer.public_key()),
body: "not streamable".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
};
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 = trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
);
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().to_string();
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: peer_route(&receiver_name, receiver.public_key()),
body: "should fail".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
};
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_add_private_trusted_peer_zero_pubkey_is_rejected_at_trust_registration() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = CommsRuntime::inproc_only(&format!("rm-private-zero-{suffix}")).unwrap();
let peer_id = PeerId::new();
let peer_name = format!("private-legacy-peer-{suffix}");
let spec = TrustedPeerDescriptor::test_only_unsigned_typed(
peer_name.clone(),
peer_id,
format!("inproc://{peer_name}"),
)
.expect("valid zero-pubkey descriptor");
let result = CoreCommsRuntime::add_private_trusted_peer(&runtime, spec).await;
match result {
Err(SendError::Validation(message)) => {
assert!(
message.contains("pubkey") && message.contains("non-zero"),
"unexpected validation message: {message}"
);
}
other => panic!("private zero-pubkey trust registration must fail, got {other:?}"),
}
let removed = CoreCommsRuntime::remove_private_trusted_peer(&runtime, &peer_id.to_string())
.await
.expect("remove after rejected private add should still validate peer id");
assert!(!removed, "rejected private peer must not enter trust");
}
#[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().to_string();
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"
);
}
#[tokio::test]
async fn test_remove_trusted_peer_rejects_pubkey_string_argument() {
let suffix = Uuid::new_v4().simple().to_string();
let sender = CommsRuntime::inproc_only(&format!("rm-pubkey-sender-{suffix}")).unwrap();
let receiver = CommsRuntime::inproc_only(&format!("rm-pubkey-receiver-{suffix}")).unwrap();
let peer_pubkey_string = receiver.public_key().to_pubkey_string();
let result = CoreCommsRuntime::remove_trusted_peer(&sender, &peer_pubkey_string).await;
assert!(
matches!(result, Err(SendError::Validation(_))),
"remove_trusted_peer should validate its argument as PeerId, got: {result:?}"
);
}
}