#[cfg(not(target_arch = "wasm32"))]
use super::comms_config::ResolvedCommsConfig;
#[cfg(not(target_arch = "wasm32"))]
use crate::InboxSender;
use crate::agent::types::CommsMessage;
#[cfg(not(target_arch = "wasm32"))]
use crate::handle_connection;
#[cfg(target_arch = "wasm32")]
use crate::tokio;
use crate::trust::{TrustEntry, TrustStore};
use crate::{InprocRegistry, Keypair, PubKey, Router, TrustedPeersView};
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, CommsTrustMutation, CommsTrustMutationResult, EventStream,
GeneratedCommsTrustAuthoritySourceKind, InputStreamMode, PeerCapabilitySet, PeerDirectoryEntry,
PeerDirectorySource, PeerId, PeerName, 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, SessionId};
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,
}
#[cfg(not(target_arch = "wasm32"))]
const PAIRING_VERSION: u32 = 1;
#[cfg(not(target_arch = "wasm32"))]
const PAIRING_CHALLENGE_KIND: &str = "meerkat_pairing_challenge";
#[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, Clone, Copy, PartialEq, Eq, serde::Deserialize)]
enum PairingIdentityKind {
#[serde(rename = "ed25519_public_key")]
Ed25519PublicKey,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone)]
struct PairingPublicKey {
raw: String,
key: PubKey,
}
#[cfg(not(target_arch = "wasm32"))]
impl<'de> serde::Deserialize<'de> for PairingPublicKey {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let raw = String::deserialize(deserializer)?;
let key = PubKey::from_pubkey_string(&raw).map_err(serde::de::Error::custom)?;
Ok(Self { raw, key })
}
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, serde::Deserialize)]
#[serde(transparent)]
struct PairingAddress(String);
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, serde::Deserialize)]
struct PairingPeerIdentity {
kind: PairingIdentityKind,
public_key: PairingPublicKey,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, serde::Deserialize)]
struct PairingPeer {
name: String,
address: PairingAddress,
identity: PairingPeerIdentity,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, serde::Deserialize)]
#[serde(tag = "kind")]
enum PairingClientMessage {
#[serde(rename = "meerkat_pairing_hello")]
Hello { version: u32 },
#[serde(rename = "meerkat_pairing_proof")]
Proof {
caller: PairingPeer,
password_proof: String,
},
}
#[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)
}
#[cfg(any(not(target_arch = "wasm32"), test))]
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_to_trust_entry(descriptor: TrustedPeerDescriptor) -> Result<TrustEntry, 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(TrustEntry {
peer_id: descriptor.peer_id,
name: descriptor.name,
pubkey,
address: descriptor.address,
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 {
fn set_outbound_content_taint(
&self,
taint: Option<meerkat_core::comms::SenderContentTaint>,
) -> Result<(), meerkat_core::comms::SendError> {
CommsRuntime::set_outbound_content_taint(self, taint);
Ok(())
}
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.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")]
{
Some(format!("inproc://{}", self.participant_name()))
}
}
fn bridge_bootstrap_token(&self) -> Option<String> {
Some(self.bridge_bootstrap_token.clone())
}
async fn apply_trust_mutation(
&self,
mutation: CommsTrustMutation,
) -> Result<CommsTrustMutationResult, SendError> {
match mutation {
CommsTrustMutation::AddTrustedPeer { peer, authority } => {
self.validate_generated_trust_authority_owner(&authority)?;
authority
.validate_public_add(self.peer_id(), &peer)
.map_err(SendError::Validation)?;
let created = self
.apply_trusted_peer_descriptor(peer, authority.trust_row_owner_kind())
.await?;
Ok(CommsTrustMutationResult::Added { created })
}
CommsTrustMutation::RemoveTrustedPeer { peer_id, authority } => {
self.validate_generated_trust_authority_owner(&authority)?;
let parsed_peer_id = meerkat_core::comms::PeerId::parse(&peer_id)
.map_err(|err| SendError::Validation(err.to_string()))?;
authority
.validate_public_remove(self.peer_id(), parsed_peer_id)
.map_err(SendError::Validation)?;
let removed = self
.apply_trusted_peer_removal(&peer_id, authority.trust_row_owner_kind())
.await?;
Ok(CommsTrustMutationResult::Removed { removed })
}
CommsTrustMutation::AddPrivateTrustedPeer { peer, authority } => {
self.validate_generated_trust_authority_owner(&authority)?;
authority
.validate_private_add(self.peer_id(), &peer)
.map_err(SendError::Validation)?;
let created = self
.apply_private_trusted_peer_descriptor(peer, authority.trust_row_owner_kind())
.await?;
Ok(CommsTrustMutationResult::Added { created })
}
CommsTrustMutation::RemovePrivateTrustedPeer { peer_id, authority } => {
self.validate_generated_trust_authority_owner(&authority)?;
let parsed_peer_id = meerkat_core::comms::PeerId::parse(&peer_id)
.map_err(|err| SendError::Validation(err.to_string()))?;
authority
.validate_private_remove(self.peer_id(), parsed_peer_id)
.map_err(SendError::Validation)?;
let removed = self
.apply_private_trusted_peer_removal(&peer_id, authority.trust_row_owner_kind())
.await?;
Ok(CommsTrustMutationResult::Removed { removed })
}
}
}
async fn install_generated_mob_trust_owner(
&self,
owner: Arc<dyn std::any::Any + Send + Sync>,
) -> Result<(), SendError> {
let mut expected = self.mob_machine_trust_owner.write();
if let Some(existing) = expected.as_ref() {
if Arc::ptr_eq(existing, &owner) {
return Ok(());
}
return Err(SendError::Validation(
"target runtime is already bound to a different generated MobMachine trust owner"
.to_string(),
));
}
*expected = Some(owner);
Ok(())
}
async fn validate_recovered_generated_mob_trust_owner(
&self,
owner: Arc<dyn std::any::Any + Send + Sync>,
) -> Result<(), SendError> {
let expected = self.mob_machine_trust_owner.read();
if let Some(existing) = expected.as_ref()
&& !Arc::ptr_eq(existing, &owner)
{
return Err(SendError::Validation(
"target runtime is already bound to a different generated MobMachine trust owner"
.to_string(),
));
}
Ok(())
}
async fn install_recovered_generated_mob_trust_owner(
&self,
owner: Arc<dyn std::any::Any + Send + Sync>,
) -> Result<(), SendError> {
let mut expected = self.mob_machine_trust_owner.write();
if let Some(existing) = expected.as_ref() {
if Arc::ptr_eq(existing, &owner) {
return Ok(());
}
return Err(SendError::Validation(
"target runtime is already bound to a different generated MobMachine trust owner"
.to_string(),
));
}
*expected = Some(owner);
Ok(())
}
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 interaction_id = Uuid::new_v4();
let injector = self.event_injector_concrete();
let content = match blocks {
Some(blocks) => meerkat_core::types::ContentInput::Blocks(blocks),
None => meerkat_core::types::ContentInput::Text(body),
};
injector
.inject_with_interaction_id(
interaction_id,
content,
PlainEventSource::from(source),
handling_mode,
None,
)
.map_err(map_event_injector_error)?;
Ok(SendReceipt::InputAccepted {
interaction_id: meerkat_core::InteractionId(interaction_id),
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,
content_taint,
handling_mode,
objective_id,
} => self
.send_peer_command(
&to,
crate::types::MessageKind::Message {
body,
blocks,
content_taint: None,
handling_mode: Some(handling_mode),
objective_id,
},
content_taint,
)
.await
.map(|outcome| SendReceipt::PeerMessageSent {
envelope_id: outcome.envelope_id,
delivery: outcome.delivery,
}),
CommsCommand::PeerLifecycle { to, kind, params } => self
.send_peer_command(
&to,
crate::types::MessageKind::Lifecycle { kind, params },
None,
)
.await
.map(|outcome| SendReceipt::PeerLifecycleSent {
envelope_id: outcome.envelope_id,
}),
CommsCommand::PeerRequest {
to,
intent,
params,
blocks,
content_taint,
handling_mode,
stream,
objective_id,
} => {
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) {
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_send_failed(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,
content_taint: None,
handling_mode: Some(handling_mode),
objective_id,
},
content_taint,
)
.await
{
Ok(outcome) => outcome.envelope_id,
Err(e) => {
if stream_reserved {
self.expire_interaction_stream_on_send_failure(corr_id, interaction_id);
}
let _ = peer_handle.request_send_failed(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,
content_taint,
handling_mode,
objective_id,
} => {
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 = crate::Status::from(status);
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 reply_terminality =
peer_handle
.classify_response_reply(core_status)
.map_err(|err| {
SendError::Validation(format!(
"DSL rejected PeerResponse reply classification for corr_id {corr_id}: {err}"
))
})?;
let is_terminal_reply = matches!(
reply_terminality,
meerkat_core::TerminalityClass::Terminal { .. }
);
if !is_terminal_reply && handling_mode.is_some() {
return Err(SendError::Validation(
"handling_mode is forbidden on progress peer responses".to_string(),
));
}
let effective_handling_mode = if is_terminal_reply {
if peer_handle.inbound_state(corr_id)
!= Some(meerkat_core::InboundPeerRequestState::Received)
{
return Err(SendError::Validation(format!(
"PeerResponse terminal reply for corr_id {corr_id} is not in a repliable inbound state; response not sent"
)));
}
let resolved =
handling_mode.or_else(|| peer_handle.inbound_handling_mode(corr_id));
if resolved.is_none() {
return Err(SendError::Validation(format!(
"PeerResponse requires machine inbound peer request handling mode for corr_id {corr_id}"
)));
}
resolved
} else {
None
};
let envelope_id = self
.send_peer_command(
&to,
crate::types::MessageKind::Response {
in_reply_to: in_reply_to.0,
status,
result,
blocks,
content_taint: None,
handling_mode: effective_handling_mode,
objective_id,
},
content_taint,
)
.await?
.envelope_id;
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}"
)));
}
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,
content_taint,
handling_mode,
stream: InputStreamMode::ReserveInteraction,
objective_id,
} => {
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) {
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_send_failed(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,
content_taint: None,
handling_mode: Some(handling_mode),
objective_id,
},
content_taint,
)
.await
{
Ok(outcome) => outcome.envelope_id,
Err(e) => {
self.expire_interaction_stream_on_send_failure(corr_id, interaction_id);
let _ = peer_handle.request_send_failed(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 } => {
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 sender_taint = envelope.kind.content_taint();
let objective_id = envelope.kind.objective_id();
let content = match envelope.kind {
MessageKind::Message {
body,
blocks,
content_taint: _,
handling_mode: _,
objective_id: _,
} => meerkat_core::InteractionContent::Message { body, blocks },
MessageKind::Request {
intent,
params,
blocks,
content_taint: _,
handling_mode: _,
objective_id: _,
} => {
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,
content_taint: _,
handling_mode: _,
objective_id: _,
} => {
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;
}
};
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,
sender_taint,
objective_id,
},
ingress,
lifecycle_peer,
response_terminality: entry.response_terminality,
})
}
crate::types::InboxItem::PlainEvent {
body,
source,
handling_mode,
objective_id,
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,
sender_taint: None,
objective_id,
},
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 trusted_peers = self.trusted_peer_descriptors(true);
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,
})
}
async fn public_trusted_peer_projection_snapshot(
&self,
) -> Result<Vec<TrustedPeerDescriptor>, meerkat_core::CommsCapabilityError> {
Ok(self.trusted_peer_descriptors(false))
}
async fn trusted_peer_projection_snapshot_for_source(
&self,
source_kind: GeneratedCommsTrustAuthoritySourceKind,
) -> Result<Vec<TrustedPeerDescriptor>, meerkat_core::CommsCapabilityError> {
Ok(self.trusted_peer_descriptors_for_source(false, source_kind))
}
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(SessionId),
#[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),
#[error("Inproc registration rejected: {0:?}")]
InprocRegistrationRejected(crate::RegistrationRejection),
}
fn observe_inproc_registration(
namespace: &str,
name: &str,
outcome: crate::RegistrationOutcome,
) -> Result<(), CommsRuntimeError> {
use crate::RegistrationOutcome;
match outcome {
RegistrationOutcome::Registered => Ok(()),
RegistrationOutcome::ReplacedPubkey { evicted_name } => {
tracing::warn!(
inproc_namespace = %namespace,
peer_name = %name,
%evicted_name,
"inproc registration replaced an existing route under this pubkey"
);
Ok(())
}
RegistrationOutcome::EvictedName { evicted_pubkey } => {
tracing::warn!(
inproc_namespace = %namespace,
peer_name = %name,
evicted_pubkey = %evicted_pubkey.to_pubkey_string(),
"inproc registration evicted a stale pubkey bound to this name"
);
Ok(())
}
RegistrationOutcome::ReplacedPubkeyAndEvictedName {
evicted_name,
evicted_pubkey,
} => {
tracing::warn!(
inproc_namespace = %namespace,
peer_name = %name,
%evicted_name,
evicted_pubkey = %evicted_pubkey.to_pubkey_string(),
"inproc registration displaced both a name and a pubkey route"
);
Ok(())
}
RegistrationOutcome::Rejected { reason } => {
Err(CommsRuntimeError::InprocRegistrationRejected(reason))
}
}
}
#[cfg(not(target_arch = "wasm32"))]
impl From<SessionClaimError> for CommsRuntimeError {
fn from(err: SessionClaimError) -> Self {
match err {
SessionClaimError::SessionIdentityInUse(sid) => Self::SessionIdentityInUse(sid),
}
}
}
pub struct CommsToolMaterial {
router: Arc<Router>,
trusted_peers: TrustedPeersView,
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) -> &TrustedPeersView {
&self.trusted_peers
}
pub fn trusted_peers_shared(&self) -> TrustedPeersView {
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,
peer_id: PeerId,
router: Arc<Router>,
trusted_peers: Arc<parking_lot::RwLock<TrustStore>>,
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,
subscriber_registry: crate::event_injector::SubscriberRegistry,
interaction_stream_registry: InteractionStreamRegistry,
peer_comms_handle: crate::classify::PeerCommsHandleSlot,
meerkat_machine_trust_owner:
parking_lot::RwLock<Option<meerkat_core::comms::GeneratedPeerCommsOwnerToken>>,
mob_machine_trust_owner: parking_lot::RwLock<Option<Arc<dyn std::any::Any + Send + Sync>>>,
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>>>,
}
impl CommsRuntime {
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"))]
fn reject_persisted_trusted_peers_seed(
path: &std::path::Path,
) -> Result<(), CommsRuntimeError> {
let peers = TrustStore::load_or_default(path)
.map_err(|err| CommsRuntimeError::TrustLoadError(err.to_string()))?;
if peers.has_peers() {
return Err(CommsRuntimeError::TrustLoadError(format!(
"persisted trusted peer seed at {} is not authoritative; generated comms trust mutation authority required",
path.display()
)));
}
Ok(())
}
#[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()))?;
Self::reject_persisted_trusted_peers_seed(&config.trusted_peers_path)?;
let trusted_peers = TrustStore::new();
let public_key = keypair.public_key();
let trusted_peers = Arc::new(parking_lot::RwLock::new(trusted_peers));
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(),
peer_comms_handle: peer_comms_handle.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,
peer_id: public_key.to_peer_id(),
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,
subscriber_registry: crate::event_injector::new_subscriber_registry(),
interaction_stream_registry: Arc::new(Mutex::new(HashMap::new())),
peer_comms_handle,
meerkat_machine_trust_owner: parking_lot::RwLock::new(None),
mob_machine_trust_owner: parking_lot::RwLock::new(None),
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),
};
let namespace = config.inproc_namespace.clone().unwrap_or_default();
let name = config.name.clone();
let outcome = InprocRegistry::global().register_with_meta_in_namespace(
&namespace,
&name,
runtime.public_key,
inbox_sender,
crate::PeerMeta::default(),
);
observe_inproc_registration(&namespace, &name, outcome)?;
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(TrustStore::new()));
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(),
peer_comms_handle: peer_comms_handle.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,
peer_id: public_key.to_peer_id(),
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,
subscriber_registry: crate::event_injector::new_subscriber_registry(),
interaction_stream_registry: Arc::new(Mutex::new(HashMap::new())),
peer_comms_handle,
meerkat_machine_trust_owner: parking_lot::RwLock::new(None),
mob_machine_trust_owner: parking_lot::RwLock::new(None),
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),
};
let namespace_ref = namespace.as_deref().unwrap_or("");
let outcome = InprocRegistry::global().register_with_meta_in_namespace(
namespace_ref,
name,
runtime.public_key,
inbox_sender,
crate::PeerMeta::default(),
);
observe_inproc_registration(namespace_ref, name, outcome)?;
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(())
}
#[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(TrustStore::new()));
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(),
peer_comms_handle: peer_comms_handle.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,
peer_id: public_key.to_peer_id(),
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,
subscriber_registry: crate::event_injector::new_subscriber_registry(),
interaction_stream_registry: Arc::new(Mutex::new(HashMap::new())),
peer_comms_handle,
meerkat_machine_trust_owner: parking_lot::RwLock::new(None),
mob_machine_trust_owner: parking_lot::RwLock::new(None),
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),
};
let namespace_ref = namespace.as_deref().unwrap_or("");
let outcome = InprocRegistry::global().register_with_meta_in_namespace(
namespace_ref,
name,
runtime.public_key,
inbox_sender,
crate::PeerMeta::default(),
);
observe_inproc_registration(namespace_ref, name, outcome)?;
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);
}
#[cfg(test)]
fn install_peer_comms_handle(&self, handle: Arc<dyn meerkat_core::handles::PeerCommsHandle>) {
*self.meerkat_machine_trust_owner.write() = None;
*self.peer_comms_handle.write() = Some(handle);
}
fn validate_generated_trust_authority_owner(
&self,
authority: &meerkat_core::comms::CommsTrustMutationAuthority,
) -> Result<(), SendError> {
let expected_meerkat = self.meerkat_machine_trust_owner.read();
let expected_mob = self.mob_machine_trust_owner.read();
authority
.validate_target_source_owner_token(expected_meerkat.as_ref(), expected_mob.as_ref())
.map_err(SendError::Validation)
}
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),
content_taint,
handling_mode,
objective_id,
} => {
self.hydrate_blocks_for_transport(&mut blocks, "comms message")
.await?;
Ok(crate::types::MessageKind::Message {
body,
blocks: Some(blocks),
content_taint,
handling_mode,
objective_id,
})
}
crate::types::MessageKind::Request {
intent,
params,
blocks: Some(mut blocks),
content_taint,
handling_mode,
objective_id,
} => {
self.hydrate_blocks_for_transport(&mut blocks, "comms request")
.await?;
Ok(crate::types::MessageKind::Request {
intent,
params,
blocks: Some(blocks),
content_taint,
handling_mode,
objective_id,
})
}
crate::types::MessageKind::Response {
in_reply_to,
status,
result,
blocks: Some(mut blocks),
content_taint,
handling_mode,
objective_id,
} => {
self.hydrate_blocks_for_transport(&mut blocks, "comms response")
.await?;
Ok(crate::types::MessageKind::Response {
in_reply_to,
status,
result,
blocks: Some(blocks),
content_taint,
handling_mode,
objective_id,
})
}
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.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 fn set_outbound_content_taint(
&self,
taint: Option<meerkat_core::comms::SenderContentTaint>,
) {
self.router.set_outbound_content_taint(taint);
}
pub fn outbound_content_taint(&self) -> Option<meerkat_core::comms::SenderContentTaint> {
self.router.outbound_content_taint()
}
fn trusted_peer_descriptors(&self, include_private: bool) -> Vec<TrustedPeerDescriptor> {
self.trusted_peer_descriptors_filtered(include_private, None)
}
fn trusted_peer_descriptors_for_source(
&self,
include_private: bool,
source_kind: GeneratedCommsTrustAuthoritySourceKind,
) -> Vec<TrustedPeerDescriptor> {
self.trusted_peer_descriptors_filtered(include_private, Some(source_kind))
}
fn trusted_peer_descriptors_filtered(
&self,
include_private: bool,
source_kind: Option<GeneratedCommsTrustAuthoritySourceKind>,
) -> Vec<TrustedPeerDescriptor> {
let peer_rows: Vec<(meerkat_core::comms::PeerId, TrustEntry)> =
if let Some(source_kind) = source_kind {
self.router.trusted_peers_for_source(source_kind)
} else {
self.trusted_peers
.read()
.entries()
.map(|entry| (entry.peer_id, entry.clone()))
.collect()
};
let mut trusted_peers: Vec<TrustedPeerDescriptor> = peer_rows
.into_iter()
.filter_map(|(peer_id, entry)| {
if !include_private && self.router.is_private_peer_id(&peer_id) {
return None;
}
Some(TrustedPeerDescriptor {
peer_id,
name: entry.name,
address: entry.address,
pubkey: *entry.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()))
});
trusted_peers
}
async fn apply_trusted_peer_descriptor(
&self,
descriptor: TrustedPeerDescriptor,
source_kind: GeneratedCommsTrustAuthoritySourceKind,
) -> Result<bool, SendError> {
let entry = descriptor_to_trust_entry(descriptor)?;
let created = self
.router
.add_trusted_peer_for_source(entry, source_kind, false)
.map_err(|err| SendError::Validation(err.to_string()))?;
Ok(created)
}
async fn apply_trusted_peer_removal(
&self,
peer_id: &str,
source_kind: GeneratedCommsTrustAuthoritySourceKind,
) -> Result<bool, SendError> {
let peer_id = meerkat_core::comms::PeerId::parse(peer_id)
.map_err(|err| SendError::Validation(err.to_string()))?;
if self.router.is_private_peer_id(&peer_id)
&& !self.router.has_trust_source(&peer_id, source_kind)
{
return Err(SendError::Validation(format!(
"public trust authority cannot remove private trusted peer {peer_id}"
)));
}
Ok(self
.router
.remove_trusted_peer_for_source(&peer_id, source_kind))
}
async fn apply_private_trusted_peer_descriptor(
&self,
descriptor: TrustedPeerDescriptor,
source_kind: GeneratedCommsTrustAuthoritySourceKind,
) -> Result<bool, SendError> {
let entry = descriptor_to_trust_entry(descriptor)?;
let created = self
.router
.add_trusted_peer_for_source(entry, source_kind, true)
.map_err(|err| SendError::Validation(err.to_string()))?;
Ok(created)
}
async fn apply_private_trusted_peer_removal(
&self,
peer_id: &str,
source_kind: GeneratedCommsTrustAuthoritySourceKind,
) -> 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_for_source(&peer_id, source_kind))
}
pub fn set_peer_meta(&self, meta: crate::PeerMeta) {
let namespace = self.inproc_namespace().unwrap_or("");
let name = self.participant_name();
let outcome = InprocRegistry::global().register_with_meta_in_namespace(
namespace,
name,
self.public_key,
self.router.inbox_sender().clone(),
meta,
);
if outcome.is_rejected() {
tracing::warn!(
inproc_namespace = %namespace,
peer_name = %name,
?outcome,
"inproc metadata refresh rejected; peer meta was not updated"
);
} else if outcome.displaced_existing() {
tracing::warn!(
inproc_namespace = %namespace,
peer_name = %name,
?outcome,
"inproc metadata refresh displaced an unexpected route"
);
}
}
async fn resolve_peer_directory(&self) -> Vec<PeerDirectoryEntry> {
let resolved = self.resolved_peers_snapshot().await;
let sendable_kinds = self.peer_request_response_sendable_kinds();
let mut peers = Vec::new();
for peer in resolved {
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(),
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
}
async fn for_each_resolved_peer<F>(&self, mut on_peer: F) -> usize
where
F: FnMut(ResolvedPeer),
{
let participant_name = self.participant_name().to_string();
let private_peer_ids = self.router.private_peer_ids();
let mut emitted = 0usize;
{
let trusted = self.trusted_peers.read();
for entry in trusted.entries() {
if entry.name.as_str() == participant_name || entry.pubkey == self.public_key {
continue;
}
if private_peer_ids.contains(&entry.peer_id) {
continue;
}
on_peer(ResolvedPeer {
name: entry.name.clone(),
peer_id: entry.peer_id,
address: entry.address.clone(),
source: PeerDirectorySource::Trusted,
meta: entry.meta.clone(),
});
emitted += 1;
}
}
emitted
}
async fn send_peer_command(
&self,
route: &PeerRoute,
kind: crate::types::MessageKind,
content_taint: Option<meerkat_core::comms::SendTaintOverride>,
) -> Result<crate::router::SendOutcome, SendError> {
self.send_peer_command_with_id(route, Uuid::new_v4(), kind, content_taint)
.await
}
async fn send_peer_command_with_id(
&self,
route: &PeerRoute,
envelope_id: Uuid,
kind: crate::types::MessageKind,
content_taint: Option<meerkat_core::comms::SendTaintOverride>,
) -> Result<crate::router::SendOutcome, SendError> {
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, content_taint)
.await;
match result {
Ok(outcome) => Ok(outcome),
Err(crate::router::SendError::PeerNotFound(peer_id)) => {
Err(SendError::PeerNotFound(peer_id.to_string()))
}
Err(crate::router::SendError::PeerOffline) => Err(SendError::PeerOffline),
Err(crate::router::SendError::AdmissionDropped { reason }) => {
Err(SendError::AdmissionDropped {
reason: reason.into(),
})
}
Err(
error @ (crate::router::SendError::Transport(_) | crate::router::SendError::Io(_)),
) => Err(SendError::Transport(error.to_string())),
}
}
fn drop_peer_interaction_projection(&self, interaction_id: Uuid) {
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) -> TrustedPeersView {
TrustedPeersView::new(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(self.event_injector_concrete())
}
fn event_injector_concrete(&self) -> crate::CommsEventInjector {
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,
objective_id: 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 mut notified = std::pin::pin!(self.inbox_notify.notified());
notified.as_mut().enable();
{
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.as_mut().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::PeerCommsInstallTarget for CommsRuntime {
fn install_generated_peer_comms_handle(
&self,
install: meerkat_core::handles::GeneratedPeerCommsInstall,
) -> Result<(), String> {
let target_peer_id = self.generated_peer_comms_target_endpoint()?.peer_id;
if install.target_peer_id() != target_peer_id {
return Err(format!(
"generated peer-comms install targets peer_id {} but runtime peer_id is {}",
install.target_peer_id(),
target_peer_id
));
}
let owner_token = install.owner_token();
let mut expected = self.meerkat_machine_trust_owner.write();
if let Some(existing) = expected.as_ref()
&& !existing.same_owner(&owner_token)
{
return Err(
"target runtime is already bound to a different generated MeerkatMachine trust owner"
.to_string(),
);
}
*expected = Some(owner_token);
*self.peer_comms_handle.write() = Some(Arc::clone(install.peer_comms_handle()));
Ok(())
}
}
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>,
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, inbox_sender) = (keypair.clone(), inbox_sender.clone());
tokio::spawn(async move {
let _ = handle_connection(stream, require_peer_auth, &keypair, &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<TrustStore>>,
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<TrustStore>>,
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, 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, inbox_sender)
.await
.map_err(|err| std::io::Error::other(err.to_string()))
}
#[cfg(not(target_arch = "wasm32"))]
async fn read_pairing_message(
reader: &mut BufReader<TcpStream>,
buf: &mut String,
) -> Result<PairingClientMessage, 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 message: {err}"),
)
})
}
#[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<TrustStore>>,
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 PairingClientMessage::Hello { version } =
read_pairing_message(&mut reader, &mut line).await?
else {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"expected meerkat pairing hello",
));
};
if version != 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 PairingClientMessage::Proof {
caller,
password_proof,
} = read_pairing_message(&mut reader, &mut line).await?
else {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"expected meerkat pairing proof",
));
};
let PairingIdentityKind::Ed25519PublicKey = caller.identity.kind;
let expected = comms_pairing_password_proof(
pairing_password,
&challenge,
&caller.identity.public_key.raw,
&caller.address.0,
);
if !constant_time_str_eq(&password_proof, &expected) {
return Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"invalid pairing password proof",
));
}
let pubkey = caller.identity.public_key.key;
let caller_name = PeerName::new(caller.name.clone())
.map_err(|err| std::io::Error::new(std::io::ErrorKind::InvalidData, err))?;
let caller_address = meerkat_core::comms::PeerAddress::parse(&caller.address.0)
.map_err(|err| std::io::Error::new(std::io::ErrorKind::InvalidData, err.to_string()))?;
let paired_peer_id = pubkey.to_peer_id();
router
.add_trusted_peer_for_source(
TrustEntry {
peer_id: paired_peer_id,
name: caller_name,
pubkey,
address: caller_address,
meta: crate::PeerMeta::default(),
},
meerkat_core::comms::GeneratedCommsTrustAuthoritySourceKind::MobMachineExternalPeerTrustWiring,
false,
)
.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 faults = Arc::new(crate::plain_listener::PlainIngressFaults::default());
let handle = tokio::spawn(async move {
while let Ok((stream, _peer)) = listener.accept().await {
let sender = inbox_sender.clone();
let sem = semaphore.clone();
let faults = faults.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,
&faults,
)
.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 faults = Arc::new(crate::plain_listener::PlainIngressFaults::default());
let handle = tokio::spawn(async move {
while let Ok((stream, _)) = listener.accept().await {
let sender = inbox_sender.clone();
let sem = semaphore.clone();
let faults = faults.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,
&faults,
)
.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::classify::test_support;
#[tokio::test]
async fn dyn_trait_taint_setter_delegates_to_router_declaration() {
let runtime =
Arc::new(CommsRuntime::inproc_only("dyn-taint-setter").expect("inproc runtime"));
let dyn_runtime: Arc<dyn CoreCommsRuntime> = runtime.clone();
dyn_runtime
.set_outbound_content_taint(Some(meerkat_core::comms::SenderContentTaint::Tainted))
.expect("concrete runtime must accept the declaration");
assert_eq!(
runtime.outbound_content_taint(),
Some(meerkat_core::comms::SenderContentTaint::Tainted)
);
dyn_runtime
.set_outbound_content_taint(None)
.expect("clearing the declaration must succeed");
assert_eq!(runtime.outbound_content_taint(), None);
}
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::{
CommsTrustMutation, CommsTrustMutationResult, InputSource, InputStreamMode,
PeerDirectorySource, PeerId, PeerName, PeerRoute, PeerSendability, StreamError,
StreamScope, TrustedPeerDescriptor,
},
interaction::InteractionId,
types::{ContentBlock, ImageData, SessionId},
};
#[test]
fn pairing_caller_identity_is_typed_and_fails_closed() {
let bad_key = serde_json::json!({
"name": "peer",
"address": "tcp://127.0.0.1:4200",
"identity": {
"kind": "ed25519_public_key",
"public_key": "not-a-valid-pubkey",
}
});
assert!(
serde_json::from_value::<PairingPeer>(bad_key).is_err(),
"a malformed pairing public key must be rejected at deserialization"
);
let valid_pubkey = Keypair::generate().public_key();
let wrong_kind = serde_json::json!({
"name": "peer",
"address": "tcp://127.0.0.1:4200",
"identity": {
"kind": "rsa_public_key",
"public_key": valid_pubkey.to_pubkey_string(),
}
});
assert!(
serde_json::from_value::<PairingPeer>(wrong_kind).is_err(),
"a non-ed25519 identity kind must be rejected at deserialization"
);
let good = serde_json::json!({
"name": "peer",
"address": "tcp://127.0.0.1:4200",
"identity": {
"kind": "ed25519_public_key",
"public_key": valid_pubkey.to_pubkey_string(),
}
});
let caller: PairingPeer =
serde_json::from_value(good).expect("well-formed pairing caller deserializes");
assert_eq!(caller.identity.public_key.key, valid_pubkey);
assert_eq!(
caller.identity.public_key.raw,
valid_pubkey.to_pubkey_string()
);
assert_eq!(caller.address.0, "tcp://127.0.0.1:4200");
assert_eq!(caller.identity.kind, PairingIdentityKind::Ed25519PublicKey);
}
#[test]
fn pairing_client_messages_parse_typed_and_fail_closed() {
assert!(
serde_json::from_value::<PairingClientMessage>(serde_json::json!({
"kind": "meerkat_pairing_evil",
"version": 1,
}))
.is_err(),
"an unknown pairing message kind must be rejected"
);
assert!(
serde_json::from_value::<PairingClientMessage>(serde_json::json!({
"kind": "meerkat_pairing_hello",
}))
.is_err(),
"a pairing hello without a version must be rejected"
);
let pubkey = Keypair::generate().public_key();
assert!(
serde_json::from_value::<PairingClientMessage>(serde_json::json!({
"kind": "meerkat_pairing_proof",
"caller": {
"name": "peer",
"address": "tcp://127.0.0.1:4200",
"identity": {
"kind": "ed25519_public_key",
"public_key": pubkey.to_pubkey_string(),
}
}
}))
.is_err(),
"a pairing proof without a password proof must be rejected"
);
let hello: PairingClientMessage = serde_json::from_value(serde_json::json!({
"kind": "meerkat_pairing_hello",
"version": 1,
}))
.expect("well-formed hello parses");
assert!(matches!(hello, PairingClientMessage::Hello { version: 1 }));
let proof: PairingClientMessage = serde_json::from_value(serde_json::json!({
"kind": "meerkat_pairing_proof",
"password_proof": "proof",
"caller": {
"name": "peer",
"address": "tcp://127.0.0.1:4200",
"identity": {
"kind": "ed25519_public_key",
"public_key": pubkey.to_pubkey_string(),
}
}
}))
.expect("well-formed proof parses");
match proof {
PairingClientMessage::Proof {
caller,
password_proof,
} => {
assert_eq!(password_proof, "proof");
assert_eq!(caller.identity.public_key.key, pubkey);
}
PairingClientMessage::Hello { .. } => panic!("expected proof variant"),
}
}
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 local_descriptor_for_runtime(runtime: &CommsRuntime) -> TrustedPeerDescriptor {
trusted_descriptor(
runtime.participant_name(),
runtime.public_key(),
&runtime.advertised_address(),
)
}
fn test_trust_entry(name: &str, pubkey: PubKey, address: &str) -> TrustEntry {
TrustEntry {
peer_id: crate::router::peer_id_from_pubkey(&pubkey),
name: PeerName::new(name.to_string()).expect("valid peer name"),
pubkey,
address: parse_peer_address(address).expect("valid peer address"),
meta: crate::PeerMeta::default(),
}
}
struct TestProjectionTrustAuthority {
dsl: Arc<meerkat_runtime::HandleDslAuthority>,
authority: meerkat_core::comms::CommsTrustMutationAuthority,
}
impl TestProjectionTrustAuthority {
fn install_owner(&self, runtime: &CommsRuntime) {
meerkat_runtime::RuntimePeerCommsHandle::install_generated_on(
Arc::clone(&self.dsl),
runtime,
)
.expect("install generated peer-comms handle");
}
}
fn try_test_projection_add_authority(
local: &TrustedPeerDescriptor,
peer: &TrustedPeerDescriptor,
epoch: u64,
) -> Result<TestProjectionTrustAuthority, String> {
let local_endpoint = meerkat_runtime::meerkat_machine::dsl::PeerEndpoint::from(local);
let endpoint = meerkat_runtime::meerkat_machine::dsl::PeerEndpoint::from(peer);
let target_epoch = epoch.max(1);
let dsl = Arc::new(meerkat_runtime::HandleDslAuthority::ephemeral());
dsl.apply_signal(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineSignal::Initialize,
"test_projection_add_authority::initialize",
)
.map_err(|err| err.to_string())?;
dsl.apply_input(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::RegisterSession {
session_id: meerkat_runtime::meerkat_machine::dsl::SessionId::from(
"meerkat-comms-runtime-test-projection",
),
},
"test_projection_add_authority::register_session",
)
.map_err(|err| err.to_string())?;
dsl.apply_input(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::PrepareBindings {
agent_runtime_id: meerkat_runtime::meerkat_machine::dsl::AgentRuntimeId::from(
"meerkat-comms-runtime-test-projection-runtime",
),
fence_token: meerkat_runtime::meerkat_machine::dsl::FenceToken::from(0),
generation: Some(meerkat_runtime::meerkat_machine::dsl::Generation::from(0)),
runtime_epoch_id: None,
session_id: meerkat_runtime::meerkat_machine::dsl::SessionId::from(
"meerkat-comms-runtime-test-projection",
),
},
"test_projection_add_authority::prepare_bindings",
)
.map_err(|err| err.to_string())?;
dsl.apply_input(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::PublishLocalEndpoint {
endpoint: local_endpoint,
},
"test_projection_add_authority::publish_local_endpoint",
)
.map_err(|err| err.to_string())?;
for filler_index in 1..target_epoch {
let filler_key = Keypair::generate().public_key();
let filler = trusted_descriptor(
&format!("filler-add-{filler_index}"),
filler_key,
&format!("inproc://filler-add-{filler_index}"),
);
let filler_endpoint =
meerkat_runtime::meerkat_machine::dsl::PeerEndpoint::from(&filler);
dsl.apply_input(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::AddDirectPeerEndpoint {
endpoint: filler_endpoint,
},
"test_projection_add_authority::add_filler_endpoint",
)
.map_err(|err| err.to_string())?;
}
let transition = dsl
.apply_input_with_transition(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::AddDirectPeerEndpoint {
endpoint: endpoint.clone(),
},
"test_projection_add_authority::add_endpoint",
)
.map_err(|err| err.to_string())?;
let mut obligations =
meerkat_runtime::protocol_comms_trust_reconcile::extract_obligations_with_freshness(
&transition,
dsl.peer_projection_freshness_authority(),
);
let obligation = obligations
.pop()
.ok_or_else(|| "generated reconcile obligation missing".to_string())?;
let authority = meerkat_runtime::protocol_comms_trust_reconcile::authority_for_endpoint(
&obligation,
&endpoint,
)?;
Ok(TestProjectionTrustAuthority { dsl, authority })
}
fn test_projection_add_authority(
local: &TrustedPeerDescriptor,
peer: &TrustedPeerDescriptor,
epoch: u64,
) -> TestProjectionTrustAuthority {
try_test_projection_add_authority(local, peer, epoch)
.expect("generated peer projection add authority")
}
fn test_projection_add_authority_with_dsl(
dsl: Arc<meerkat_runtime::HandleDslAuthority>,
peer: &TrustedPeerDescriptor,
) -> TestProjectionTrustAuthority {
let endpoint = meerkat_runtime::meerkat_machine::dsl::PeerEndpoint::from(peer);
let transition = dsl
.apply_input_with_transition(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::AddDirectPeerEndpoint {
endpoint: endpoint.clone(),
},
"test_projection_add_authority_with_dsl::add_endpoint",
)
.expect("AddDirectPeerEndpoint input");
let mut obligations =
meerkat_runtime::protocol_comms_trust_reconcile::extract_obligations_with_freshness(
&transition,
dsl.peer_projection_freshness_authority(),
);
let obligation = obligations.pop().expect("generated reconcile obligation");
let authority = meerkat_runtime::protocol_comms_trust_reconcile::authority_for_endpoint(
&obligation,
&endpoint,
)
.expect("generated peer projection add authority");
TestProjectionTrustAuthority { dsl, authority }
}
fn test_projection_remove_authority(
local: &TrustedPeerDescriptor,
peer: &TrustedPeerDescriptor,
epoch: u64,
) -> TestProjectionTrustAuthority {
let local_endpoint = meerkat_runtime::meerkat_machine::dsl::PeerEndpoint::from(local);
let endpoint = meerkat_runtime::meerkat_machine::dsl::PeerEndpoint::from(peer);
let target_epoch = epoch.max(2);
let dsl = Arc::new(meerkat_runtime::HandleDslAuthority::ephemeral());
dsl.apply_signal(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineSignal::Initialize,
"test_projection_remove_authority::initialize",
)
.expect("Initialize signal");
dsl.apply_input(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::RegisterSession {
session_id: meerkat_runtime::meerkat_machine::dsl::SessionId::from(
"meerkat-comms-runtime-test-projection",
),
},
"test_projection_remove_authority::register_session",
)
.expect("RegisterSession input");
dsl.apply_input(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::PrepareBindings {
agent_runtime_id: meerkat_runtime::meerkat_machine::dsl::AgentRuntimeId::from(
"meerkat-comms-runtime-test-projection-runtime",
),
fence_token: meerkat_runtime::meerkat_machine::dsl::FenceToken::from(0),
generation: Some(meerkat_runtime::meerkat_machine::dsl::Generation::from(0)),
runtime_epoch_id: None,
session_id: meerkat_runtime::meerkat_machine::dsl::SessionId::from(
"meerkat-comms-runtime-test-projection",
),
},
"test_projection_remove_authority::prepare_bindings",
)
.expect("PrepareBindings input");
dsl.apply_input(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::PublishLocalEndpoint {
endpoint: local_endpoint,
},
"test_projection_remove_authority::publish_local_endpoint",
)
.expect("PublishLocalEndpoint input");
dsl.apply_input(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::AddDirectPeerEndpoint {
endpoint: endpoint.clone(),
},
"test_projection_remove_authority::add_endpoint",
)
.expect("AddDirectPeerEndpoint input");
for filler_index in 2..target_epoch {
let filler_key = Keypair::generate().public_key();
let filler = trusted_descriptor(
&format!("filler-remove-{filler_index}"),
filler_key,
&format!("inproc://filler-remove-{filler_index}"),
);
let filler_endpoint =
meerkat_runtime::meerkat_machine::dsl::PeerEndpoint::from(&filler);
dsl.apply_input(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::AddDirectPeerEndpoint {
endpoint: filler_endpoint,
},
"test_projection_remove_authority::add_filler_endpoint",
)
.expect("AddDirectPeerEndpoint filler input");
}
let transition = dsl
.apply_input_with_transition(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::RemoveDirectPeerEndpoint {
endpoint: endpoint.clone(),
},
"test_projection_remove_authority::remove_endpoint",
)
.expect("RemoveDirectPeerEndpoint input");
let mut obligations =
meerkat_runtime::protocol_comms_trust_reconcile::extract_obligations_with_freshness(
&transition,
dsl.peer_projection_freshness_authority(),
);
let obligation = obligations.pop().expect("generated reconcile obligation");
let authority =
meerkat_runtime::protocol_comms_trust_reconcile::removal_authority_for_peer_id(
&obligation,
&endpoint.peer_id.0,
)
.expect("generated peer projection remove authority");
TestProjectionTrustAuthority { dsl, authority }
}
fn test_projection_remove_authority_with_dsl(
dsl: Arc<meerkat_runtime::HandleDslAuthority>,
peer: &TrustedPeerDescriptor,
) -> TestProjectionTrustAuthority {
let endpoint = meerkat_runtime::meerkat_machine::dsl::PeerEndpoint::from(peer);
let transition = dsl
.apply_input_with_transition(
meerkat_runtime::meerkat_machine::dsl::MeerkatMachineInput::RemoveDirectPeerEndpoint {
endpoint: endpoint.clone(),
},
"test_projection_remove_authority_with_dsl::remove_endpoint",
)
.expect("RemoveDirectPeerEndpoint input");
let mut obligations =
meerkat_runtime::protocol_comms_trust_reconcile::extract_obligations_with_freshness(
&transition,
dsl.peer_projection_freshness_authority(),
);
let obligation = obligations.pop().expect("generated reconcile obligation");
let authority =
meerkat_runtime::protocol_comms_trust_reconcile::removal_authority_for_peer_id(
&obligation,
&endpoint.peer_id.0,
)
.expect("generated peer projection remove authority");
TestProjectionTrustAuthority { dsl, authority }
}
async fn try_add_trusted_peer_with_generated_authority(
runtime: &CommsRuntime,
peer: TrustedPeerDescriptor,
) -> Result<Arc<meerkat_runtime::HandleDslAuthority>, SendError> {
let projection =
try_test_projection_add_authority(&local_descriptor_for_runtime(runtime), &peer, 0)
.map_err(SendError::Validation)?;
let dsl = Arc::clone(&projection.dsl);
projection.install_owner(runtime);
CoreCommsRuntime::apply_trust_mutation(
runtime,
CommsTrustMutation::AddTrustedPeer {
authority: projection.authority,
peer,
},
)
.await
.map(|_| dsl)
}
async fn add_trusted_peer_with_generated_authority(
runtime: &CommsRuntime,
peer: TrustedPeerDescriptor,
) -> Arc<meerkat_runtime::HandleDslAuthority> {
try_add_trusted_peer_with_generated_authority(runtime, peer)
.await
.expect("generated add trusted peer should apply")
}
async fn remove_trusted_peer_with_generated_authority(
runtime: &CommsRuntime,
peer: &TrustedPeerDescriptor,
) -> Result<bool, SendError> {
let projection =
test_projection_remove_authority(&local_descriptor_for_runtime(runtime), peer, 0);
projection.install_owner(runtime);
let result = CoreCommsRuntime::apply_trust_mutation(
runtime,
CommsTrustMutation::RemoveTrustedPeer {
authority: projection.authority,
peer_id: peer.peer_id.to_string(),
},
)
.await?;
match result {
CommsTrustMutationResult::Removed { removed } => Ok(removed),
other => Err(SendError::Validation(format!(
"unexpected generated trust mutation result: {other:?}"
))),
}
}
#[tokio::test]
async fn generated_peer_projection_trust_requires_canonical_owner() {
let runtime = CommsRuntime::inproc_only("trust-owner-mismatch").expect("runtime");
let local = local_descriptor_for_runtime(&runtime);
let peer_key = Keypair::generate().public_key();
let descriptor = trusted_descriptor("peer", peer_key, "inproc://peer");
let authority = test_projection_add_authority(&local, &descriptor, 3);
let missing_owner = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::AddTrustedPeer {
peer: descriptor.clone(),
authority: authority.authority.clone(),
},
)
.await
.expect_err("generated authority without runtime owner must fail closed");
assert!(
matches!(missing_owner, SendError::Validation(ref message) if message.contains("requires the target runtime's generated owner token")),
"unexpected missing-owner rejection: {missing_owner:?}"
);
let foreign_owner = test_projection_add_authority(&local, &descriptor, 4);
foreign_owner.install_owner(&runtime);
let mismatched_owner = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::AddTrustedPeer {
peer: descriptor,
authority: authority.authority,
},
)
.await
.expect_err("foreign generated owner must fail closed");
assert!(
matches!(mismatched_owner, SendError::Validation(ref message) if message.contains("different generated owner")),
"unexpected owner-mismatch rejection: {mismatched_owner:?}"
);
assert!(
runtime.peers().await.is_empty(),
"owner mismatch must not mutate peer directory projection"
);
}
#[tokio::test]
async fn recovered_mob_trust_owner_cannot_replace_live_owner() {
let runtime = CommsRuntime::inproc_only("mob-trust-owner-recovery-guard").expect("runtime");
let live_owner: Arc<dyn std::any::Any + Send + Sync> = Arc::new("live-owner");
let recovered_owner: Arc<dyn std::any::Any + Send + Sync> = Arc::new("recovered-owner");
CoreCommsRuntime::validate_recovered_generated_mob_trust_owner(
&runtime,
Arc::clone(&recovered_owner),
)
.await
.expect("validation without installed owner should pass without binding");
CoreCommsRuntime::install_generated_mob_trust_owner(&runtime, Arc::clone(&live_owner))
.await
.expect("install live mob owner");
let validation_err = CoreCommsRuntime::validate_recovered_generated_mob_trust_owner(
&runtime,
Arc::clone(&recovered_owner),
)
.await
.expect_err("read-only validation must not have installed recovered owner");
assert!(
matches!(validation_err, SendError::Validation(ref message) if message.contains("different generated MobMachine trust owner")),
"unexpected recovered owner validation rejection: {validation_err:?}"
);
CoreCommsRuntime::install_recovered_generated_mob_trust_owner(
&runtime,
Arc::clone(&live_owner),
)
.await
.expect("same recovered owner is idempotent");
let err = CoreCommsRuntime::install_recovered_generated_mob_trust_owner(
&runtime,
recovered_owner,
)
.await
.expect_err("recovered owner must not replace a different live owner");
assert!(
matches!(err, SendError::Validation(ref message) if message.contains("different generated MobMachine trust owner")),
"unexpected recovered owner rejection: {err:?}"
);
}
#[tokio::test]
async fn generated_trust_mutation_updates_projection() {
let runtime = CommsRuntime::inproc_only("trust-mutator-generated").expect("runtime");
let local = local_descriptor_for_runtime(&runtime);
let peer_key = Keypair::generate().public_key();
let descriptor = trusted_descriptor("peer", peer_key, "inproc://peer");
let peer_id = descriptor.peer_id.to_string();
let add_authority = test_projection_add_authority(&local, &descriptor, 7);
add_authority.install_owner(&runtime);
let add_authority_replay = add_authority.authority.clone();
let add = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::AddTrustedPeer {
peer: descriptor.clone(),
authority: add_authority.authority,
},
)
.await
.expect("generated add should apply");
assert_eq!(add, CommsTrustMutationResult::Added { created: true });
assert_eq!(runtime.peers().await.len(), 1);
let replay = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::AddTrustedPeer {
peer: descriptor.clone(),
authority: add_authority_replay,
},
)
.await
.expect_err("cloned generated add authority must be one-use");
assert!(
matches!(replay, SendError::Validation(ref message) if message.contains("already consumed")),
"unexpected replay rejection: {replay:?}"
);
let wrong_peer_key = Keypair::generate().public_key();
let wrong_peer = trusted_descriptor(
"wrong-operation-peer",
wrong_peer_key,
"inproc://wrong-operation-peer",
);
let wrong_operation_authority =
test_projection_add_authority_with_dsl(Arc::clone(&add_authority.dsl), &wrong_peer);
let wrong_operation = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::RemoveTrustedPeer {
peer_id: peer_id.clone(),
authority: wrong_operation_authority.authority,
},
)
.await
.expect_err("generated add authority must not remove trust");
assert!(
matches!(wrong_operation, SendError::Validation(ref message) if message.contains("PublicAdd") && message.contains("remove a public trusted peer")),
"unexpected wrong-operation rejection: {wrong_operation:?}"
);
assert_eq!(
runtime.peers().await.len(),
1,
"wrong-operation authority must not mutate trust projection",
);
let remove_authority =
test_projection_remove_authority_with_dsl(Arc::clone(&add_authority.dsl), &descriptor);
let remove = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::RemoveTrustedPeer {
peer_id: peer_id.clone(),
authority: remove_authority.authority,
},
)
.await
.expect("generated remove should apply");
assert_eq!(remove, CommsTrustMutationResult::Removed { removed: true });
assert!(runtime.peers().await.is_empty());
}
#[tokio::test]
async fn generated_trust_add_reports_existing_source_without_rewriting_descriptor() {
let runtime = CommsRuntime::inproc_only("trust-mutator-existing-source").expect("runtime");
let peer_key = Keypair::generate().public_key();
let descriptor = trusted_descriptor("peer", peer_key, "inproc://peer");
let peer = descriptor_to_trust_entry(descriptor.clone()).expect("valid descriptor");
runtime
.router
.add_trusted_peer_for_source(
peer,
GeneratedCommsTrustAuthoritySourceKind::MeerkatMachinePeerProjection,
false,
)
.expect("test setup installs MeerkatMachine-owned trust");
let same_peer = descriptor_to_trust_entry(descriptor.clone()).expect("valid descriptor");
let same = runtime
.router
.add_trusted_peer_for_source(
same_peer,
GeneratedCommsTrustAuthoritySourceKind::MeerkatMachinePeerProjection,
false,
)
.expect("same generated source row should be idempotent");
assert!(!same);
let conflicting_descriptor =
trusted_descriptor("peer-renamed", peer_key, "inproc://peer-renamed");
let conflicting_peer =
descriptor_to_trust_entry(conflicting_descriptor).expect("valid descriptor");
let conflict = runtime
.router
.add_trusted_peer_for_source(
conflicting_peer,
GeneratedCommsTrustAuthoritySourceKind::MeerkatMachinePeerProjection,
false,
)
.expect_err("trust repair must not rewrite an existing generated source row");
assert!(
matches!(
conflict,
crate::trust::TrustError::ConflictingGeneratedTrustSource { .. }
),
"unexpected conflicting-source rejection: {conflict:?}"
);
let public_snapshot = CoreCommsRuntime::public_trusted_peer_projection_snapshot(&runtime)
.await
.expect("public trust snapshot");
assert_eq!(public_snapshot.len(), 1);
assert_eq!(public_snapshot[0].name.as_str(), "peer");
assert_eq!(public_snapshot[0].address.to_string(), "inproc://peer");
}
#[tokio::test]
async fn generated_remove_does_not_remove_trust_owned_by_other_source() {
let runtime =
CommsRuntime::inproc_only("trust-mutator-source-scoped-remove").expect("runtime");
let local = local_descriptor_for_runtime(&runtime);
let peer_key = Keypair::generate().public_key();
let descriptor = trusted_descriptor("mob-owned-peer", peer_key, "inproc://mob-owned-peer");
let peer_id = descriptor.peer_id.to_string();
let peer = descriptor_to_trust_entry(descriptor.clone()).expect("valid descriptor");
runtime
.router
.add_trusted_peer_for_source(
peer,
GeneratedCommsTrustAuthoritySourceKind::MobMachineMemberTrustWiring,
false,
)
.expect("test setup installs mob-owned public trust");
let meerkat_projection = CoreCommsRuntime::trusted_peer_projection_snapshot_for_source(
&runtime,
GeneratedCommsTrustAuthoritySourceKind::MeerkatMachinePeerProjection,
)
.await
.expect("source-scoped snapshot");
assert!(
meerkat_projection.is_empty(),
"MeerkatMachine peer projection must not see MobMachine-owned trust as its canonical rows",
);
let remove_authority = test_projection_remove_authority(&local, &descriptor, 8);
remove_authority.install_owner(&runtime);
let remove = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::RemoveTrustedPeer {
peer_id,
authority: remove_authority.authority,
},
)
.await
.expect("generated remove for an unowned source should apply as no-op");
assert_eq!(remove, CommsTrustMutationResult::Removed { removed: false });
let public_snapshot = CoreCommsRuntime::public_trusted_peer_projection_snapshot(&runtime)
.await
.expect("public trust snapshot");
assert_eq!(
public_snapshot.len(),
1,
"MobMachine-owned trust must remain after a MeerkatMachine removal no-op",
);
assert_eq!(public_snapshot[0].peer_id, descriptor.peer_id);
}
#[tokio::test]
async fn source_scoped_snapshot_uses_source_owned_descriptor() {
let runtime =
CommsRuntime::inproc_only("trust-mutator-source-scoped-descriptor").expect("runtime");
let peer_key = Keypair::generate().public_key();
let projection_descriptor =
trusted_descriptor("projection-owned", peer_key, "inproc://projection-owned");
let mob_descriptor = trusted_descriptor("mob-owned", peer_key, "inproc://mob-owned");
let peer_id = projection_descriptor.peer_id;
let projection_peer =
descriptor_to_trust_entry(projection_descriptor.clone()).expect("valid descriptor");
let mob_peer = descriptor_to_trust_entry(mob_descriptor.clone()).expect("valid descriptor");
runtime
.router
.add_trusted_peer_for_source(
projection_peer,
GeneratedCommsTrustAuthoritySourceKind::MeerkatMachinePeerProjection,
false,
)
.expect("test setup installs MeerkatMachine-owned trust");
runtime
.router
.add_trusted_peer_for_source(
mob_peer,
GeneratedCommsTrustAuthoritySourceKind::MobMachineMemberTrustWiring,
false,
)
.expect("test setup installs MobMachine-owned trust");
let projection_snapshot = CoreCommsRuntime::trusted_peer_projection_snapshot_for_source(
&runtime,
GeneratedCommsTrustAuthoritySourceKind::MeerkatMachinePeerProjection,
)
.await
.expect("projection source snapshot");
assert_eq!(projection_snapshot, vec![projection_descriptor.clone()]);
let mob_snapshot = CoreCommsRuntime::trusted_peer_projection_snapshot_for_source(
&runtime,
GeneratedCommsTrustAuthoritySourceKind::MobMachineMemberTrustWiring,
)
.await
.expect("mob source snapshot");
assert_eq!(mob_snapshot, vec![mob_descriptor.clone()]);
let remove = runtime.router.remove_trusted_peer_for_source(
&peer_id,
GeneratedCommsTrustAuthoritySourceKind::MobMachineMemberTrustWiring,
);
assert!(remove, "test setup removal should remove MobMachine source");
let public_snapshot = CoreCommsRuntime::public_trusted_peer_projection_snapshot(&runtime)
.await
.expect("public trust snapshot");
assert_eq!(
public_snapshot,
vec![projection_descriptor],
"union projection should fall back to the remaining source-owned descriptor",
);
}
#[tokio::test]
async fn generated_add_authority_rejects_descriptor_drift() {
let runtime = CommsRuntime::inproc_only("trust-mutator-descriptor-drift").expect("runtime");
let local = local_descriptor_for_runtime(&runtime);
let peer_key = Keypair::generate().public_key();
let descriptor = trusted_descriptor("peer", peer_key, "inproc://peer");
let drifted_descriptor =
trusted_descriptor("peer", peer_key, "inproc://peer-drifted-address");
let add_authority = test_projection_add_authority(&local, &descriptor, 11);
add_authority.install_owner(&runtime);
let err = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::AddTrustedPeer {
peer: drifted_descriptor,
authority: add_authority.authority,
},
)
.await
.expect_err("generated add authority must bind the full peer descriptor");
assert!(
matches!(err, SendError::Validation(ref message) if message.contains("descriptor")),
"unexpected descriptor drift rejection: {err:?}"
);
assert!(
runtime.peers().await.is_empty(),
"descriptor drift must not mutate public trust"
);
}
#[tokio::test]
async fn public_trust_projection_excludes_and_cannot_remove_private_peer() {
let runtime = CommsRuntime::inproc_only("trust-mutator-private-split").expect("runtime");
let local = local_descriptor_for_runtime(&runtime);
let peer_key = Keypair::generate().public_key();
let descriptor = trusted_descriptor("private-peer", peer_key, "inproc://private-peer");
let peer_id = descriptor.peer_id.to_string();
runtime
.apply_private_trusted_peer_descriptor(
descriptor.clone(),
GeneratedCommsTrustAuthoritySourceKind::MeerkatMachineSupervisorPublish,
)
.await
.expect("test setup installs private trust");
let public_snapshot = CoreCommsRuntime::public_trusted_peer_projection_snapshot(&runtime)
.await
.expect("public trust snapshot");
assert!(
public_snapshot.is_empty(),
"private trust must not be exported as public reconciliation truth"
);
let remove_authority = test_projection_remove_authority(&local, &descriptor, 12);
remove_authority.install_owner(&runtime);
let err = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::RemoveTrustedPeer {
peer_id,
authority: remove_authority.authority,
},
)
.await
.expect_err("public remove authority must not remove private trust");
assert!(
matches!(err, SendError::Validation(ref message) if message.contains("public trust authority cannot remove private trusted peer")),
"unexpected private removal rejection: {err:?}"
);
let runtime_snapshot = CoreCommsRuntime::peer_ingress_runtime_snapshot(&runtime)
.await
.expect("runtime snapshot");
assert_eq!(
runtime_snapshot.trusted_peers.len(),
1,
"private trust should remain installed for admission"
);
}
fn install_test_peer_comms_handle(runtime: &CommsRuntime) {
runtime.install_peer_comms_handle(test_support::runtime_peer_comms_handle());
}
fn inproc_only_with_test_peer_authority(name: &str) -> CommsRuntime {
let runtime = CommsRuntime::inproc_only(name).unwrap();
install_test_peer_comms_handle(&runtime);
runtime
}
async fn new_with_test_peer_authority(config: ResolvedCommsConfig) -> CommsRuntime {
let runtime = CommsRuntime::new(config).await.unwrap();
install_test_peer_comms_handle(&runtime);
runtime
}
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>>,
inbound_modes:
Mutex<HashMap<meerkat_core::PeerCorrelationId, meerkat_core::types::HandlingMode>>,
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,
) -> 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 response_rejected(
&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::response_rejected",
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_send_failed(
&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_send_failed",
corr_id,
));
}
self.notify_cleanup(corr_id);
Ok(())
}
fn request_received(
&self,
corr_id: meerkat_core::PeerCorrelationId,
handling_mode: meerkat_core::types::HandlingMode,
) -> 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);
self.inbound_modes.lock().insert(corr_id, handling_mode);
Ok(())
}
fn classify_response_reply(
&self,
status: meerkat_core::ResponseStatus,
) -> Result<meerkat_core::TerminalityClass, meerkat_core::handles::DslTransitionError>
{
let generated = meerkat_runtime::handles::RuntimePeerInteractionHandle::ephemeral();
meerkat_core::handles::PeerInteractionHandle::classify_response_reply(
&generated, status,
)
}
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,
));
}
self.inbound_modes.lock().remove(&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 inbound_handling_mode(
&self,
corr_id: meerkat_core::PeerCorrelationId,
) -> Option<meerkat_core::types::HandlingMode> {
self.inbound_modes.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": "meerkat_pairing_hello",
"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": "meerkat_pairing_proof",
"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 = TrustStore::load(&runtime.config.trusted_peers_path)
.await
.unwrap();
assert!(persisted.entries().any(|entry| {
entry.name.as_str() == "hive" && entry.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": "meerkat_pairing_hello",
"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());
config.require_peer_auth = false;
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());
add_trusted_peer_with_generated_authority(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await;
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
objective_id: None,
content_taint: None,
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();
install_test_peer_comms_handle(&runtime);
let sender = Keypair::generate();
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor("sender", sender.public_key(), "tcp://127.0.0.1:4200"),
)
.await;
let request_id = Uuid::new_v4();
let reply_to = Uuid::new_v4();
let msg = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Message {
objective_id: None,
content_taint: None,
blocks: None,
body: "hello".to_string(),
handling_mode: None,
},
);
let req = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
objective_id: None,
intent: "review".to_string(),
params: serde_json::json!({"pr": 19}),
blocks: None,
content_taint: 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 {
objective_id: None,
in_reply_to: reply_to,
status: Status::Completed,
result: serde_json::json!({"ok": true}),
blocks: None,
content_taint: 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 drain_inbox_interactions_carries_sender_content_taint() {
use meerkat_core::comms::SenderContentTaint;
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("taint-drain", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
install_test_peer_comms_handle(&runtime);
let sender = Keypair::generate();
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor("sender", sender.public_key(), "tcp://127.0.0.1:4200"),
)
.await;
let tainted = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Message {
objective_id: None,
body: "tainted body".to_string(),
blocks: None,
content_taint: Some(SenderContentTaint::Tainted),
handling_mode: None,
},
);
let clean = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
objective_id: None,
intent: "review".to_string(),
params: serde_json::json!({"pr": 7}),
blocks: None,
content_taint: Some(SenderContentTaint::Clean),
handling_mode: None,
},
);
let undeclared = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Message {
objective_id: None,
body: "undeclared body".to_string(),
blocks: None,
content_taint: None,
handling_mode: None,
},
);
for envelope in [tainted, clean, undeclared] {
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope })
.into_result()
.unwrap();
}
let interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 3);
let taint_for = |predicate: &dyn Fn(&meerkat_core::InboxInteraction) -> bool| {
interactions
.iter()
.find(|interaction| predicate(interaction))
.expect("expected drained interaction")
.sender_taint
};
assert_eq!(
taint_for(&|i| matches!(
&i.content,
meerkat_core::InteractionContent::Message { body, .. } if body == "tainted body"
)),
Some(SenderContentTaint::Tainted)
);
assert_eq!(
taint_for(&|i| matches!(
&i.content,
meerkat_core::InteractionContent::Request { intent, .. } if intent == "review"
)),
Some(SenderContentTaint::Clean)
);
assert_eq!(
taint_for(&|i| matches!(
&i.content,
meerkat_core::InteractionContent::Message { body, .. } if body == "undeclared body"
)),
None,
"no declaration must stay None, never coalesced into Clean"
);
}
#[tokio::test]
async fn lifecycle_envelopes_remain_non_responseable_request_projections() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("lifecycle-no-response-cache", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
install_test_peer_comms_handle(&runtime);
let sender = Keypair::generate();
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor("sender", sender.public_key(), "tcp://127.0.0.1:4200"),
)
.await;
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!({"peer": "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_eq!(
interactions[0].handling_mode,
meerkat_core::types::HandlingMode::Queue
);
}
#[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();
install_test_peer_comms_handle(&runtime);
let sender = Keypair::generate();
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor("sender", sender.public_key(), "tcp://127.0.0.1:4200"),
)
.await;
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 {
objective_id: None,
body: "please inspect this".to_string(),
blocks: Some(blocks.clone()),
content_taint: None,
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 peer_message_body_dismiss_does_not_control_lifecycle() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("dismiss-body-is-not-lifecycle", &tmp);
let runtime = CommsRuntime::new(config).await.unwrap();
install_test_peer_comms_handle(&runtime);
let sender = Keypair::generate();
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor("sender", sender.public_key(), "tcp://127.0.0.1:4242"),
)
.await;
let msg = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Message {
objective_id: None,
body: "DISMISS".to_string(),
blocks: None,
content_taint: None,
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,
"a 'DISMISS' body must not be swallowed by a lifecycle trigger"
);
assert!(
matches!(
&interactions[0].content,
meerkat_core::InteractionContent::Message { body, .. } if body == "DISMISS"
),
"the 'DISMISS' body must drain as an ordinary peer message"
);
assert!(
!CoreCommsRuntime::dismiss_received(&runtime),
"a peer message body must never raise a runtime dismiss signal"
);
}
#[cfg(not(target_arch = "wasm32"))]
#[tokio::test]
async fn peer_request_send_failure_is_typed_send_error_not_timeout() {
let suffix = Uuid::new_v4().simple().to_string();
let peer_name = format!("unreachable-peer-{suffix}");
let runtime_name = format!("send-fail-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);
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor(&peer_name, peer.public_key(), "tcp://127.0.0.1:1"),
)
.await;
let result = CoreCommsRuntime::send(
runtime.as_ref(),
CommsCommand::PeerRequest {
objective_id: None,
content_taint: None,
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;
match result {
Err(SendError::Transport(_) | SendError::PeerOffline | SendError::PeerNotFound(_)) => {}
Ok(receipt) => {
panic!("transport send failure must not produce a success receipt, got {receipt:?}")
}
Err(other) => panic!(
"transport send failure must be a typed connectivity send error, 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();
install_test_peer_comms_handle(&runtime);
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 = inproc_only_with_test_peer_authority(&peer_name);
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(test_trust_entry(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
))
.expect("valid test peer should upsert");
}
add_trusted_peer_with_generated_authority(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await;
let receipt = CoreCommsRuntime::send(
runtime.as_ref(),
CommsCommand::PeerRequest {
objective_id: None,
content_taint: None,
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();
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await;
let result = CoreCommsRuntime::send(
&runtime,
CommsCommand::PeerRequest {
objective_id: None,
content_taint: None,
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();
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await;
let result = CoreCommsRuntime::send_and_stream(
&runtime,
CommsCommand::PeerRequest {
objective_id: None,
content_taint: None,
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();
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await;
let result = CoreCommsRuntime::send(
&runtime,
CommsCommand::PeerResponse {
objective_id: None,
content_taint: None,
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,
meerkat_core::types::HandlingMode::Queue,
)
.expect("test authority should seed inbound request state");
let result = CoreCommsRuntime::send(
runtime.as_ref(),
CommsCommand::PeerResponse {
objective_id: None,
content_taint: None,
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 terminal_peer_response_default_handling_mode_projects_from_machine_authority() {
let suffix = Uuid::new_v4().simple().to_string();
let sender_name = format!("response-default-sender-{suffix}");
let receiver_name = format!("response-default-receiver-{suffix}");
let sender = Arc::new(CommsRuntime::inproc_only(&sender_name).unwrap());
let receiver = Arc::new(CommsRuntime::inproc_only(&receiver_name).unwrap());
let peer_handle = Arc::new(TestPeerInteractionHandle::default());
sender.install_peer_request_response_authority(PeerRequestResponseAuthority::new(
peer_handle.clone(),
Arc::new(TestInteractionStreamHandle::default()),
));
install_test_peer_comms_handle(receiver.as_ref());
let receiver_peer_handle = Arc::new(TestPeerInteractionHandle::default());
receiver.install_peer_request_response_authority(PeerRequestResponseAuthority::new(
receiver_peer_handle.clone(),
Arc::new(TestInteractionStreamHandle::default()),
));
add_trusted_peer_with_generated_authority(
sender.as_ref(),
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
receiver.as_ref(),
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await;
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,
meerkat_core::types::HandlingMode::Queue,
)
.expect("machine authority should record inbound request lane");
assert_eq!(
meerkat_core::handles::PeerInteractionHandle::inbound_handling_mode(
peer_handle.as_ref(),
corr_id
),
Some(meerkat_core::types::HandlingMode::Queue)
);
meerkat_core::handles::PeerInteractionHandle::request_sent(
receiver_peer_handle.as_ref(),
corr_id,
)
.expect("receiver machine authority should expect the response");
let receipt = CoreCommsRuntime::send(
sender.as_ref(),
CommsCommand::PeerResponse {
objective_id: None,
content_taint: None,
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}),
blocks: None,
handling_mode: None,
},
)
.await
.expect("terminal response send should succeed");
assert!(matches!(receipt, SendReceipt::PeerResponseSent { .. }));
let interactions = CoreCommsRuntime::drain_inbox_interactions(receiver.as_ref()).await;
assert_eq!(interactions.len(), 1);
assert_eq!(
interactions[0].handling_mode,
meerkat_core::types::HandlingMode::Queue
);
assert!(
meerkat_core::handles::PeerInteractionHandle::inbound_handling_mode(
peer_handle.as_ref(),
corr_id
)
.is_none()
);
assert!(matches!(
&interactions[0].content,
meerkat_core::InteractionContent::Response { in_reply_to, .. }
if in_reply_to.0 == interaction_id
));
}
#[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();
add_trusted_peer_with_generated_authority(
&runtime,
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await;
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()));
add_trusted_peer_with_generated_authority(
runtime.as_ref(),
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await;
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()));
add_trusted_peer_with_generated_authority(
runtime.as_ref(),
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await;
let result = CoreCommsRuntime::send(
runtime.as_ref(),
CommsCommand::PeerRequest {
objective_id: None,
content_taint: None,
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());
add_trusted_peer_with_generated_authority(
runtime.as_ref(),
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await;
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,
meerkat_core::types::HandlingMode::Queue,
)
.expect("seed inbound request state");
let result = CoreCommsRuntime::send(
runtime.as_ref(),
CommsCommand::PeerResponse {
objective_id: None,
content_taint: None,
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);
add_trusted_peer_with_generated_authority(
runtime.as_ref(),
trusted_descriptor(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
),
)
.await;
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 = inproc_only_with_test_peer_authority(&receiver_name);
add_trusted_peer_with_generated_authority(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await;
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",
}),
},
)
.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();
install_test_peer_comms_handle(&runtime);
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 = inproc_only_with_test_peer_authority(&format!("phase1-bridge-{suffix}"));
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(inproc_only_with_test_peer_authority(&format!(
"local-transport-authority-{suffix}"
)));
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 = inproc_only_with_test_peer_authority(&format!("phase1-complete-{suffix}"));
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();
install_test_peer_comms_handle(&runtime);
let interaction_id = Uuid::new_v4();
runtime
.router
.inbox_sender()
.send_classified(InboxItem::PlainEvent {
objective_id: None,
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();
install_test_peer_comms_handle(&runtime);
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 {
objective_id: None,
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();
install_test_peer_comms_handle(&runtime);
runtime
.router
.inbox_sender()
.send_classified(InboxItem::PlainEvent {
objective_id: None,
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();
install_test_peer_comms_handle(&runtime);
let sender = Keypair::generate();
let envelope = signed_envelope(
&sender,
runtime.public_key,
crate::types::MessageKind::Request {
objective_id: None,
intent: "review".to_string(),
params: serde_json::json!({"pr": 42}),
blocks: None,
content_taint: None,
handling_mode: None,
},
);
runtime
.router
.inbox_sender()
.send_classified(InboxItem::External { envelope })
.into_result()
.unwrap();
runtime
.router
.inbox_sender()
.send_classified(InboxItem::PlainEvent {
objective_id: None,
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();
install_test_peer_comms_handle(&runtime);
let peer_key = Keypair::generate();
let trusted_peer = trusted_descriptor("ally", peer_key.public_key(), "inproc://ally");
add_trusted_peer_with_generated_authority(&runtime, trusted_peer.clone()).await;
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_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();
install_test_peer_comms_handle(&runtime);
let sender = Keypair::generate();
let envelope = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
objective_id: None,
intent: "review".to_string(),
params: serde_json::json!({"scope": "peer"}),
blocks: None,
content_taint: 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();
install_test_peer_comms_handle(&runtime);
let sender = Keypair::generate();
let pubkey_for_authority_check = sender.public_key();
let envelope = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
objective_id: None,
intent: "review".to_string(),
params: serde_json::json!({"scope": "peer"}),
blocks: None,
content_taint: 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");
let trust_dsl =
add_trusted_peer_with_generated_authority(&runtime, trusted_peer.clone()).await;
{
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 remove_authority =
test_projection_remove_authority_with_dsl(Arc::clone(&trust_dsl), &trusted_peer);
let removed = CoreCommsRuntime::apply_trust_mutation(
&runtime,
CommsTrustMutation::RemoveTrustedPeer {
peer_id: trusted_peer.peer_id.to_string(),
authority: remove_authority.authority,
},
)
.await
.expect("remove_trusted_peer should succeed");
let CommsTrustMutationResult::Removed { removed } = removed else {
panic!("unexpected generated trust mutation result: {removed:?}");
};
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();
install_test_peer_comms_handle(&runtime);
let sender = Keypair::generate();
let envelope = signed_envelope(
&sender,
runtime.public_key(),
MessageKind::Request {
objective_id: None,
intent: "status".to_string(),
params: serde_json::json!({}),
blocks: None,
content_taint: 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 = inproc_only_with_test_peer_authority(&format!("input-nores-{suffix}"));
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_no_reservation_returns_correlated_id() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = inproc_only_with_test_peer_authority(&format!("input-corr-{suffix}"));
let cmd = CommsCommand::Input {
blocks: None,
session_id: SessionId::new(),
body: "correlated input".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
source: InputSource::Rpc,
stream: InputStreamMode::None,
allow_self_session: true,
};
let returned_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 interactions = CoreCommsRuntime::drain_inbox_interactions(&runtime).await;
assert_eq!(interactions.len(), 1);
assert_eq!(
interactions[0].id, returned_id,
"returned InteractionId must match the injected event's interaction id"
);
}
#[tokio::test]
async fn test_core_send_input_reserves_stream() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime = inproc_only_with_test_peer_authority(&format!("input-res-{suffix}"));
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 = inproc_only_with_test_peer_authority(&format!("dup-attach-{suffix}"));
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 = inproc_only_with_test_peer_authority(&format!("pre-send-{suffix}"));
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 = inproc_only_with_test_peer_authority(&format!("sas-{suffix}"));
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 {
objective_id: None,
content_taint: None,
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 = inproc_only_with_test_peer_authority(&receiver_name);
add_trusted_peer_with_generated_authority(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await;
let cmd = CommsCommand::PeerMessage {
objective_id: None,
content_taint: None,
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: _,
delivery,
}) => assert_eq!(
delivery,
meerkat_core::comms::PeerDeliveryOutcome::HandedOff
),
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 = inproc_only_with_test_peer_authority(&receiver_name);
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);
add_trusted_peer_with_generated_authority(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await;
let cmd = CommsCommand::PeerMessage {
objective_id: None,
content_taint: None,
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 = inproc_only_with_test_peer_authority(&receiver_name);
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);
add_trusted_peer_with_generated_authority(
sender.as_ref(),
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await;
let cmd = CommsCommand::PeerRequest {
objective_id: None,
content_taint: None,
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 = inproc_only_with_test_peer_authority(&receiver_name);
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()),
));
add_trusted_peer_with_generated_authority(
sender.as_ref(),
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await;
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,
meerkat_core::types::HandlingMode::Queue,
)
.expect("seed inbound request state");
let cmd = CommsCommand::PeerResponse {
objective_id: None,
content_taint: None,
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 {
objective_id: None,
content_taint: None,
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_requires_generated_trust_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 = new_with_test_peer_authority(receiver_config).await;
let sender_pubkey = sender.public_key();
let receiver_pubkey = receiver.public_key();
let cmd = CommsCommand::PeerMessage {
objective_id: None,
content_taint: None,
blocks: None,
to: peer_route(&receiver_name, receiver.public_key()),
body: "hello without trusted".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
};
let expected_peer_id = receiver.public_key().to_peer_id().to_string();
let receipt = CoreCommsRuntime::send(&sender, cmd).await;
match receipt {
Err(SendError::PeerNotFound(peer_id)) if peer_id == expected_peer_id => {}
other => panic!("Expected peer-not-found without generated trust, got: {other:?}"),
}
let interactions = CoreCommsRuntime::drain_inbox_interactions(&receiver).await;
assert_eq!(interactions.len(), 0);
assert!(crate::InprocRegistry::global().unregister(&sender_pubkey));
assert!(crate::InprocRegistry::global().unregister(&receiver_pubkey));
}
#[tokio::test]
async fn test_core_peers_excludes_inproc_without_generated_trust_when_auth_disabled() {
let suffix = Uuid::new_v4().simple().to_string();
let runtime_name = format!("directory-no-auth-runtime-{suffix}");
let peer_name = format!("directory-no-auth-peer-{suffix}");
let runtime_tmp = tempfile::tempdir().unwrap();
let peer_tmp = tempfile::tempdir().unwrap();
let runtime_config = ResolvedCommsConfig {
enabled: true,
name: runtime_name.clone(),
inproc_namespace: None,
identity_dir: runtime_tmp.path().join("identity"),
trusted_peers_path: runtime_tmp.path().join("trusted_peers.json"),
comms_config: crate::CommsConfig::default(),
auth: meerkat_core::CommsAuthMode::Open,
require_peer_auth: false,
listen_uds: None,
listen_tcp: None,
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
allow_external_unauthenticated: false,
advertise_address: None,
pairing_password: None,
};
let peer_config = ResolvedCommsConfig {
enabled: true,
name: peer_name.clone(),
inproc_namespace: None,
identity_dir: peer_tmp.path().join("identity"),
trusted_peers_path: peer_tmp.path().join("trusted_peers.json"),
comms_config: crate::CommsConfig::default(),
auth: meerkat_core::CommsAuthMode::Open,
require_peer_auth: false,
listen_uds: None,
listen_tcp: None,
event_listen_tcp: None,
#[cfg(unix)]
event_listen_uds: None,
allow_external_unauthenticated: false,
advertise_address: None,
pairing_password: None,
};
let runtime = CommsRuntime::new(runtime_config).await.unwrap();
let peer = CommsRuntime::new(peer_config).await.unwrap();
let runtime_pubkey = runtime.public_key();
let peer_pubkey = peer.public_key();
let peers = CoreCommsRuntime::peers(&runtime).await;
assert!(
peers.iter().all(|entry| entry.name.as_str() != peer_name),
"inproc registry membership must not enter the peer directory without generated trust"
);
assert!(crate::InprocRegistry::global().unregister(&runtime_pubkey));
assert!(crate::InprocRegistry::global().unregister(&peer_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 = inproc_only_with_test_peer_authority(&receiver_name);
let peer_spec = trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
);
add_trusted_peer_with_generated_authority(&sender, peer_spec).await;
let reverse_spec = trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
);
add_trusted_peer_with_generated_authority(&receiver, reverse_spec).await;
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 {
objective_id: None,
content_taint: None,
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_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 =
try_add_trusted_peer_with_generated_authority(&sender, invalid_peer_spec).await;
assert!(matches!(result, Err(SendError::Validation(_))));
}
#[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:?}"),
}
}
#[tokio::test]
async fn test_runtime_startup_rejects_persisted_trusted_peer_seed() {
let tmp = tempfile::TempDir::new().unwrap();
let config = test_runtime_config("persisted-trust-seed", &tmp);
let peer_key = Keypair::generate().public_key();
tokio::fs::write(
&config.trusted_peers_path,
serde_json::json!({
"peers": [{
"name": "persisted-peer",
"pubkey": peer_key.to_pubkey_string(),
"addr": "inproc://persisted-peer"
}]
})
.to_string(),
)
.await
.unwrap();
let error = match CommsRuntime::new(config).await {
Ok(_) => panic!("runtime startup must reject persisted trust seeds"),
Err(error) => error,
};
match error {
CommsRuntimeError::TrustLoadError(message) => {
assert!(
message.contains("generated comms trust mutation authority"),
"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();
add_trusted_peer_with_generated_authority(
&sender,
trusted_descriptor(
"session-scoped-trust-peer",
peer.public_key(),
&format!("inproc://{}", peer.participant_name()),
),
)
.await;
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(test_trust_entry(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
))
.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_eq!(peer.source, PeerDirectorySource::Trusted);
assert_eq!(peer.address.to_string(), format!("inproc://{peer_name}"));
assert_eq!(peer.sendable_kinds, vec![PeerSendability::PeerMessage]);
}
#[tokio::test]
async fn test_core_peers_send_success_keeps_peer_directory_projection() {
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 = inproc_only_with_test_peer_authority(&receiver_name);
add_trusted_peer_with_generated_authority(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await;
add_trusted_peer_with_generated_authority(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
&format!("inproc://{sender_name}"),
),
)
.await;
CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
objective_id: None,
content_taint: None,
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.peer_id, receiver.public_key().to_peer_id());
}
#[tokio::test]
async fn test_core_peers_transport_failure_keeps_directory_projection() {
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();
add_trusted_peer_with_generated_authority(
&sender,
trusted_descriptor(&peer_name, peer_pubkey, "tcp://127.0.0.1:9"),
)
.await;
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
objective_id: None,
content_taint: None,
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::Transport(_))),
"transport failure should surface as a typed transport 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.peer_id, peer_pubkey.to_peer_id());
}
#[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 = inproc_only_with_test_peer_authority(&receiver_name);
add_trusted_peer_with_generated_authority(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
&format!("inproc://{receiver_name}"),
),
)
.await;
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
objective_id: None,
content_taint: None,
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;
peers
.iter()
.find(|listed| listed.name.as_str() == receiver_name)
.expect("trusted peer should remain listed after admission drop");
}
#[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 {
objective_id: None,
content_taint: None,
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_keeps_directory_projection() {
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();
add_trusted_peer_with_generated_authority(
&sender,
trusted_descriptor(
&missing_name,
missing_pubkey,
&format!("inproc://{missing_name}"),
),
)
.await;
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
objective_id: None,
content_taint: None,
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.peer_id, missing_pubkey.to_peer_id());
}
#[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(test_trust_entry(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
))
.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 {
objective_id: None,
content_taint: None,
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 {
objective_id: None,
content_taint: None,
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 {
objective_id: None,
content_taint: None,
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 = inproc_only_with_test_peer_authority(&format!("dup-close-{suffix}"));
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 = inproc_only_with_test_peer_authority(&format!("complete-{suffix}"));
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 = inproc_only_with_test_peer_authority(&peer_name);
let runtime = CommsRuntime::inproc_only(&runtime_name).unwrap();
{
let mut trusted = runtime.trusted_peers.write();
trusted
.upsert(test_trust_entry(
&peer_name,
peer.public_key(),
&format!("inproc://{peer_name}"),
))
.expect("valid test peer should upsert");
}
add_trusted_peer_with_generated_authority(
&peer,
trusted_descriptor(
&runtime_name,
runtime.public_key(),
&format!("inproc://{runtime_name}"),
),
)
.await;
let cmd = CommsCommand::PeerMessage {
objective_id: None,
content_taint: None,
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 = inproc_only_with_test_peer_authority(&format!("m5-none-{suffix}"));
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 = inproc_only_with_test_peer_authority(&format!("post-comp-{suffix}"));
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}"),
);
let trust_dsl = add_trusted_peer_with_generated_authority(&sender, peer_spec.clone()).await;
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 remove_authority =
test_projection_remove_authority_with_dsl(Arc::clone(&trust_dsl), &peer_spec);
let removed = CoreCommsRuntime::apply_trust_mutation(
&sender,
CommsTrustMutation::RemoveTrustedPeer {
peer_id: peer_spec.peer_id.to_string(),
authority: remove_authority.authority,
},
)
.await
.expect("remove_trusted_peer should succeed");
let CommsTrustMutationResult::Removed { removed } = removed else {
panic!("unexpected generated trust mutation result: {removed:?}");
};
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 {
objective_id: None,
content_taint: None,
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_fails_closed_without_generated_authority() {
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;
assert!(
matches!(result, Err(SendError::Unsupported(ref message)) if message.contains("generated comms private trust mutation authority")),
"raw private zero-pubkey add must fail closed before mutation: {result:?}"
);
let removed = CoreCommsRuntime::remove_private_trusted_peer(&runtime, &peer_id.to_string())
.await
.expect_err("raw private removal must fail closed");
assert!(
matches!(removed, SendError::Unsupported(ref message) if message.contains("generated comms private trust mutation authority")),
"raw private removal must fail closed: {removed:?}"
);
}
#[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 peer = trusted_descriptor(
&format!("rm-absent-other-{suffix}"),
other.public_key(),
&format!("inproc://rm-absent-other-{suffix}"),
);
let removed = remove_trusted_peer_with_generated_authority(&sender, &peer)
.await
.expect("remove_trusted_peer should succeed even for absent peer");
assert_eq!(peer.peer_id.to_string(), peer_id);
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 peer = trusted_descriptor(
&format!("rm-pubkey-receiver-{suffix}"),
receiver.public_key(),
&format!("inproc://rm-pubkey-receiver-{suffix}"),
);
let remove_authority =
test_projection_remove_authority(&local_descriptor_for_runtime(&sender), &peer, 0);
remove_authority.install_owner(&sender);
let result = CoreCommsRuntime::apply_trust_mutation(
&sender,
CommsTrustMutation::RemoveTrustedPeer {
authority: remove_authority.authority,
peer_id: peer_pubkey_string,
},
)
.await;
assert!(
matches!(result, Err(SendError::Validation(_))),
"remove_trusted_peer should validate its argument as PeerId, got: {result:?}"
);
}
}