use std::collections::{BTreeMap, BTreeSet};
use std::time::Duration;
use base64::Engine as _;
use serde::{Deserialize, Serialize};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::{TcpListener, TcpStream};
#[cfg(unix)]
use tokio::net::{UnixListener, UnixStream};
use super::cross_mob_remote::{RemoteEndpoint, RemoteMobError};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ControlListenAddr {
Tcp(String),
Uds(String),
}
impl ControlListenAddr {
pub fn parse(s: &str) -> Result<Self, String> {
match crate::contact_directory::parse_transport(s) {
Some(crate::contact_directory::MobTransport::Tcp(addr)) => Ok(Self::Tcp(addr)),
Some(crate::contact_directory::MobTransport::Uds(path)) => Ok(Self::Uds(path)),
Some(crate::contact_directory::MobTransport::Inproc) => Err(
"control listener requires a cross-process address (tcp://host:port or \
uds:///path); 'inproc' has no listener"
.to_string(),
),
None => Err(format!(
"invalid control listen address '{s}': expected tcp://host:port or uds:///path"
)),
}
}
}
impl std::fmt::Display for ControlListenAddr {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Tcp(addr) => write!(f, "tcp://{addr}"),
Self::Uds(path) => write!(f, "uds://{path}"),
}
}
}
pub enum BoundControlListener {
Tcp {
listener: TcpListener,
advertised: String,
},
#[cfg(unix)]
Uds {
listener: UnixListener,
advertised: String,
},
}
impl BoundControlListener {
pub async fn bind(addr: &ControlListenAddr) -> Result<Self, std::io::Error> {
match addr {
ControlListenAddr::Tcp(spec) => {
let listener = TcpListener::bind(spec).await?;
let local = listener.local_addr()?;
Ok(Self::Tcp {
listener,
advertised: format!("tcp://{local}"),
})
}
#[cfg(unix)]
ControlListenAddr::Uds(path) => {
let listener = UnixListener::bind(std::path::Path::new(path))?;
Ok(Self::Uds {
listener,
advertised: format!("uds://{path}"),
})
}
#[cfg(not(unix))]
ControlListenAddr::Uds(_) => Err(std::io::Error::new(
std::io::ErrorKind::Unsupported,
"unix domain sockets are not supported on this platform",
)),
}
}
pub fn advertised_address(&self) -> &str {
match self {
Self::Tcp { advertised, .. } => advertised,
#[cfg(unix)]
Self::Uds { advertised, .. } => advertised,
}
}
pub async fn serve(
self,
handler: std::sync::Arc<dyn ControlHandler>,
signer: ControlSignerSlot,
) {
self.serve_with_authorizer(
handler,
signer,
std::sync::Arc::new(ControlAuthorizer::open()),
)
.await;
}
pub async fn serve_with_authorizer(
self,
handler: std::sync::Arc<dyn ControlHandler>,
signer: ControlSignerSlot,
authorizer: std::sync::Arc<ControlAuthorizer>,
) {
match self {
Self::Tcp { listener, .. } => {
serve_tcp_control_with_authorizer(listener, handler, signer, authorizer).await;
}
#[cfg(unix)]
Self::Uds { listener, .. } => {
serve_uds_control_with_authorizer(listener, handler, signer, authorizer).await;
}
}
}
}
const MAX_CONTROL_PAYLOAD: u32 = 64 * 1024;
pub const DEFAULT_CONTROL_TIMEOUT: Duration = Duration::from_secs(5);
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(tag = "op", rename_all = "snake_case")]
pub enum ControlRequest {
Wire {
remote_member: String,
local_peer_spec_address: String,
local_comms_name: String,
local_peer_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
local_pubkey_b64: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
nonce: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
caller: Option<ControlCaller>,
},
Unwire {
remote_member: String,
local_peer_spec_address: String,
local_comms_name: String,
local_peer_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
local_pubkey_b64: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
nonce: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
caller: Option<ControlCaller>,
},
Inject {
remote_member: String,
content: serde_json::Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
nonce: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
caller: Option<ControlCaller>,
},
LookupMember {
remote_member: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
nonce: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
caller: Option<ControlCaller>,
},
HostDescribe {
#[serde(default, skip_serializing_if = "Option::is_none")]
nonce: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
caller: Option<ControlCaller>,
},
HostHealth {
#[serde(default, skip_serializing_if = "Option::is_none")]
nonce: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
caller: Option<ControlCaller>,
},
}
impl ControlRequest {
pub fn verb(&self) -> ControlVerb {
match self {
Self::Wire { .. } => ControlVerb::Wire,
Self::Unwire { .. } => ControlVerb::Unwire,
Self::Inject { .. } => ControlVerb::Inject,
Self::LookupMember { .. } => ControlVerb::LookupMember,
Self::HostDescribe { .. } => ControlVerb::HostDescribe,
Self::HostHealth { .. } => ControlVerb::HostHealth,
}
}
pub fn remote_member(&self) -> &str {
match self {
Self::Wire { remote_member, .. }
| Self::Unwire { remote_member, .. }
| Self::Inject { remote_member, .. }
| Self::LookupMember { remote_member, .. } => remote_member,
Self::HostDescribe { .. } | Self::HostHealth { .. } => HOST_PLANE_MEMBER,
}
}
pub fn caller(&self) -> Option<&ControlCaller> {
match self {
Self::Wire { caller, .. }
| Self::Unwire { caller, .. }
| Self::Inject { caller, .. }
| Self::LookupMember { caller, .. }
| Self::HostDescribe { caller, .. }
| Self::HostHealth { caller, .. } => caller.as_ref(),
}
}
fn caller_mut(&mut self) -> Option<&mut ControlCaller> {
match self {
Self::Wire { caller, .. }
| Self::Unwire { caller, .. }
| Self::Inject { caller, .. }
| Self::LookupMember { caller, .. }
| Self::HostDescribe { caller, .. }
| Self::HostHealth { caller, .. } => caller.as_mut(),
}
}
pub fn set_caller(&mut self, value: Option<ControlCaller>) {
match self {
Self::Wire { caller, .. }
| Self::Unwire { caller, .. }
| Self::Inject { caller, .. }
| Self::LookupMember { caller, .. }
| Self::HostDescribe { caller, .. }
| Self::HostHealth { caller, .. } => *caller = value,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(tag = "result", rename_all = "snake_case")]
pub enum ControlResponse {
Ok {
#[serde(default, skip_serializing_if = "Option::is_none")]
sig_b64: Option<String>,
},
Injected {
session_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
sig_b64: Option<String>,
},
Member {
peer_id: String,
comms_name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pubkey_b64: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
advertised_address: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
sig_b64: Option<String>,
},
Host {
facts: super::remote_host::HostFacts,
#[serde(default, skip_serializing_if = "Option::is_none")]
sig_b64: Option<String>,
},
HostHealth {
health: meerkat_contracts::RuntimeHostHealth,
#[serde(default, skip_serializing_if = "Option::is_none")]
sig_b64: Option<String>,
},
Err { code: String, message: String },
}
pub const CONTROL_SIG_CONTEXT: &str = "mobkit-cross-mob-control-v1";
pub type ControlSignerSlot = std::sync::Arc<
std::sync::RwLock<Option<std::sync::Arc<crate::auth::peer_keys::GatewayPeerKeys>>>,
>;
pub type HostFactsProviderSlot = std::sync::Arc<
std::sync::RwLock<Option<std::sync::Arc<dyn super::remote_host::HostFactsProvider>>>,
>;
pub fn unsigned_control_signer() -> ControlSignerSlot {
std::sync::Arc::new(std::sync::RwLock::new(None))
}
pub fn empty_host_facts_provider() -> HostFactsProviderSlot {
std::sync::Arc::new(std::sync::RwLock::new(None))
}
fn control_request_digest_hex(request_bytes: &[u8]) -> String {
sha256_hex(request_bytes)
}
pub(crate) fn sha256_hex(bytes: &[u8]) -> String {
use sha2::Digest;
let mut hasher = sha2::Sha256::new();
hasher.update(bytes);
let digest = hasher.finalize();
let mut hex = String::with_capacity(digest.len() * 2);
for byte in digest {
use std::fmt::Write;
let _ = write!(hex, "{byte:02x}");
}
hex
}
pub fn control_response_signing_payload(
request_bytes: &[u8],
response: &ControlResponse,
) -> Option<String> {
let facts = match response {
ControlResponse::Ok { .. } => "ok".to_string(),
ControlResponse::Injected { session_id, .. } => format!("injected\n{session_id}"),
ControlResponse::Member {
peer_id,
comms_name,
pubkey_b64,
advertised_address,
..
} => format!(
"member\n{peer_id}\n{comms_name}\n{}\n{}",
pubkey_b64.as_deref().unwrap_or(""),
advertised_address.as_deref().unwrap_or("")
),
ControlResponse::Host { facts, .. } => {
format!("host\n{}", facts.signing_digest_hex())
}
ControlResponse::HostHealth { health, .. } => format!(
"host_health\n{}",
super::remote_host::host_health_digest_hex(health)
),
ControlResponse::Err { .. } => return None,
};
Some(format!(
"{CONTROL_SIG_CONTEXT}\n{}\n{facts}",
control_request_digest_hex(request_bytes)
))
}
fn sign_control_response(
signer: &crate::auth::peer_keys::GatewayPeerKeys,
request_bytes: &[u8],
response: &mut ControlResponse,
) {
use ed25519_dalek::Signer;
let Some(payload) = control_response_signing_payload(request_bytes, response) else {
return;
};
let signature = signer.signing_key().sign(payload.as_bytes());
let encoded = base64::engine::general_purpose::STANDARD.encode(signature.to_bytes());
match response {
ControlResponse::Ok { sig_b64 }
| ControlResponse::Injected { sig_b64, .. }
| ControlResponse::Member { sig_b64, .. }
| ControlResponse::Host { sig_b64, .. }
| ControlResponse::HostHealth { sig_b64, .. } => *sig_b64 = Some(encoded),
ControlResponse::Err { .. } => {}
}
}
pub fn verify_control_response(
pinned_pubkey: &[u8; 32],
request_bytes: &[u8],
response: &ControlResponse,
) -> Result<(), String> {
let Some(payload) = control_response_signing_payload(request_bytes, response) else {
return Ok(());
};
let sig_b64 = match response {
ControlResponse::Ok { sig_b64 }
| ControlResponse::Injected { sig_b64, .. }
| ControlResponse::Member { sig_b64, .. }
| ControlResponse::Host { sig_b64, .. }
| ControlResponse::HostHealth { sig_b64, .. } => sig_b64.as_deref(),
ControlResponse::Err { .. } => None,
};
let Some(sig_b64) = sig_b64 else {
return Err(
"response is unsigned; the peer gateway has no signing keys installed or predates \
signed control responses"
.to_string(),
);
};
let sig_bytes = base64::engine::general_purpose::STANDARD
.decode(sig_b64)
.map_err(|err| format!("signature is not valid base64: {err}"))?;
let sig_bytes: [u8; 64] = sig_bytes
.try_into()
.map_err(|_| "signature must be 64 bytes".to_string())?;
let signature = ed25519_dalek::Signature::from_bytes(&sig_bytes);
let verifying_key = ed25519_dalek::VerifyingKey::from_bytes(pinned_pubkey)
.map_err(|err| format!("pinned pubkey is not a valid Ed25519 key: {err}"))?;
verifying_key
.verify_strict(payload.as_bytes(), &signature)
.map_err(|_| "signature does not verify against the pinned gateway pubkey".to_string())
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct ControlCaller {
pub pubkey_b64: String,
pub sig_b64: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub audience: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ControlVerb {
Wire,
Unwire,
Inject,
LookupMember,
HostDescribe,
HostHealth,
}
impl ControlVerb {
pub fn as_str(&self) -> &'static str {
match self {
Self::Wire => "wire",
Self::Unwire => "unwire",
Self::Inject => "inject",
Self::LookupMember => "lookup_member",
Self::HostDescribe => "host_describe",
Self::HostHealth => "host_health",
}
}
pub fn parse(value: &str) -> Option<Self> {
match value {
"wire" => Some(Self::Wire),
"unwire" => Some(Self::Unwire),
"inject" => Some(Self::Inject),
"lookup_member" => Some(Self::LookupMember),
"host_describe" => Some(Self::HostDescribe),
"host_health" => Some(Self::HostHealth),
_ => None,
}
}
pub fn all() -> [Self; 6] {
[
Self::Wire,
Self::Unwire,
Self::Inject,
Self::LookupMember,
Self::HostDescribe,
Self::HostHealth,
]
}
pub fn member_plane() -> [Self; 4] {
[Self::Wire, Self::Unwire, Self::Inject, Self::LookupMember]
}
pub fn is_host_plane(&self) -> bool {
match self {
Self::HostDescribe | Self::HostHealth => true,
Self::Wire | Self::Unwire | Self::Inject | Self::LookupMember => false,
}
}
}
pub const HOST_PLANE_MEMBER: &str = "";
impl std::fmt::Display for ControlVerb {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
pub const CONTROL_CALLER_SIG_CONTEXT: &str = "mobkit-cross-mob-control-caller-v1";
pub(crate) fn push_signed_field(out: &mut String, value: &str) {
use std::fmt::Write;
let _ = writeln!(out, "{}:{value}", value.len());
}
fn control_content_digest_hex(content: &serde_json::Value) -> String {
control_request_digest_hex(&serde_json::to_vec(content).unwrap_or_default())
}
pub fn control_request_signing_payload(request: &ControlRequest) -> String {
let mut payload = String::new();
payload.push_str(CONTROL_CALLER_SIG_CONTEXT);
payload.push('\n');
push_signed_field(&mut payload, request.verb().as_str());
match request {
ControlRequest::Wire {
remote_member,
local_peer_spec_address,
local_comms_name,
local_peer_id,
local_pubkey_b64,
nonce,
caller: _caller,
}
| ControlRequest::Unwire {
remote_member,
local_peer_spec_address,
local_comms_name,
local_peer_id,
local_pubkey_b64,
nonce,
caller: _caller,
} => {
push_signed_field(&mut payload, remote_member);
push_signed_field(&mut payload, local_peer_spec_address);
push_signed_field(&mut payload, local_comms_name);
push_signed_field(&mut payload, local_peer_id);
push_signed_field(&mut payload, local_pubkey_b64.as_deref().unwrap_or(""));
push_signed_field(&mut payload, nonce.as_deref().unwrap_or(""));
}
ControlRequest::Inject {
remote_member,
content,
nonce,
caller: _caller,
} => {
push_signed_field(&mut payload, remote_member);
push_signed_field(&mut payload, &control_content_digest_hex(content));
push_signed_field(&mut payload, nonce.as_deref().unwrap_or(""));
}
ControlRequest::LookupMember {
remote_member,
nonce,
caller: _caller,
} => {
push_signed_field(&mut payload, remote_member);
push_signed_field(&mut payload, nonce.as_deref().unwrap_or(""));
}
ControlRequest::HostDescribe {
nonce,
caller: _caller,
}
| ControlRequest::HostHealth {
nonce,
caller: _caller,
} => {
push_signed_field(&mut payload, nonce.as_deref().unwrap_or(""));
}
}
push_signed_field(
&mut payload,
request
.caller()
.and_then(|caller| caller.audience.as_deref())
.unwrap_or(""),
);
payload
}
pub fn sign_control_request_as_caller(
keys: &crate::auth::peer_keys::GatewayPeerKeys,
audience: Option<&str>,
request: &mut ControlRequest,
) {
use ed25519_dalek::Signer;
request.set_caller(Some(ControlCaller {
pubkey_b64: keys.pubkey_b64(),
sig_b64: String::new(),
audience: audience.map(str::to_string),
}));
let payload = control_request_signing_payload(request);
let signature = keys.signing_key().sign(payload.as_bytes());
let encoded = base64::engine::general_purpose::STANDARD.encode(signature.to_bytes());
if let Some(caller) = request.caller_mut() {
caller.sig_b64 = encoded;
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ControlAuthzDenial {
UnauthenticatedCaller { verb: ControlVerb },
InvalidCallerCredential { reason: String },
InvalidCallerSignature { pubkey_b64: String },
CallerNotGranted { pubkey_b64: String },
VerbNotGranted { label: String, verb: ControlVerb },
MemberNotGranted {
label: String,
verb: ControlVerb,
member: String,
},
AudienceMismatch {
label: String,
expected: String,
presented: String,
},
}
pub const CONTROL_AUTHZ_DENIAL_CODES: [&str; 7] = [
"unauthenticated_caller",
"invalid_caller_credential",
"invalid_caller_signature",
"caller_not_granted",
"verb_not_granted",
"member_not_granted",
"audience_mismatch",
];
impl ControlAuthzDenial {
pub fn code(&self) -> &'static str {
match self {
Self::UnauthenticatedCaller { .. } => "unauthenticated_caller",
Self::InvalidCallerCredential { .. } => "invalid_caller_credential",
Self::InvalidCallerSignature { .. } => "invalid_caller_signature",
Self::CallerNotGranted { .. } => "caller_not_granted",
Self::VerbNotGranted { .. } => "verb_not_granted",
Self::MemberNotGranted { .. } => "member_not_granted",
Self::AudienceMismatch { .. } => "audience_mismatch",
}
}
pub fn is_denial_code(code: &str) -> bool {
CONTROL_AUTHZ_DENIAL_CODES.contains(&code)
}
}
impl std::fmt::Display for ControlAuthzDenial {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::UnauthenticatedCaller { verb } => write!(
f,
"control verb '{verb}' requires a signed caller credential on this listener; \
the calling gateway must install a keypair (set_gateway_peer_keys) and be \
listed in this gateway's control grant table"
),
Self::InvalidCallerCredential { reason } => {
write!(f, "caller credential is unusable: {reason}")
}
Self::InvalidCallerSignature { pubkey_b64 } => write!(
f,
"caller signature does not verify against the presented pubkey '{pubkey_b64}'"
),
Self::CallerNotGranted { pubkey_b64 } => write!(
f,
"caller '{pubkey_b64}' holds no control grant on this gateway"
),
Self::VerbNotGranted { label, verb } if verb.is_host_plane() => write!(
f,
"caller '{label}' is not granted host-plane control verb '{verb}' on this \
gateway; host-plane verbs are never covered by verbs = [\"*\"] and must be \
named explicitly"
),
Self::VerbNotGranted { label, verb } => write!(
f,
"caller '{label}' is not granted control verb '{verb}' on this gateway"
),
Self::MemberNotGranted {
label,
verb,
member,
} => write!(
f,
"caller '{label}' is not granted control verb '{verb}' on member '{member}'"
),
Self::AudienceMismatch {
label,
expected,
presented,
} => write!(
f,
"caller '{label}' presented a request minted for '{presented}', but this \
gateway answers for '{expected}'"
),
}
}
}
impl std::error::Error for ControlAuthzDenial {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ControlMemberScope {
All,
Members(BTreeSet<String>),
}
impl ControlMemberScope {
pub fn members<I, S>(aliases: I) -> Self
where
I: IntoIterator<Item = S>,
S: AsRef<str>,
{
let mut set = BTreeSet::new();
for alias in aliases {
let alias = alias.as_ref();
if alias == "*" {
return Self::All;
}
set.insert(normalize_member_alias(alias));
}
Self::Members(set)
}
pub fn contains(&self, member: &str) -> bool {
match self {
Self::All => true,
Self::Members(members) => members.contains(&normalize_member_alias(member)),
}
}
}
fn normalize_member_alias(member: &str) -> String {
crate::member_comms_id::runtime_alias_str(member).into_owned()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ControlGrant {
label: String,
verbs: BTreeSet<ControlVerb>,
members: ControlMemberScope,
}
impl ControlGrant {
pub fn new<I>(label: impl Into<String>, verbs: I, members: ControlMemberScope) -> Self
where
I: IntoIterator<Item = ControlVerb>,
{
Self {
label: label.into(),
verbs: verbs.into_iter().collect(),
members,
}
}
pub fn label(&self) -> &str {
&self.label
}
pub fn verbs(&self) -> &BTreeSet<ControlVerb> {
&self.verbs
}
pub fn members(&self) -> &ControlMemberScope {
&self.members
}
fn permits(&self, verb: ControlVerb, member: &str) -> Result<(), ControlAuthzDenial> {
if !self.verbs.contains(&verb) {
return Err(ControlAuthzDenial::VerbNotGranted {
label: self.label.clone(),
verb,
});
}
if verb.is_host_plane() {
return Ok(());
}
if !self.members.contains(member) {
return Err(ControlAuthzDenial::MemberNotGranted {
label: self.label.clone(),
verb,
member: member.to_string(),
});
}
Ok(())
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ControlGrantTable {
grants: BTreeMap<[u8; 32], ControlGrant>,
}
impl ControlGrantTable {
pub fn new() -> Self {
Self::default()
}
pub fn insert(&mut self, caller_pubkey: [u8; 32], grant: ControlGrant) -> Option<ControlGrant> {
self.grants.insert(caller_pubkey, grant)
}
pub fn get(&self, caller_pubkey: &[u8; 32]) -> Option<&ControlGrant> {
self.grants.get(caller_pubkey)
}
pub fn len(&self) -> usize {
self.grants.len()
}
pub fn is_empty(&self) -> bool {
self.grants.is_empty()
}
pub fn from_toml(text: &str) -> Result<Option<Self>, ControlGrantConfigError> {
let table: toml::Value =
toml::from_str(text).map_err(|err| ControlGrantConfigError::Parse(err.to_string()))?;
let Some(section) = table.get("control_grants") else {
return Ok(None);
};
let section = section.as_table().ok_or_else(|| {
ControlGrantConfigError::Parse(
"control_grants must be a table of caller labels, e.g. [control_grants.ops-mob]"
.to_string(),
)
})?;
let mut out = Self::new();
let mut labels_by_pubkey: BTreeMap<[u8; 32], String> = BTreeMap::new();
for (label, value) in section {
let entry = value
.as_table()
.ok_or_else(|| ControlGrantConfigError::InvalidField {
label: label.clone(),
detail: "each caller must be a table with pubkey/verbs/members".to_string(),
})?;
for (key, _) in entry {
if !matches!(key.as_str(), "pubkey" | "verbs" | "members") {
return Err(ControlGrantConfigError::InvalidField {
label: label.clone(),
detail: format!(
"unknown key '{key}'; a control grant accepts only pubkey, verbs, \
and members (did you mean 'members'?)"
),
});
}
}
let pubkey_text = entry
.get("pubkey")
.and_then(|value| value.as_str())
.ok_or_else(|| ControlGrantConfigError::InvalidField {
label: label.clone(),
detail: "missing 'pubkey' (base64 Ed25519, optional ed25519: prefix)"
.to_string(),
})?;
let pubkey = crate::auth::peer_keys::decode_pubkey_b64(pubkey_text).map_err(|err| {
ControlGrantConfigError::InvalidPubkey {
label: label.clone(),
reason: err.to_string(),
}
})?;
if let Some(other) = labels_by_pubkey.get(&pubkey) {
return Err(ControlGrantConfigError::DuplicatePubkey {
label: label.clone(),
other: other.clone(),
});
}
let verb_values = entry
.get("verbs")
.and_then(|value| value.as_array())
.ok_or_else(|| ControlGrantConfigError::InvalidField {
label: label.clone(),
detail: "missing 'verbs' array; a caller with no verbs reaches nothing"
.to_string(),
})?;
let mut verbs = BTreeSet::new();
for verb_value in verb_values {
let verb_text =
verb_value
.as_str()
.ok_or_else(|| ControlGrantConfigError::InvalidField {
label: label.clone(),
detail: "'verbs' entries must be strings".to_string(),
})?;
if verb_text == "*" {
verbs.extend(ControlVerb::member_plane());
continue;
}
let verb = ControlVerb::parse(verb_text).ok_or_else(|| {
ControlGrantConfigError::UnknownVerb {
label: label.clone(),
verb: verb_text.to_string(),
}
})?;
verbs.insert(verb);
}
if verbs.is_empty() {
return Err(ControlGrantConfigError::InvalidField {
label: label.clone(),
detail: "'verbs' is empty; remove the caller instead of granting nothing"
.to_string(),
});
}
let members = match entry.get("members") {
None => ControlMemberScope::All,
Some(value) => {
let items =
value
.as_array()
.ok_or_else(|| ControlGrantConfigError::InvalidField {
label: label.clone(),
detail: "'members' must be an array of member aliases".to_string(),
})?;
let mut aliases = Vec::with_capacity(items.len());
for item in items {
let alias =
item.as_str()
.ok_or_else(|| ControlGrantConfigError::InvalidField {
label: label.clone(),
detail: "'members' entries must be strings".to_string(),
})?;
aliases.push(alias.to_string());
}
ControlMemberScope::members(aliases)
}
};
labels_by_pubkey.insert(pubkey, label.clone());
out.insert(pubkey, ControlGrant::new(label.clone(), verbs, members));
}
Ok(Some(out))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ControlGrantConfigError {
Parse(String),
InvalidPubkey { label: String, reason: String },
DuplicatePubkey { label: String, other: String },
UnknownVerb { label: String, verb: String },
InvalidField { label: String, detail: String },
}
impl std::fmt::Display for ControlGrantConfigError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Parse(reason) => write!(f, "control grant parse error: {reason}"),
Self::InvalidPubkey { label, reason } => {
write!(f, "control grant '{label}' has an invalid pubkey: {reason}")
}
Self::DuplicatePubkey { label, other } => write!(
f,
"control grants '{label}' and '{other}' name the same caller pubkey; one caller \
must have exactly one grant"
),
Self::UnknownVerb { label, verb } => write!(
f,
"control grant '{label}' names unknown verb '{verb}'; expected one of \
wire, unwire, inject, lookup_member (or '*')"
),
Self::InvalidField { label, detail } => {
write!(f, "control grant '{label}': {detail}")
}
}
}
}
impl std::error::Error for ControlGrantConfigError {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AuthorizedCaller {
pub label: String,
pub pubkey: [u8; 32],
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ControlAuthorizer {
Open,
Grants {
table: ControlGrantTable,
expected_audience: Option<String>,
},
}
impl ControlAuthorizer {
pub fn open() -> Self {
Self::Open
}
pub fn with_grants(table: ControlGrantTable) -> Self {
Self::Grants {
table,
expected_audience: None,
}
}
pub fn with_grants_for_audience(table: ControlGrantTable, audience: impl Into<String>) -> Self {
Self::Grants {
table,
expected_audience: Some(audience.into()),
}
}
pub fn from_toml(text: &str) -> Result<Self, ControlGrantConfigError> {
Ok(match ControlGrantTable::from_toml(text)? {
Some(table) => Self::with_grants(table),
None => Self::Open,
})
}
pub fn is_enforcing(&self) -> bool {
matches!(self, Self::Grants { .. })
}
pub fn authorize(
&self,
request: &ControlRequest,
) -> Result<Option<AuthorizedCaller>, ControlAuthzDenial> {
let Self::Grants {
table,
expected_audience,
} = self
else {
return Ok(None);
};
let verb = request.verb();
let Some(caller) = request.caller() else {
return Err(ControlAuthzDenial::UnauthenticatedCaller { verb });
};
let pubkey =
crate::auth::peer_keys::decode_pubkey_b64(&caller.pubkey_b64).map_err(|err| {
ControlAuthzDenial::InvalidCallerCredential {
reason: format!("caller pubkey: {err}"),
}
})?;
verify_caller_signature(&pubkey, caller, request)?;
let Some(grant) = table.get(&pubkey) else {
return Err(ControlAuthzDenial::CallerNotGranted {
pubkey_b64: caller.pubkey_b64.clone(),
});
};
if let Some(expected) = expected_audience.as_deref() {
let presented = caller.audience.as_deref().unwrap_or("");
if presented != expected {
return Err(ControlAuthzDenial::AudienceMismatch {
label: grant.label().to_string(),
expected: expected.to_string(),
presented: presented.to_string(),
});
}
}
grant.permits(verb, request.remote_member())?;
Ok(Some(AuthorizedCaller {
label: grant.label().to_string(),
pubkey,
}))
}
}
fn verify_caller_signature(
pubkey: &[u8; 32],
caller: &ControlCaller,
request: &ControlRequest,
) -> Result<(), ControlAuthzDenial> {
let sig_bytes = base64::engine::general_purpose::STANDARD
.decode(&caller.sig_b64)
.map_err(|err| ControlAuthzDenial::InvalidCallerCredential {
reason: format!("caller signature is not valid base64: {err}"),
})?;
let sig_bytes: [u8; 64] =
sig_bytes
.try_into()
.map_err(|_| ControlAuthzDenial::InvalidCallerCredential {
reason: "caller signature must be 64 bytes".to_string(),
})?;
let signature = ed25519_dalek::Signature::from_bytes(&sig_bytes);
let verifying_key = ed25519_dalek::VerifyingKey::from_bytes(pubkey).map_err(|err| {
ControlAuthzDenial::InvalidCallerCredential {
reason: format!("caller pubkey is not a valid Ed25519 key: {err}"),
}
})?;
let payload = control_request_signing_payload(request);
verifying_key
.verify_strict(payload.as_bytes(), &signature)
.map_err(|_| ControlAuthzDenial::InvalidCallerSignature {
pubkey_b64: caller.pubkey_b64.clone(),
})
}
enum ControlStream {
Tcp(TcpStream),
#[cfg(unix)]
Uds(UnixStream),
}
impl ControlStream {
async fn write_frame(&mut self, payload: &[u8]) -> Result<(), std::io::Error> {
let len = u32::try_from(payload.len()).map_err(|_| {
std::io::Error::new(std::io::ErrorKind::InvalidInput, "payload too large")
})?;
let header = len.to_be_bytes();
match self {
Self::Tcp(s) => {
s.write_all(&header).await?;
s.write_all(payload).await?;
s.flush().await
}
#[cfg(unix)]
Self::Uds(s) => {
s.write_all(&header).await?;
s.write_all(payload).await?;
s.flush().await
}
}
}
async fn read_frame(&mut self) -> Result<Vec<u8>, std::io::Error> {
let mut header = [0u8; 4];
match self {
Self::Tcp(s) => s.read_exact(&mut header).await?,
#[cfg(unix)]
Self::Uds(s) => s.read_exact(&mut header).await?,
};
let len = u32::from_be_bytes(header);
if len > MAX_CONTROL_PAYLOAD {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("frame too large: {len} bytes"),
));
}
let mut buf = vec![0u8; len as usize];
match self {
Self::Tcp(s) => s.read_exact(&mut buf).await?,
#[cfg(unix)]
Self::Uds(s) => s.read_exact(&mut buf).await?,
};
Ok(buf)
}
}
pub struct RemoteControlClient;
impl RemoteControlClient {
pub async fn send(
endpoint: &RemoteEndpoint,
request: &ControlRequest,
timeout: Duration,
) -> Result<ControlResponse, RemoteMobError> {
let payload =
serde_json::to_vec(request).map_err(|err| encode_error(endpoint, err.to_string()))?;
Self::send_payload(endpoint, &payload, timeout).await
}
pub async fn send_payload(
endpoint: &RemoteEndpoint,
payload: &[u8],
timeout: Duration,
) -> Result<ControlResponse, RemoteMobError> {
tokio::time::timeout(timeout, Self::send_inner(endpoint, payload))
.await
.map_err(|_| RemoteMobError::ControlChannelUnavailable {
mob_id: String::new(),
endpoint: endpoint.comms_address(),
operation: "timeout",
})?
}
async fn send_inner(
endpoint: &RemoteEndpoint,
payload: &[u8],
) -> Result<ControlResponse, RemoteMobError> {
let mut stream = match endpoint {
RemoteEndpoint::Tcp(addr) => ControlStream::Tcp(
TcpStream::connect(addr)
.await
.map_err(|err| io_error("connect", endpoint, err))?,
),
#[cfg(unix)]
RemoteEndpoint::Uds(path) => ControlStream::Uds(
UnixStream::connect(std::path::Path::new(path))
.await
.map_err(|err| io_error("connect", endpoint, err))?,
),
#[cfg(not(unix))]
RemoteEndpoint::Uds(_) => {
return Err(RemoteMobError::UnsupportedTransport {
mob_id: String::new(),
transport: endpoint.comms_address(),
});
}
};
stream
.write_frame(payload)
.await
.map_err(|err| io_error("write", endpoint, err))?;
let response_payload = stream
.read_frame()
.await
.map_err(|err| io_error("read", endpoint, err))?;
serde_json::from_slice::<ControlResponse>(&response_payload)
.map_err(|err| decode_error(endpoint, err.to_string()))
}
}
fn io_error(stage: &'static str, endpoint: &RemoteEndpoint, err: std::io::Error) -> RemoteMobError {
RemoteMobError::ControlChannelUnavailable {
mob_id: String::new(),
endpoint: endpoint.comms_address(),
operation: match stage {
"connect" => "connect",
"write" => "write",
"read" => "read",
_ => "io",
},
}
.with_context(err.to_string())
}
fn encode_error(endpoint: &RemoteEndpoint, message: String) -> RemoteMobError {
RemoteMobError::Encode {
endpoint: endpoint.comms_address(),
message,
}
}
fn decode_error(endpoint: &RemoteEndpoint, message: String) -> RemoteMobError {
RemoteMobError::Decode {
endpoint: endpoint.comms_address(),
message,
}
}
pub trait ControlHandler: Send + Sync + 'static {
fn handle(
&self,
request: ControlRequest,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ControlResponse> + Send + '_>>;
}
pub struct MobHandleControlHandler {
handle: meerkat_mob::MobHandle,
identity_authority: ControlIdentityAuthority,
session_service: Option<std::sync::Arc<dyn meerkat_mob::MobSessionService>>,
host_facts: HostFactsProviderSlot,
}
enum ControlIdentityAuthority {
None,
Fixed(std::sync::Arc<crate::identity_first::IdentityRuntime>),
Shared(
std::sync::Arc<
std::sync::RwLock<Option<std::sync::Arc<crate::identity_first::IdentityRuntime>>>,
>,
),
}
impl ControlIdentityAuthority {
fn current(&self) -> Option<std::sync::Arc<crate::identity_first::IdentityRuntime>> {
match self {
Self::None => None,
Self::Fixed(runtime) => Some(std::sync::Arc::clone(runtime)),
Self::Shared(slot) => slot
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone(),
}
}
}
impl MobHandleControlHandler {
pub fn new(handle: meerkat_mob::MobHandle) -> Self {
Self {
handle,
identity_authority: ControlIdentityAuthority::None,
session_service: None,
host_facts: empty_host_facts_provider(),
}
}
pub fn with_identity_runtime(
handle: meerkat_mob::MobHandle,
identity_runtime: std::sync::Arc<crate::identity_first::IdentityRuntime>,
) -> Self {
Self {
handle,
identity_authority: ControlIdentityAuthority::Fixed(identity_runtime),
session_service: None,
host_facts: empty_host_facts_provider(),
}
}
pub fn with_shared_identity_authority(
handle: meerkat_mob::MobHandle,
identity_slot: std::sync::Arc<
std::sync::RwLock<Option<std::sync::Arc<crate::identity_first::IdentityRuntime>>>,
>,
) -> Self {
Self {
handle,
identity_authority: ControlIdentityAuthority::Shared(identity_slot),
session_service: None,
host_facts: empty_host_facts_provider(),
}
}
#[must_use]
pub fn with_session_service(
mut self,
session_service: std::sync::Arc<dyn meerkat_mob::MobSessionService>,
) -> Self {
self.session_service = Some(session_service);
self
}
#[must_use]
pub fn with_host_facts(
self,
provider: std::sync::Arc<dyn super::remote_host::HostFactsProvider>,
) -> Self {
*self
.host_facts
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(provider);
self
}
#[must_use]
pub fn with_host_facts_slot(mut self, slot: HostFactsProviderSlot) -> Self {
self.host_facts = slot;
self
}
}
fn host_plane_unavailable(verb: ControlVerb) -> ControlResponse {
ControlResponse::Err {
code: super::remote_host::HOST_PLANE_UNAVAILABLE_CODE.to_string(),
message: format!(
"this gateway serves no runtime-host facts, so '{verb}' has nothing to answer; \
install one with MobHandleControlHandler::with_host_facts"
),
}
}
fn host_describe_response(
provider: Option<&dyn super::remote_host::HostFactsProvider>,
) -> ControlResponse {
match provider {
Some(provider) => ControlResponse::Host {
facts: provider.describe(),
sig_b64: None,
},
None => host_plane_unavailable(ControlVerb::HostDescribe),
}
}
fn host_health_response(
provider: Option<&dyn super::remote_host::HostFactsProvider>,
) -> ControlResponse {
match provider {
Some(provider) => ControlResponse::HostHealth {
health: provider.health(),
sig_b64: None,
},
None => host_plane_unavailable(ControlVerb::HostHealth),
}
}
fn current_host_facts(
slot: &HostFactsProviderSlot,
) -> Option<std::sync::Arc<dyn super::remote_host::HostFactsProvider>> {
slot.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
impl ControlHandler for MobHandleControlHandler {
fn handle(
&self,
request: ControlRequest,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ControlResponse> + Send + '_>> {
let handle = self.handle.clone();
let identity_runtime = self.identity_authority.current();
let session_service = self.session_service.clone();
let host_facts = self.host_facts.clone();
Box::pin(async move {
match request {
ControlRequest::Wire {
remote_member,
local_peer_spec_address,
local_comms_name,
local_peer_id,
local_pubkey_b64,
nonce: _,
caller: _,
} => {
handle_wire(
&handle,
WireControlParams {
remote_member: &remote_member,
local_peer_spec_address: &local_peer_spec_address,
local_comms_name: &local_comms_name,
local_peer_id: &local_peer_id,
local_pubkey_b64: local_pubkey_b64.as_deref(),
wire: true,
},
identity_runtime.as_ref(),
)
.await
}
ControlRequest::Unwire {
remote_member,
local_peer_spec_address,
local_comms_name,
local_peer_id,
local_pubkey_b64,
nonce: _,
caller: _,
} => {
handle_wire(
&handle,
WireControlParams {
remote_member: &remote_member,
local_peer_spec_address: &local_peer_spec_address,
local_comms_name: &local_comms_name,
local_peer_id: &local_peer_id,
local_pubkey_b64: local_pubkey_b64.as_deref(),
wire: false,
},
identity_runtime.as_ref(),
)
.await
}
ControlRequest::Inject {
remote_member,
content,
nonce: _,
caller: _,
} => {
handle_inject(&handle, &remote_member, content, identity_runtime.as_ref()).await
}
ControlRequest::LookupMember {
remote_member,
nonce: _,
caller: _,
} => {
handle_lookup_member(
&handle,
session_service.as_ref(),
&remote_member,
identity_runtime.as_ref(),
)
.await
}
ControlRequest::HostDescribe {
nonce: _,
caller: _,
} => {
let provider = current_host_facts(&host_facts);
host_describe_response(provider.as_deref())
}
ControlRequest::HostHealth {
nonce: _,
caller: _,
} => {
let provider = current_host_facts(&host_facts);
host_health_response(provider.as_deref())
}
}
})
}
}
fn control_identity_error(error: crate::identity_first::IdentityRuntimeError) -> (String, String) {
let code = match error {
crate::identity_first::IdentityRuntimeError::StaleRuntimeAlias { .. } => {
"stale_runtime_alias"
}
crate::identity_first::IdentityRuntimeError::UnknownIdentity(_) => "unknown_identity",
crate::identity_first::IdentityRuntimeError::InvalidState { .. } => "identity_not_active",
_ => "identity_error",
};
(code.to_string(), error.to_string())
}
const TRACKED_CONTROL_ERROR_PREFIX: &str = "mobkit-control-operation:";
fn tracked_control_error(error: crate::identity_first::IdentityRuntimeError) -> (String, String) {
if let crate::identity_first::IdentityRuntimeError::Internal(message) = &error
&& let Some(encoded) = message.strip_prefix(TRACKED_CONTROL_ERROR_PREFIX)
&& let Ok((code, message)) = serde_json::from_str::<(String, String)>(encoded)
{
return (code, message);
}
control_identity_error(error)
}
async fn run_control_member_operation<T, F, Fut>(
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
member_alias: &str,
operation: F,
) -> Result<T, (String, String)>
where
T: Send + 'static,
F: FnOnce() -> Fut + Send + 'static,
Fut: std::future::Future<Output = Result<T, (String, String)>> + Send + 'static,
{
let member_alias = crate::member_comms_id::runtime_alias_str(member_alias).into_owned();
if let Some(runtime) = identity_runtime {
let target = runtime
.member_alias_lifecycle_target(&member_alias)
.await
.map_err(control_identity_error)?;
if let Some(target) = target {
return crate::identity_first::IdentityRuntime::run_member_alias_targets_operation_tracked(
vec![target],
move || async move {
operation().await.map_err(|error| {
format!(
"{TRACKED_CONTROL_ERROR_PREFIX}{}",
serde_json::to_string(&error)
.unwrap_or_else(|_| "[\"identity_error\",\"unserializable control error\"]".to_string())
)
})
},
)
.await
.map_err(tracked_control_error);
}
}
if crate::member_comms_id::is_reserved_generated_alias(&member_alias) {
return Err((
"identity_authority_unavailable".to_string(),
format!("generated member alias '{member_alias}' requires the owning IdentityRuntime"),
));
}
operation().await
}
struct WireControlParams<'a> {
remote_member: &'a str,
local_peer_spec_address: &'a str,
local_comms_name: &'a str,
local_peer_id: &'a str,
local_pubkey_b64: Option<&'a str>,
wire: bool,
}
fn build_remote_wire_descriptor(
params: &WireControlParams<'_>,
) -> Result<meerkat_core::comms::TrustedPeerDescriptor, (String, String)> {
if params.local_peer_spec_address.starts_with("inproc://") {
return Err((
"inproc_address_rejected".to_string(),
format!(
"peer address '{}' is inproc:// and unreachable from another process; the \
calling gateway must advertise a tcp:// or uds:// member comms address",
params.local_peer_spec_address
),
));
}
let Some(pubkey_b64) = params.local_pubkey_b64.filter(|value| !value.is_empty()) else {
return Err((
"missing_pubkey".to_string(),
"remote wire requests must carry the calling member's Ed25519 transport pubkey; \
an unsigned descriptor would admit any sender at comms ingress"
.to_string(),
));
};
let pubkey = crate::auth::peer_keys::decode_pubkey_b64(pubkey_b64)
.map_err(|err| ("decode".to_string(), format!("local_pubkey_b64: {err}")))?;
meerkat_core::comms::TrustedPeerDescriptor::unsigned_with_pubkey(
params.local_comms_name,
params.local_peer_id,
pubkey,
params.local_peer_spec_address,
)
.map_err(|err| ("peer_spec".to_string(), err))
}
async fn handle_wire(
handle: &meerkat_mob::MobHandle,
params: WireControlParams<'_>,
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
) -> ControlResponse {
let spec = match build_remote_wire_descriptor(¶ms) {
Ok(spec) => spec,
Err((code, message)) => return ControlResponse::Err { code, message },
};
let operation_handle = handle.clone();
let operation_member =
crate::member_comms_id::runtime_alias_str(params.remote_member).into_owned();
let wire = params.wire;
let result =
run_control_member_operation(identity_runtime, params.remote_member, move || async move {
let mid = crate::member_comms_id::mob_member_id(&operation_member);
let result = if wire {
operation_handle
.wire(mid, meerkat_mob::PeerTarget::External(spec))
.await
} else {
operation_handle
.unwire(mid, meerkat_mob::PeerTarget::External(spec))
.await
};
result.map_err(|error| ("mob_error".to_string(), error.to_string()))
})
.await;
match result {
Ok(()) => ControlResponse::Ok { sig_b64: None },
Err((code, message)) => ControlResponse::Err { code, message },
}
}
async fn handle_inject(
handle: &meerkat_mob::MobHandle,
remote_member: &str,
content: serde_json::Value,
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
) -> ControlResponse {
let content_input: meerkat_core::ContentInput = match serde_json::from_value(content) {
Ok(c) => c,
Err(err) => {
return ControlResponse::Err {
code: "decode".to_string(),
message: format!("content: {err}"),
};
}
};
let operation_handle = handle.clone();
let operation_member = crate::member_comms_id::runtime_alias_str(remote_member).into_owned();
match run_control_member_operation(identity_runtime, remote_member, move || async move {
let mid = crate::member_comms_id::mob_member_id(&operation_member);
operation_handle
.member(&mid)
.await
.map_err(|error| ("unknown_member".to_string(), error.to_string()))?
.send(content_input, meerkat_core::types::HandlingMode::Queue)
.await
.map_err(|error| ("mob_error".to_string(), error.to_string()))?;
operation_handle
.resolve_bridge_session_id(&mid)
.await
.map(|session_id| session_id.to_string())
.ok_or_else(|| {
(
"no_session".to_string(),
format!("member '{operation_member}' has no bound bridge session"),
)
})
})
.await
{
Ok(session_id) => ControlResponse::Injected {
session_id,
sig_b64: None,
},
Err((code, message)) => ControlResponse::Err { code, message },
}
}
async fn handle_lookup_member(
handle: &meerkat_mob::MobHandle,
session_service: Option<&std::sync::Arc<dyn meerkat_mob::MobSessionService>>,
remote_member: &str,
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
) -> ControlResponse {
let operation_handle = handle.clone();
let operation_service = session_service.cloned();
let operation_member = crate::member_comms_id::runtime_alias_str(remote_member).into_owned();
match run_control_member_operation(identity_runtime, remote_member, move || async move {
handle_lookup_member_raw(
&operation_handle,
operation_service.as_ref(),
&operation_member,
)
.await
})
.await
{
Ok(response) => response,
Err((code, message)) => ControlResponse::Err { code, message },
}
}
async fn handle_lookup_member_raw(
handle: &meerkat_mob::MobHandle,
session_service: Option<&std::sync::Arc<dyn meerkat_mob::MobSessionService>>,
remote_member: &str,
) -> Result<ControlResponse, (String, String)> {
let mid = crate::member_comms_id::mob_member_id(remote_member);
let mob_id = handle.mob_id().to_string();
let entry = match handle.get_member(&mid).await {
Ok(Some(e)) => e,
Ok(None) => {
return Err((
"unknown_member".to_string(),
format!("member '{remote_member}' not in mob '{mob_id}'"),
));
}
Err(err) => {
return Err(("mob_error".to_string(), err.to_string()));
}
};
let peer_id = match entry.peer_id() {
Some(p) => p.to_string(),
None => {
return Err((
"no_comms".to_string(),
format!("member '{remote_member}' has no comms runtime"),
));
}
};
let comms_name = match meerkat_core::MemberCommsName::new(
mob_id.as_str(),
entry.role.as_str(),
mid.as_str(),
) {
Ok(name) => name.to_string(),
Err(err) => {
return Err((
"invalid_comms_name".to_string(),
format!(
"member '{remote_member}' in mob '{mob_id}' has an invalid comms name component: {err}"
),
));
}
};
let pubkey_b64 = entry.transport_public_key().map(str::to_string);
let mut advertised_address = None;
if let Some(service) = session_service
&& let Some(session_id) = handle.resolve_bridge_session_id(&mid).await
&& let Some(comms) = service.comms_runtime(&session_id).await
{
advertised_address = comms.advertised_address();
}
Ok(ControlResponse::Member {
peer_id,
comms_name,
pubkey_b64,
advertised_address,
sig_b64: None,
})
}
pub async fn serve_tcp_control(
listener: TcpListener,
handler: std::sync::Arc<dyn ControlHandler>,
signer: ControlSignerSlot,
) {
serve_tcp_control_with_authorizer(
listener,
handler,
signer,
std::sync::Arc::new(ControlAuthorizer::open()),
)
.await;
}
pub async fn serve_tcp_control_with_authorizer(
listener: TcpListener,
handler: std::sync::Arc<dyn ControlHandler>,
signer: ControlSignerSlot,
authorizer: std::sync::Arc<ControlAuthorizer>,
) {
loop {
let (stream, _peer_addr) = match listener.accept().await {
Ok(pair) => pair,
Err(err) => {
tracing::warn!(error = %err, "control listener accept failed; exiting");
return;
}
};
let handler = handler.clone();
let signer = std::sync::Arc::clone(&signer);
let authorizer = std::sync::Arc::clone(&authorizer);
tokio::spawn(serve_one_tcp(stream, handler, signer, authorizer));
}
}
#[cfg(unix)]
pub async fn serve_uds_control(
listener: UnixListener,
handler: std::sync::Arc<dyn ControlHandler>,
signer: ControlSignerSlot,
) {
serve_uds_control_with_authorizer(
listener,
handler,
signer,
std::sync::Arc::new(ControlAuthorizer::open()),
)
.await;
}
#[cfg(unix)]
pub async fn serve_uds_control_with_authorizer(
listener: UnixListener,
handler: std::sync::Arc<dyn ControlHandler>,
signer: ControlSignerSlot,
authorizer: std::sync::Arc<ControlAuthorizer>,
) {
loop {
let (stream, _peer_addr) = match listener.accept().await {
Ok(pair) => pair,
Err(err) => {
tracing::warn!(error = %err, "uds control listener accept failed; exiting");
return;
}
};
let handler = handler.clone();
let signer = std::sync::Arc::clone(&signer);
let authorizer = std::sync::Arc::clone(&authorizer);
tokio::spawn(serve_one_uds(stream, handler, signer, authorizer));
}
}
async fn serve_one_tcp(
stream: TcpStream,
handler: std::sync::Arc<dyn ControlHandler>,
signer: ControlSignerSlot,
authorizer: std::sync::Arc<ControlAuthorizer>,
) {
let mut s = ControlStream::Tcp(stream);
serve_one(&mut s, handler, signer, authorizer).await;
}
#[cfg(unix)]
async fn serve_one_uds(
stream: UnixStream,
handler: std::sync::Arc<dyn ControlHandler>,
signer: ControlSignerSlot,
authorizer: std::sync::Arc<ControlAuthorizer>,
) {
let mut s = ControlStream::Uds(stream);
serve_one(&mut s, handler, signer, authorizer).await;
}
async fn serve_one(
stream: &mut ControlStream,
handler: std::sync::Arc<dyn ControlHandler>,
signer: ControlSignerSlot,
authorizer: std::sync::Arc<ControlAuthorizer>,
) {
let payload = match stream.read_frame().await {
Ok(buf) => buf,
Err(err) => {
tracing::debug!(error = %err, "control listener: read failed");
return;
}
};
let request = match serde_json::from_slice::<ControlRequest>(&payload) {
Ok(req) => req,
Err(err) => {
let response = ControlResponse::Err {
code: "decode".to_string(),
message: err.to_string(),
};
let response_payload = serde_json::to_vec(&response).unwrap_or_default();
let _ = stream.write_frame(&response_payload).await;
return;
}
};
let verb = request.verb();
match authorizer.authorize(&request) {
Ok(None) => {}
Ok(Some(caller)) => {
tracing::debug!(
caller = %caller.label,
verb = %verb,
"cross-mob control request admitted by grant"
);
}
Err(denial) => {
tracing::warn!(
code = denial.code(),
verb = %verb,
reason = %denial,
"cross-mob control request refused by grant policy"
);
let response = ControlResponse::Err {
code: denial.code().to_string(),
message: denial.to_string(),
};
let response_payload = serde_json::to_vec(&response).unwrap_or_default();
let _ = stream.write_frame(&response_payload).await;
return;
}
}
let mut response = handler.handle(request).await;
let keys = signer
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
if let Some(keys) = keys {
sign_control_response(&keys, &payload, &mut response);
}
let response_payload = serde_json::to_vec(&response).unwrap_or_default();
let _ = stream.write_frame(&response_payload).await;
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
mod tests {
use super::*;
use crate::identity_first::{
AgentAddressability, AgentIdentity, AgentRuntimeId, CheckpointVersion,
ContinuityGeneration, ContinuityRecord, DurabilityPolicy, DurableAgentSpec, FencingToken,
IdentityLifecycleState, IdentityRuntime, IdentityRuntimeConfig, LeaseGrant,
LocalContinuityStore, LocalLeaseProvider,
};
use std::collections::BTreeMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
struct EchoHandler;
async fn identity_runtime_with_current_alias()
-> Result<Arc<IdentityRuntime>, Box<dyn std::error::Error + Send + Sync>> {
let runtime = Arc::new(IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: Arc::new(LocalContinuityStore::in_memory()?),
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: "cross-mob-control-test".to_string(),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
}));
let identity = AgentIdentity::parse("worker")?;
runtime
.register(
DurableAgentSpec {
identity: identity.clone(),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
placement: None,
},
IdentityLifecycleState::Active,
Some(ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:worker:0")?,
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
}),
Some(LeaseGrant {
identity,
fencing_token: FencingToken::new(1),
ttl: Duration::from_mins(1),
}),
)
.await;
Ok(runtime)
}
impl ControlHandler for EchoHandler {
fn handle(
&self,
request: ControlRequest,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ControlResponse> + Send + '_>>
{
Box::pin(async move {
match request {
ControlRequest::Wire { .. } | ControlRequest::Unwire { .. } => {
ControlResponse::Ok { sig_b64: None }
}
ControlRequest::Inject { remote_member, .. } => ControlResponse::Injected {
session_id: format!("session-for-{remote_member}"),
sig_b64: None,
},
ControlRequest::LookupMember { remote_member, .. } => ControlResponse::Member {
peer_id: format!("peer-id-for-{remote_member}"),
comms_name: format!("mob/role/{remote_member}"),
pubkey_b64: None,
advertised_address: None,
sig_b64: None,
},
ControlRequest::HostDescribe { .. } => ControlResponse::Host {
facts: echo_host_facts(),
sig_b64: None,
},
ControlRequest::HostHealth { .. } => ControlResponse::HostHealth {
health: echo_host_health(),
sig_b64: None,
},
}
})
}
}
fn echo_host_capabilities() -> meerkat_contracts::RuntimeHostCapabilities {
meerkat_contracts::RuntimeHostCapabilities {
contract_version: meerkat_contracts::ContractVersion::CURRENT,
features: meerkat_contracts::RuntimeHostFeatureFlags {
runtime_backed_sessions: true,
mobs: true,
mcp_live: false,
comms: true,
blobs: false,
session_events: true,
session_streams: false,
schedules: false,
skills: false,
event_replay: false,
artifacts: false,
approvals: false,
external_members: false,
secure_remote_rpc: false,
multi_host_mobs: false,
durable_jobs: false,
},
}
}
fn echo_host_facts() -> super::super::remote_host::HostFacts {
super::super::remote_host::HostFacts::new(
"echo-host",
TEST_PUBKEY_B64,
"tcp://127.0.0.1:7801",
echo_host_capabilities(),
)
}
fn echo_host_health() -> meerkat_contracts::RuntimeHostHealth {
meerkat_contracts::RuntimeHostHealth {
contract_version: meerkat_contracts::ContractVersion::CURRENT,
status: meerkat_contracts::RuntimeHostHealthStatus::Ok,
checks: BTreeMap::new(),
}
}
const TEST_PUBKEY_B64: &str = "KioqKioqKioqKioqKioqKioqKioqKioqKioqKioqKio=";
fn wire_params<'a>(
address: &'a str,
peer_id: &'a str,
pubkey_b64: Option<&'a str>,
) -> WireControlParams<'a> {
WireControlParams {
remote_member: "bob",
local_peer_spec_address: address,
local_comms_name: "mob-a/worker/alice",
local_peer_id: peer_id,
local_pubkey_b64: pubkey_b64,
wire: true,
}
}
#[test]
fn remote_wire_descriptor_rejects_inproc_addresses() {
let params = wire_params(
"inproc://mob-a/worker/alice",
"00000000-0000-4000-8000-000000000001",
Some(TEST_PUBKEY_B64),
);
let (code, _) = build_remote_wire_descriptor(¶ms).expect_err("inproc must fail");
assert_eq!(code, "inproc_address_rejected");
}
#[test]
fn remote_wire_descriptor_requires_pubkey() {
for pubkey in [None, Some("")] {
let params = wire_params(
"tcp://127.0.0.1:9001",
"00000000-0000-4000-8000-000000000001",
pubkey,
);
let (code, _) =
build_remote_wire_descriptor(¶ms).expect_err("missing pubkey must fail");
assert_eq!(code, "missing_pubkey");
}
}
#[test]
fn remote_wire_descriptor_builds_signed_descriptor() {
let derived_peer_id =
meerkat_core::comms::PeerId::from_ed25519_pubkey(&[42u8; 32]).to_string();
let params = wire_params(
"tcp://127.0.0.1:9001",
&derived_peer_id,
Some(TEST_PUBKEY_B64),
);
let spec = build_remote_wire_descriptor(¶ms).expect("descriptor");
assert_eq!(spec.pubkey, [42u8; 32]);
assert_eq!(spec.address.endpoint(), "127.0.0.1:9001");
}
#[test]
fn remote_wire_descriptor_rejects_mismatched_pubkey() {
let params = wire_params(
"tcp://127.0.0.1:9001",
"00000000-0000-4000-8000-000000000001",
Some(TEST_PUBKEY_B64),
);
let (code, _) =
build_remote_wire_descriptor(¶ms).expect_err("mismatched pubkey must fail");
assert_eq!(code, "peer_spec");
}
fn signer_with(keys: crate::auth::peer_keys::GatewayPeerKeys) -> ControlSignerSlot {
Arc::new(std::sync::RwLock::new(Some(Arc::new(keys))))
}
#[tokio::test]
async fn signed_responses_verify_against_pinned_gateway_key() {
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let pinned = keys.pubkey_bytes();
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(serve_tcp_control(listener, handler, signer_with(keys)));
let request = ControlRequest::LookupMember {
remote_member: "alice".to_string(),
nonce: Some("nonce-1".to_string()),
caller: None,
};
let payload = serde_json::to_vec(&request).expect("encode");
let endpoint = RemoteEndpoint::Tcp(addr.to_string());
let response =
RemoteControlClient::send_payload(&endpoint, &payload, DEFAULT_CONTROL_TIMEOUT)
.await
.expect("control rpc");
verify_control_response(&pinned, &payload, &response).expect("signature verifies");
let replayed_request = ControlRequest::LookupMember {
remote_member: "alice".to_string(),
nonce: Some("nonce-2".to_string()),
caller: None,
};
let replayed_payload = serde_json::to_vec(&replayed_request).expect("encode");
assert!(
verify_control_response(&pinned, &replayed_payload, &response).is_err(),
"a response must not verify for different request bytes (replay)"
);
let other_key = crate::auth::peer_keys::GatewayPeerKeys::ephemeral().pubkey_bytes();
assert!(
verify_control_response(&other_key, &payload, &response).is_err(),
"a response must not verify against a different pinned key"
);
server.abort();
}
#[test]
fn unsigned_response_fails_verification_when_pinned() {
let pinned = crate::auth::peer_keys::GatewayPeerKeys::ephemeral().pubkey_bytes();
let response = ControlResponse::Ok { sig_b64: None };
let err =
verify_control_response(&pinned, b"{}", &response).expect_err("unsigned must fail");
assert!(err.contains("unsigned"), "got: {err}");
}
#[test]
fn error_responses_skip_verification() {
let pinned = crate::auth::peer_keys::GatewayPeerKeys::ephemeral().pubkey_bytes();
let response = ControlResponse::Err {
code: "unknown_member".to_string(),
message: "nope".to_string(),
};
verify_control_response(&pinned, b"{}", &response).expect("Err passes through");
}
#[tokio::test]
async fn pinned_contact_requires_signed_control_responses() {
use crate::contact_directory::{ContactEntry, MobTransport};
use crate::runtime::cross_mob_remote::RemoteMobProxy;
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(serve_tcp_control(
listener,
handler,
unsigned_control_signer(),
));
let pinned = crate::auth::peer_keys::GatewayPeerKeys::ephemeral().pubkey_bytes();
let strict_entry = ContactEntry {
mob_id: "remote".to_string(),
transport: MobTransport::Tcp(addr.to_string()),
pubkey: Some(pinned),
require_signed_control: None,
};
let proxy = RemoteMobProxy::from_entry(&strict_entry)
.expect("tcp ok")
.expect("some");
let err = proxy
.lookup_member("alice")
.await
.expect_err("unsigned response must fail a pinned contact");
assert!(
matches!(err, RemoteMobError::ControlResponseUnauthenticated { .. }),
"got {err:?}"
);
let opt_out_entry = ContactEntry {
require_signed_control: Some(false),
..strict_entry
};
let proxy = RemoteMobProxy::from_entry(&opt_out_entry)
.expect("tcp ok")
.expect("some");
proxy
.lookup_member("alice")
.await
.expect("explicit opt-out accepts unsigned responses");
server.abort();
}
#[tokio::test]
async fn pinned_contact_accepts_signed_responses() {
use crate::contact_directory::{ContactEntry, MobTransport};
use crate::runtime::cross_mob_remote::RemoteMobProxy;
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let pinned = keys.pubkey_bytes();
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(serve_tcp_control(listener, handler, signer_with(keys)));
let entry = ContactEntry {
mob_id: "remote".to_string(),
transport: MobTransport::Tcp(addr.to_string()),
pubkey: Some(pinned),
require_signed_control: None,
};
let proxy = RemoteMobProxy::from_entry(&entry)
.expect("tcp ok")
.expect("some");
let info = proxy
.lookup_member("alice")
.await
.expect("signed lookup verifies");
assert_eq!(info.peer_id, "peer-id-for-alice");
proxy
.wire_remote(
"alice",
"tcp://127.0.0.1:9001",
"demo/role/alice",
"00000000-0000-4000-8000-000000000001",
None,
)
.await
.expect("signed wire ack verifies");
server.abort();
}
#[test]
fn control_listen_addr_parses_contact_directory_spellings() {
assert_eq!(
ControlListenAddr::parse("tcp://127.0.0.1:9001"),
Ok(ControlListenAddr::Tcp("127.0.0.1:9001".to_string())),
);
assert_eq!(
ControlListenAddr::parse("uds:///var/run/mob.sock"),
Ok(ControlListenAddr::Uds("/var/run/mob.sock".to_string())),
);
assert!(ControlListenAddr::parse("inproc").is_err());
assert!(ControlListenAddr::parse("ftp://nope").is_err());
assert!(ControlListenAddr::parse("127.0.0.1:9001").is_err());
}
#[test]
fn control_listen_addr_display_round_trips() {
for spec in ["tcp://127.0.0.1:9001", "uds:///var/run/mob.sock"] {
let addr = ControlListenAddr::parse(spec).expect("parse");
assert_eq!(addr.to_string(), spec);
}
}
#[tokio::test]
async fn bound_listener_reports_real_port_and_serves() {
let addr = ControlListenAddr::parse("tcp://127.0.0.1:0").expect("parse");
let bound = BoundControlListener::bind(&addr).await.expect("bind");
let advertised = bound.advertised_address().to_string();
let dial = advertised.strip_prefix("tcp://").expect("tcp scheme");
assert!(
!dial.ends_with(":0"),
"advertised address must carry the kernel-assigned port: {advertised}"
);
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(bound.serve(handler, unsigned_control_signer()));
let endpoint = RemoteEndpoint::Tcp(dial.to_string());
let request = ControlRequest::LookupMember {
remote_member: "alice".to_string(),
nonce: None,
caller: None,
};
let response = RemoteControlClient::send(&endpoint, &request, DEFAULT_CONTROL_TIMEOUT)
.await
.expect("control rpc");
assert!(matches!(response, ControlResponse::Member { .. }));
server.abort();
}
#[cfg(unix)]
#[tokio::test]
async fn bound_uds_listener_serves() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("control.sock");
let addr = ControlListenAddr::Uds(path.display().to_string());
let bound = BoundControlListener::bind(&addr).await.expect("bind");
assert_eq!(
bound.advertised_address(),
format!("uds://{}", path.display())
);
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(bound.serve(handler, unsigned_control_signer()));
let endpoint = RemoteEndpoint::Uds(path.display().to_string());
let request = ControlRequest::Inject {
remote_member: "bob".to_string(),
content: serde_json::json!({"text": "hi"}),
nonce: None,
caller: None,
};
let response = RemoteControlClient::send(&endpoint, &request, DEFAULT_CONTROL_TIMEOUT)
.await
.expect("control rpc");
assert_eq!(
response,
ControlResponse::Injected {
session_id: "session-for-bob".to_string(),
sig_b64: None,
},
);
server.abort();
}
#[tokio::test]
async fn tcp_round_trip_inject() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(serve_tcp_control(
listener,
handler,
unsigned_control_signer(),
));
let endpoint = RemoteEndpoint::Tcp(addr.to_string());
let request = ControlRequest::Inject {
remote_member: "alice".to_string(),
content: serde_json::json!({"text": "hello"}),
nonce: None,
caller: None,
};
let response = RemoteControlClient::send(&endpoint, &request, DEFAULT_CONTROL_TIMEOUT)
.await
.expect("control rpc");
assert_eq!(
response,
ControlResponse::Injected {
session_id: "session-for-alice".to_string(),
sig_b64: None,
},
);
server.abort();
}
#[tokio::test]
async fn tcp_round_trip_wire() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(serve_tcp_control(
listener,
handler,
unsigned_control_signer(),
));
let endpoint = RemoteEndpoint::Tcp(addr.to_string());
let request = ControlRequest::Wire {
remote_member: "bob".to_string(),
local_peer_spec_address: "tcp://127.0.0.1:9001".to_string(),
local_comms_name: "demo/role/alice".to_string(),
local_peer_id: "00000000-0000-4000-8000-000000000001".to_string(),
local_pubkey_b64: None,
nonce: None,
caller: None,
};
let response = RemoteControlClient::send(&endpoint, &request, DEFAULT_CONTROL_TIMEOUT)
.await
.expect("control rpc");
assert_eq!(response, ControlResponse::Ok { sig_b64: None });
server.abort();
}
#[tokio::test]
async fn malformed_request_returns_decode_error() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let _server = tokio::spawn(serve_tcp_control(
listener,
handler,
unsigned_control_signer(),
));
let mut stream = TcpStream::connect(addr).await.expect("connect");
stream
.write_all(&u32::to_be_bytes(5))
.await
.expect("write header");
stream.write_all(b"hello").await.expect("write payload");
stream.flush().await.expect("flush");
let mut header = [0u8; 4];
stream.read_exact(&mut header).await.expect("read header");
let len = u32::from_be_bytes(header) as usize;
let mut buf = vec![0u8; len];
stream.read_exact(&mut buf).await.expect("read payload");
let response: ControlResponse = serde_json::from_slice(&buf).expect("decode response");
match response {
ControlResponse::Err { code, .. } => assert_eq!(code, "decode"),
other => panic!("expected decode error, got {other:?}"),
}
}
#[tokio::test]
async fn identity_control_rejects_stale_and_encoded_aliases_before_mutation()
-> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let runtime = identity_runtime_with_current_alias().await?;
let calls = Arc::new(AtomicUsize::new(0));
for stale_alias in [
"rt:worker:1".to_string(),
crate::member_comms_id::mob_member_id_str("rt:worker:1").into_owned(),
] {
let operation_calls = Arc::clone(&calls);
let error =
run_control_member_operation(Some(&runtime), &stale_alias, move || async move {
operation_calls.fetch_add(1, Ordering::SeqCst);
Ok::<_, (String, String)>(())
})
.await
.expect_err("stale generation must fail before the lower plane");
assert_eq!(error.0, "stale_runtime_alias");
}
assert_eq!(calls.load(Ordering::SeqCst), 0);
let operation_calls = Arc::clone(&calls);
run_control_member_operation(Some(&runtime), "rt:worker:0", move || async move {
operation_calls.fetch_add(1, Ordering::SeqCst);
Ok::<_, (String, String)>(())
})
.await
.expect("current generation is admitted");
assert_eq!(calls.load(Ordering::SeqCst), 1);
Ok(())
}
#[tokio::test]
async fn generated_control_alias_without_authority_fails_closed() {
let calls = Arc::new(AtomicUsize::new(0));
let operation_calls = Arc::clone(&calls);
let encoded = crate::member_comms_id::mob_member_id_str("rt:worker:0").into_owned();
let error = run_control_member_operation(None, &encoded, move || async move {
operation_calls.fetch_add(1, Ordering::SeqCst);
Ok::<_, (String, String)>(())
})
.await
.expect_err("generated aliases require identity authority");
assert_eq!(error.0, "identity_authority_unavailable");
assert_eq!(calls.load(Ordering::SeqCst), 0);
}
fn lookup_request(member: &str) -> ControlRequest {
ControlRequest::LookupMember {
remote_member: member.to_string(),
nonce: Some("nonce".to_string()),
caller: None,
}
}
fn inject_request(member: &str) -> ControlRequest {
ControlRequest::Inject {
remote_member: member.to_string(),
content: serde_json::json!({"text": "hi"}),
nonce: Some("nonce".to_string()),
caller: None,
}
}
fn signed_as(
keys: &crate::auth::peer_keys::GatewayPeerKeys,
audience: Option<&str>,
mut request: ControlRequest,
) -> ControlRequest {
sign_control_request_as_caller(keys, audience, &mut request);
request
}
fn table_for(
keys: &crate::auth::peer_keys::GatewayPeerKeys,
verbs: &[ControlVerb],
members: ControlMemberScope,
) -> ControlGrantTable {
let mut table = ControlGrantTable::new();
table.insert(
keys.pubkey_bytes(),
ControlGrant::new("peer-mob", verbs.iter().copied(), members),
);
table
}
#[test]
fn control_verb_config_spellings_round_trip() {
for verb in ControlVerb::all() {
assert_eq!(ControlVerb::parse(verb.as_str()), Some(verb));
}
assert_eq!(ControlVerb::parse("delete_everything"), None);
}
#[test]
fn open_authorizer_admits_unauthenticated_requests() {
let authorizer = ControlAuthorizer::open();
assert!(!authorizer.is_enforcing());
assert_eq!(
authorizer.authorize(&inject_request("bob")),
Ok(None),
"open mode authorizes nothing and refuses nothing"
);
}
#[test]
fn empty_grant_table_denies_every_caller() {
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let authorizer = ControlAuthorizer::with_grants(ControlGrantTable::new());
assert!(authorizer.is_enforcing());
let denial = authorizer
.authorize(&signed_as(&keys, None, lookup_request("bob")))
.expect_err("an empty table grants nothing");
assert_eq!(denial.code(), "caller_not_granted");
}
#[test]
fn unauthenticated_request_is_refused_when_grants_are_enforced() {
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let authorizer = ControlAuthorizer::with_grants(table_for(
&keys,
&ControlVerb::all(),
ControlMemberScope::All,
));
let denial = authorizer
.authorize(&inject_request("bob"))
.expect_err("no credential must be refused");
assert_eq!(denial.code(), "unauthenticated_caller");
}
#[test]
fn granted_caller_is_admitted() {
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let authorizer = ControlAuthorizer::with_grants(table_for(
&keys,
&[ControlVerb::Inject],
ControlMemberScope::members(["bob"]),
));
let admitted = authorizer
.authorize(&signed_as(&keys, None, inject_request("bob")))
.expect("granted verb and member")
.expect("enforcing mode names the caller");
assert_eq!(admitted.label, "peer-mob");
assert_eq!(admitted.pubkey, keys.pubkey_bytes());
}
#[test]
fn verb_and_member_outside_the_grant_are_refused() {
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let authorizer = ControlAuthorizer::with_grants(table_for(
&keys,
&[ControlVerb::Inject],
ControlMemberScope::members(["bob"]),
));
let verb_denial = authorizer
.authorize(&signed_as(&keys, None, lookup_request("bob")))
.expect_err("lookup_member is not granted");
assert_eq!(verb_denial.code(), "verb_not_granted");
let member_denial = authorizer
.authorize(&signed_as(&keys, None, inject_request("carol")))
.expect_err("carol is outside the member scope");
assert_eq!(member_denial.code(), "member_not_granted");
}
#[test]
fn unknown_caller_is_refused() {
let granted = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let stranger = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let authorizer = ControlAuthorizer::with_grants(table_for(
&granted,
&ControlVerb::all(),
ControlMemberScope::All,
));
let denial = authorizer
.authorize(&signed_as(&stranger, None, inject_request("bob")))
.expect_err("a stranger holds no grant");
assert_eq!(denial.code(), "caller_not_granted");
}
#[test]
fn member_scope_matches_the_encoded_roster_alias() {
let scope = ControlMemberScope::members(["rt:worker:0"]);
let encoded = crate::member_comms_id::mob_member_id_str("rt:worker:0").into_owned();
assert_ne!(encoded, "rt:worker:0", "fixture must exercise the encoding");
assert!(scope.contains("rt:worker:0"));
assert!(scope.contains(&encoded));
assert!(!scope.contains("rt:worker:1"));
let encoded_scope = ControlMemberScope::members([encoded.as_str()]);
assert!(encoded_scope.contains("rt:worker:0"));
}
#[test]
fn wildcard_member_entry_widens_to_all() {
assert_eq!(
ControlMemberScope::members(["bob", "*"]),
ControlMemberScope::All
);
}
#[test]
fn tampering_with_the_request_breaks_the_caller_signature() {
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let authorizer = ControlAuthorizer::with_grants(table_for(
&keys,
&ControlVerb::all(),
ControlMemberScope::All,
));
let mut request = signed_as(&keys, None, inject_request("bob"));
let caller = request.caller().cloned().expect("signed");
request = inject_request("carol");
request.set_caller(Some(caller));
let denial = authorizer
.authorize(&request)
.expect_err("a moved signature must not verify");
assert_eq!(denial.code(), "invalid_caller_signature");
}
#[test]
fn claimed_pubkey_without_the_matching_key_is_refused() {
let victim = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let attacker = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let authorizer = ControlAuthorizer::with_grants(table_for(
&victim,
&ControlVerb::all(),
ControlMemberScope::All,
));
let mut request = signed_as(&attacker, None, inject_request("bob"));
if let Some(caller) = request.caller().cloned() {
request.set_caller(Some(ControlCaller {
pubkey_b64: victim.pubkey_b64(),
..caller
}));
}
let denial = authorizer
.authorize(&request)
.expect_err("claiming the granted pubkey must not admit the attacker");
assert_eq!(denial.code(), "invalid_caller_signature");
}
#[test]
fn audience_binding_refuses_a_request_minted_for_another_gateway() {
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let authorizer = ControlAuthorizer::with_grants_for_audience(
table_for(&keys, &ControlVerb::all(), ControlMemberScope::All),
"mob-b",
);
authorizer
.authorize(&signed_as(&keys, Some("mob-b"), inject_request("bob")))
.expect("matching audience is admitted");
let denial = authorizer
.authorize(&signed_as(&keys, Some("mob-c"), inject_request("bob")))
.expect_err("a request minted for mob-c must not spend here");
assert_eq!(denial.code(), "audience_mismatch");
let missing = authorizer
.authorize(&signed_as(&keys, None, inject_request("bob")))
.expect_err("an audience-less request must not spend on a bound listener");
assert_eq!(missing.code(), "audience_mismatch");
}
#[test]
fn inject_signing_payload_survives_the_json_round_trip() {
let content = serde_json::json!({
"text": "hi",
"nested": {
"b": 1.5,
"a": [1, 2, {"z": true, "n": null}],
"unicode": "naïve\nwith newline\tand tab",
},
"big": 9_007_199_254_740_993_u64,
});
let request = ControlRequest::Inject {
remote_member: "bob".to_string(),
content,
nonce: Some("nonce".to_string()),
caller: None,
};
let bytes = serde_json::to_vec(&request).expect("encode");
let decoded: ControlRequest = serde_json::from_slice(&bytes).expect("decode");
assert_eq!(
control_request_signing_payload(&request),
control_request_signing_payload(&decoded),
"a signature minted by the caller must verify against the decoded request"
);
}
#[test]
fn signing_payload_is_not_field_injectable() {
let split = ControlRequest::LookupMember {
remote_member: "a".to_string(),
nonce: Some("b".to_string()),
caller: None,
};
let merged = ControlRequest::LookupMember {
remote_member: "a\n1:b".to_string(),
nonce: None,
caller: None,
};
assert_ne!(
control_request_signing_payload(&split),
control_request_signing_payload(&merged)
);
}
#[test]
fn caller_and_response_signature_contexts_are_distinct() {
assert_ne!(CONTROL_CALLER_SIG_CONTEXT, CONTROL_SIG_CONTEXT);
assert!(
control_request_signing_payload(&lookup_request("bob"))
.starts_with(CONTROL_CALLER_SIG_CONTEXT)
);
}
#[test]
fn grant_toml_parses_verbs_and_member_scope() {
let table = ControlGrantTable::from_toml(&format!(
r#"
[control_grants.ops-mob]
pubkey = "{TEST_PUBKEY_B64}"
verbs = ["inject", "lookup_member"]
members = ["bob"]
"#
))
.expect("parse")
.expect("section is present");
let grant = table.get(&[42u8; 32]).expect("filed under the pubkey");
assert_eq!(grant.label(), "ops-mob");
assert_eq!(
grant.verbs(),
&[ControlVerb::Inject, ControlVerb::LookupMember]
.into_iter()
.collect::<std::collections::BTreeSet<_>>()
);
assert!(grant.members().contains("bob"));
assert!(!grant.members().contains("carol"));
}
#[test]
fn absent_grant_section_leaves_the_listener_open() {
assert_eq!(
ControlGrantTable::from_toml("[mobs]\nremote = \"inproc\"\n").expect("parse"),
None
);
assert!(matches!(
ControlAuthorizer::from_toml("[mobs]\nremote = \"inproc\"\n").expect("parse"),
ControlAuthorizer::Open
));
}
#[test]
fn grant_toml_omitting_members_grants_every_member() {
let table = ControlGrantTable::from_toml(&format!(
r#"
[control_grants.ops-mob]
pubkey = "{TEST_PUBKEY_B64}"
verbs = ["*"]
"#
))
.expect("parse")
.expect("section is present");
let grant = table.get(&[42u8; 32]).expect("filed");
assert_eq!(grant.verbs().len(), ControlVerb::member_plane().len());
assert_eq!(grant.members(), &ControlMemberScope::All);
}
#[test]
fn star_verbs_never_widen_to_the_host_plane() {
let table = ControlGrantTable::from_toml(&format!(
r#"
[control_grants.ops-mob]
pubkey = "{TEST_PUBKEY_B64}"
verbs = ["*"]
"#
))
.expect("parse")
.expect("section is present");
let grant = table.get(&[42u8; 32]).expect("filed");
for verb in ControlVerb::member_plane() {
assert!(
grant.verbs().contains(&verb),
"'*' must still grant member-plane verb {verb}"
);
}
for verb in ControlVerb::all() {
if verb.is_host_plane() {
assert!(
!grant.verbs().contains(&verb),
"'*' must not grant host-plane verb {verb}"
);
}
}
}
#[test]
fn host_plane_verbs_are_named_explicitly_and_ignore_member_scope() {
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let mut table = ControlGrantTable::new();
table.insert(
keys.pubkey_bytes(),
ControlGrant::new(
"placement-controller",
[ControlVerb::HostDescribe],
ControlMemberScope::members(["nobody"]),
),
);
let authorizer = ControlAuthorizer::with_grants(table);
let mut describe = ControlRequest::HostDescribe {
nonce: Some("n1".to_string()),
caller: None,
};
sign_control_request_as_caller(&keys, None, &mut describe);
let admitted = authorizer
.authorize(&describe)
.expect("named host verb is admitted");
assert_eq!(
admitted.map(|caller| caller.label),
Some("placement-controller".to_string())
);
let mut health = ControlRequest::HostHealth {
nonce: Some("n2".to_string()),
caller: None,
};
sign_control_request_as_caller(&keys, None, &mut health);
assert!(matches!(
authorizer.authorize(&health),
Err(ControlAuthzDenial::VerbNotGranted {
verb: ControlVerb::HostHealth,
..
})
));
let mut inject = ControlRequest::Inject {
remote_member: "nobody".to_string(),
content: serde_json::json!({"text": "hi"}),
nonce: Some("n3".to_string()),
caller: None,
};
sign_control_request_as_caller(&keys, None, &mut inject);
assert!(matches!(
authorizer.authorize(&inject),
Err(ControlAuthzDenial::VerbNotGranted { .. })
));
}
#[test]
fn signed_host_facts_do_not_verify_after_tampering() {
let keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let request = ControlRequest::HostDescribe {
nonce: Some("probe-1".to_string()),
caller: None,
};
let request_bytes = serde_json::to_vec(&request).expect("encode request");
let mut response = ControlResponse::Host {
facts: echo_host_facts(),
sig_b64: None,
};
sign_control_response(&keys, &request_bytes, &mut response);
assert!(
verify_control_response(&keys.pubkey_bytes(), &request_bytes, &response).is_ok(),
"positive control: the untampered signed response must verify"
);
let ControlResponse::Host { sig_b64, .. } = &response else {
panic!("expected a Host response");
};
let signature = sig_b64.clone();
let mut tampered_facts = echo_host_facts();
tampered_facts
.placement_labels
.insert("zone".to_string(), "attacker".to_string());
let tampered = ControlResponse::Host {
facts: tampered_facts,
sig_b64: signature.clone(),
};
assert!(
verify_control_response(&keys.pubkey_bytes(), &request_bytes, &tampered).is_err(),
"a tampered placement label must invalidate the host signature"
);
let unsigned = ControlResponse::Host {
facts: echo_host_facts(),
sig_b64: None,
};
assert!(
verify_control_response(&keys.pubkey_bytes(), &request_bytes, &unsigned).is_err(),
"an unsigned host response must never verify"
);
let crossed = ControlResponse::HostHealth {
health: echo_host_health(),
sig_b64: signature,
};
assert!(
verify_control_response(&keys.pubkey_bytes(), &request_bytes, &crossed).is_err(),
"a facts signature must not verify as a health signature"
);
}
#[test]
fn host_plane_is_refused_without_a_facts_provider() {
let refused = host_describe_response(None);
match &refused {
ControlResponse::Err { code, .. } => {
assert_eq!(code, super::super::remote_host::HOST_PLANE_UNAVAILABLE_CODE);
}
other => panic!("expected a typed refusal, got {other:?}"),
}
match host_health_response(None) {
ControlResponse::Err { code, .. } => {
assert_eq!(code, super::super::remote_host::HOST_PLANE_UNAVAILABLE_CODE);
}
other => panic!("expected a typed refusal, got {other:?}"),
}
let provider =
super::super::remote_host::StaticHostFacts::new(echo_host_facts(), echo_host_health());
assert!(matches!(
host_describe_response(Some(&provider)),
ControlResponse::Host { .. }
));
assert!(matches!(
host_health_response(Some(&provider)),
ControlResponse::HostHealth { .. }
));
}
#[test]
fn host_plane_unavailable_is_not_an_authorization_denial() {
assert!(!ControlAuthzDenial::is_denial_code(
super::super::remote_host::HOST_PLANE_UNAVAILABLE_CODE
));
assert!(ControlAuthzDenial::is_denial_code("caller_not_granted"));
}
#[test]
fn grant_toml_fails_closed_on_bad_policy() {
let unknown_verb = ControlGrantTable::from_toml(&format!(
r#"
[control_grants.ops-mob]
pubkey = "{TEST_PUBKEY_B64}"
verbs = ["inject", "drop_database"]
"#
))
.expect_err("unknown verb must fail closed");
assert!(matches!(
unknown_verb,
ControlGrantConfigError::UnknownVerb { .. }
));
let empty_verbs = ControlGrantTable::from_toml(&format!(
r#"
[control_grants.ops-mob]
pubkey = "{TEST_PUBKEY_B64}"
verbs = []
"#
))
.expect_err("an empty verb list is a typo, not a policy");
assert!(matches!(
empty_verbs,
ControlGrantConfigError::InvalidField { .. }
));
let missing_verbs = ControlGrantTable::from_toml(&format!(
r#"
[control_grants.ops-mob]
pubkey = "{TEST_PUBKEY_B64}"
"#
))
.expect_err("verbs is mandatory");
assert!(matches!(
missing_verbs,
ControlGrantConfigError::InvalidField { .. }
));
let duplicate = ControlGrantTable::from_toml(&format!(
r#"
[control_grants.ops-mob]
pubkey = "{TEST_PUBKEY_B64}"
verbs = ["inject"]
[control_grants.other-mob]
pubkey = "{TEST_PUBKEY_B64}"
verbs = ["wire"]
"#
))
.expect_err("one caller must have exactly one grant");
assert!(matches!(
duplicate,
ControlGrantConfigError::DuplicatePubkey { .. }
));
let bad_key = ControlGrantTable::from_toml(
r#"
[control_grants.ops-mob]
pubkey = "not-base64!!"
verbs = ["inject"]
"#,
)
.expect_err("an undecodable pubkey must fail closed");
assert!(matches!(
bad_key,
ControlGrantConfigError::InvalidPubkey { .. }
));
}
#[test]
fn grant_toml_rejects_unknown_keys_rather_than_widening() {
let typo = ControlGrantTable::from_toml(&format!(
r#"
[control_grants.ops-mob]
pubkey = "{TEST_PUBKEY_B64}"
verbs = ["inject"]
member = ["bob"]
"#
))
.expect_err("a singular 'member' typo must not silently grant every member");
assert!(matches!(typo, ControlGrantConfigError::InvalidField { .. }));
let table = ControlGrantTable::from_toml(&format!(
r#"
[control_grants.ops-mob]
pubkey = "{TEST_PUBKEY_B64}"
verbs = ["inject"]
members = ["bob"]
"#
))
.expect("parse")
.expect("section is present");
let grant = table.get(&[42u8; 32]).expect("filed");
assert!(grant.members().contains("bob"));
assert!(!grant.members().contains("carol"));
}
#[tokio::test]
async fn grant_enforcing_listener_refuses_ungranted_verbs_over_tcp() {
use crate::contact_directory::{ContactEntry, MobTransport};
use crate::runtime::cross_mob_remote::{RemoteMobError, RemoteMobProxy};
let caller_keys = Arc::new(crate::auth::peer_keys::GatewayPeerKeys::ephemeral());
let authorizer = Arc::new(ControlAuthorizer::with_grants(table_for(
&caller_keys,
&[ControlVerb::LookupMember],
ControlMemberScope::members(["alice"]),
)));
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(serve_tcp_control_with_authorizer(
listener,
handler,
unsigned_control_signer(),
authorizer,
));
let entry = ContactEntry {
mob_id: "remote".to_string(),
transport: MobTransport::Tcp(addr.to_string()),
pubkey: None,
require_signed_control: None,
};
let proxy = RemoteMobProxy::from_entry_with_caller(&entry, Some(Arc::clone(&caller_keys)))
.expect("tcp ok")
.expect("some");
proxy
.lookup_member("alice")
.await
.expect("granted verb on a granted member is admitted");
let member_denied = proxy
.lookup_member("bob")
.await
.expect_err("bob is outside the member scope");
assert!(
matches!(
member_denied,
RemoteMobError::ControlRequestUnauthorized { ref code, .. }
if code == "member_not_granted"
),
"got {member_denied:?}"
);
let verb_denied = proxy
.inject_message("alice", serde_json::json!({"text": "hi"}))
.await
.expect_err("inject is outside the verb scope");
assert!(
matches!(
verb_denied,
RemoteMobError::ControlRequestUnauthorized { ref code, .. }
if code == "verb_not_granted"
),
"got {verb_denied:?}"
);
let anonymous = RemoteMobProxy::from_entry(&entry)
.expect("tcp ok")
.expect("some");
let anonymous_denied = anonymous
.lookup_member("alice")
.await
.expect_err("an unauthenticated caller must be refused");
assert!(
matches!(
anonymous_denied,
RemoteMobError::ControlRequestUnauthorized { ref code, .. }
if code == "unauthenticated_caller"
),
"got {anonymous_denied:?}"
);
server.abort();
}
#[tokio::test]
async fn open_listener_serves_authenticated_callers_unchanged() {
use crate::contact_directory::{ContactEntry, MobTransport};
use crate::runtime::cross_mob_remote::RemoteMobProxy;
let caller_keys = Arc::new(crate::auth::peer_keys::GatewayPeerKeys::ephemeral());
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(serve_tcp_control(
listener,
handler,
unsigned_control_signer(),
));
let entry = ContactEntry {
mob_id: "remote".to_string(),
transport: MobTransport::Tcp(addr.to_string()),
pubkey: None,
require_signed_control: None,
};
let proxy = RemoteMobProxy::from_entry_with_caller(&entry, Some(caller_keys))
.expect("tcp ok")
.expect("some");
let info = proxy
.lookup_member("alice")
.await
.expect("an open listener ignores the credential");
assert_eq!(info.peer_id, "peer-id-for-alice");
server.abort();
}
#[tokio::test]
async fn signed_caller_and_signed_response_compose() {
use crate::contact_directory::{ContactEntry, MobTransport};
use crate::runtime::cross_mob_remote::RemoteMobProxy;
let caller_keys = Arc::new(crate::auth::peer_keys::GatewayPeerKeys::ephemeral());
let server_keys = crate::auth::peer_keys::GatewayPeerKeys::ephemeral();
let pinned = server_keys.pubkey_bytes();
let authorizer = Arc::new(ControlAuthorizer::with_grants(table_for(
&caller_keys,
&ControlVerb::all(),
ControlMemberScope::All,
)));
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handler: Arc<dyn ControlHandler> = Arc::new(EchoHandler);
let server = tokio::spawn(serve_tcp_control_with_authorizer(
listener,
handler,
signer_with(server_keys),
authorizer,
));
let entry = ContactEntry {
mob_id: "remote".to_string(),
transport: MobTransport::Tcp(addr.to_string()),
pubkey: Some(pinned),
require_signed_control: None,
};
let proxy = RemoteMobProxy::from_entry_with_caller(&entry, Some(caller_keys))
.expect("tcp ok")
.expect("some");
let info = proxy
.lookup_member("alice")
.await
.expect("authorized request, authenticated answer");
assert_eq!(info.peer_id, "peer-id-for-alice");
server.abort();
}
}