use std::time::Duration;
use crate::{Error, Result};
use zenkey::origin::{HostId, ServiceOrigin};
use zenkey::qos::QosProfile;
use zenkey::{Declared, Fanout, ProcedureKind};
use zenoh::Session;
use crate::bus::query::FleetAnswer;
use crate::model::registry::SliceSet;
use crate::report::{
CallAnswer, CallError, CallOutcome, CallReport, ConcurrentLane, HlcReference, TRACE_CHAIN_RULE,
TRACE_EXCLUDED, TraceReport,
};
pub struct Publication {
publisher: zenoh::pubsub::Publisher<'static>,
encoding: Option<String>,
}
impl std::fmt::Debug for Publication {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Publication")
.field("key", &self.publisher.key_expr().as_str())
.finish_non_exhaustive()
}
}
pub async fn declare_publication(
session: &Session,
key: &str,
qos: QosProfile,
encoding: Option<&str>,
) -> Result<Publication> {
let publisher = session
.declare_publisher(key.to_string())
.reliability(qos.reliability())
.congestion_control(qos.congestion_control())
.priority(qos.priority())
.express(qos.express())
.await
.map_err(|e| Error::bus("declare publisher", key, e))?;
Ok(Publication {
publisher,
encoding: encoding.map(str::to_string),
})
}
impl Publication {
pub async fn send(&self, payload: Vec<u8>, attachment: Option<Vec<u8>>) -> Result<()> {
self.send_stamped(payload, attachment, None).await
}
pub async fn send_stamped(
&self,
payload: Vec<u8>,
attachment: Option<Vec<u8>>,
timestamp: Option<zenoh::time::Timestamp>,
) -> Result<()> {
let put = self.publisher.put(payload);
let put = match &self.encoding {
Some(e) => put.encoding(e.as_str()),
None => put,
};
let put = match attachment {
Some(a) => put.attachment(a),
None => put,
};
let put = match timestamp {
Some(ts) => put.timestamp(ts),
None => put,
};
put.await
.map_err(|e| Error::bus("put", self.publisher.key_expr().as_str(), e))
}
pub async fn retire(&self) -> Result<()> {
self.publisher
.delete()
.await
.map_err(|e| Error::bus("delete", self.publisher.key_expr().as_str(), e))
}
pub async fn undeclare(self) -> Result<()> {
self.publisher
.undeclare()
.await
.map_err(|e| Error::bus("undeclare publisher", "", e))
}
pub async fn matching_status(&self) -> Result<bool> {
self.publisher
.matching_status()
.await
.map(|s| s.matching())
.map_err(|e| Error::bus("matching status", "", e))
}
pub async fn matching_events(&self) -> Result<MatchingEvents> {
let listener = self
.publisher
.matching_listener()
.await
.map_err(|e| Error::bus("matching listener", "", e))?;
Ok(MatchingEvents { listener })
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RetireClass {
State {
registered: bool,
ttl_s: Option<i64>,
},
NonState { class: String },
Unclassified { reason: String },
}
pub fn check_retire(
base: &str,
key: &str,
slices: Option<&SliceSet>,
force: bool,
) -> Result<RetireClass> {
if key.contains('*') || key.contains('$') {
return Err(Error::unaskable(
key,
"is a wildcard — a tombstone is addressed to one concrete key; a \
wildcard delete is not an operator act, it is a blast radius \
(RFC 04 §1.2, v1.12). Not overridable.",
));
}
let facts = crate::model::facts::describe_key(base, key, slices).facts;
use crate::model::facts::{ClassKind, KeyShape, Registration};
match &facts.shape {
KeyShape::V1(v) if v.class_kind == ClassKind::State => {
let (registered, ttl_s) = match &facts.registration {
Registration::Registered(s) => (true, s.ttl_s),
_ => (false, None),
};
Ok(RetireClass::State { registered, ttl_s })
}
KeyShape::V1(v) if matches!(v.class_kind, ClassKind::Telemetry | ClassKind::Events) => {
if force {
return Ok(RetireClass::NonState {
class: v.class.clone(),
});
}
Err(Error::unaskable(
key,
format!(
"is {}-shaped — RFC 04 §1: a delete there is meaningless and \
MUST NOT be sent by the class's publisher. Retiring it anyway \
is an operator cleanup (RFC 04 §1.2, v1.12) — pass --i-know \
to mean it.",
v.class
),
))
}
KeyShape::V1(v) => {
if force {
return Ok(RetireClass::NonState {
class: v.class.clone(),
});
}
Err(Error::unaskable(
key,
format!(
"sits on the {} plane — a plane key answers GETs or carries \
frames; a tombstone there is at most a storage purge \
(RFC 04 §1.2, v1.12) — pass --i-know to mean it.",
v.class
),
))
}
KeyShape::NotUnderBase | KeyShape::Unparsed { .. } => {
let reason = match &facts.shape {
KeyShape::Unparsed { reason } => reason.clone(),
_ => format!("not under base {base:?}"),
};
if force {
return Ok(RetireClass::Unclassified { reason });
}
Err(Error::unaskable(
key,
format!(
"cannot be classified under base {base:?} ({reason}) — 'not \
asked' is not 'state' (RFC 09 §5.1 O4); pass --i-know to \
retire an unclassified key."
),
))
}
}
}
pub struct MatchingEvents {
listener: zenoh::matching::MatchingListener<
zenoh::handlers::FifoChannelHandler<zenoh::matching::MatchingStatus>,
>,
}
impl MatchingEvents {
pub(crate) async fn for_querier(querier: &zenoh::query::Querier<'_>) -> Result<Self> {
let listener = querier
.matching_listener()
.await
.map_err(|e| Error::bus("matching listener", "", e))?;
Ok(MatchingEvents { listener })
}
pub async fn recv(&self) -> Option<bool> {
self.listener.recv_async().await.ok().map(|s| s.matching())
}
pub fn stream(&self) -> impl futures_core::Stream<Item = bool> + '_ {
futures_util::StreamExt::map(self.listener.stream(), |s| s.matching())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CallTarget {
Host(HostId),
Fleet,
Service(ServiceOrigin),
}
impl CallTarget {
pub fn parse(s: &str) -> Result<CallTarget> {
if s == "*" {
return Ok(CallTarget::Fleet);
}
if s.starts_with('@') {
return Ok(CallTarget::Service(ServiceOrigin::new(s)?));
}
HostId::parse(s).map(CallTarget::Host).map_err(|e| {
Error::unaskable(
"origin",
format!("{e} — a hostname is not an origin; resolve it first (RFC 06 §6)"),
)
})
}
}
fn attachment_value(bytes: &[u8]) -> serde_json::Value {
if let Ok(v) = serde_json::from_slice::<serde_json::Value>(bytes) {
v
} else if let Ok(s) = std::str::from_utf8(bytes) {
serde_json::Value::String(s.to_string())
} else {
serde_json::Value::String(format!("<{} bytes>", bytes.len()))
}
}
pub struct CallSpec<'a> {
pub target: &'a CallTarget,
pub producer: &'a str,
pub procedure: &'a str,
pub params: &'a [String],
pub body: Option<Vec<u8>>,
pub attachment: Option<Vec<u8>>,
pub timeout: Duration,
pub slices: Option<&'a SliceSet>,
}
pub async fn call(fleet: &crate::Fleet<'_>, spec: CallSpec<'_>) -> Result<CallReport> {
let (key, timeout, answers) = call_answers(fleet, spec).await?;
Ok(project_call(key, timeout, &answers))
}
async fn call_answers(
fleet: &crate::Fleet<'_>,
spec: CallSpec<'_>,
) -> Result<(String, Duration, Vec<FleetAnswer>)> {
let CallSpec {
target,
producer,
procedure,
params,
body,
attachment,
timeout,
slices,
} = spec;
if matches!(target, CallTarget::Fleet)
&& let Some(slices) = slices
&& let Some(slice) = slices.get(producer)
&& let Some(proc_decl) = slice.procedures.iter().find(|p| p.path == procedure)
{
let forbidden = match proc_decl.fanout.as_ref().and_then(Declared::known) {
Some(Fanout::Forbidden) => true,
Some(Fanout::Allowed) => false,
None => matches!(
proc_decl.kind.as_ref().and_then(Declared::known),
Some(ProcedureKind::Write)
),
};
if forbidden {
let declared = match proc_decl.fanout.as_ref() {
Some(f) if f.is(&Fanout::Forbidden) => {
"declares fanout = \"forbidden\"".to_string()
}
Some(f) => format!(
"declares fanout = {:?}, a token this build does not know — \
RFC 08 §2 defaults a write to forbidden and an unreadable \
spelling is not a licence",
f.token()
),
None => "is a write with no declared fanout, which defaults to forbidden \
(RFC 08 §2)"
.to_string(),
};
return Err(Error::unaskable(
format!("procedure {producer}/{procedure}"),
format!(
"{declared} — a fleet (`*`) call to it is refused \
(RFC 05 §2.1); name one origin"
),
));
}
}
let segments: Vec<&str> = procedure.split('/').collect();
let relative = match target {
CallTarget::Host(id) => {
let origin = zenkey::origin::RemoteOrigin::from_host(id.clone());
zenkey::selector::rpc_at(&origin, producer, &segments).to_string()
}
CallTarget::Fleet => zenkey::selector::fleet_rpc(producer, &segments).to_string(),
CallTarget::Service(origin) => zenkey::selector::service_rpc(origin, &segments).to_string(),
};
let mut key = fleet.wire(relative);
if !params.is_empty() {
key.push('?');
key.push_str(¶ms.join(";"));
}
let answers = crate::bus::query::fleet_get(
fleet,
&key,
&crate::bus::query::GetOpts::new(timeout)
.payload(body)
.attachment(attachment),
)
.await?;
Ok((key, timeout, answers))
}
fn project_call(key: String, timeout: Duration, answers: &[FleetAnswer]) -> CallReport {
CallReport {
key,
timeout_s: timeout.as_secs_f64(),
answers: answers
.iter()
.map(|a| {
let (att, att_bytes) = match &a.attachment {
Some(z) => {
let bytes = z.to_bytes();
(Some(attachment_value(&bytes)), Some(bytes.len()))
}
None => (None, None),
};
let outcome = match &a.answer {
crate::bus::query::Answer::Value(bytes) => {
let bytes = bytes.to_bytes();
match serde_json::from_slice::<serde_json::Value>(&bytes) {
Ok(v) => CallOutcome::Ok {
value: Some(v),
text: None,
},
Err(_) => CallOutcome::Ok {
value: None,
text: Some(String::from_utf8_lossy(&bytes).to_string()),
},
}
}
crate::bus::query::Answer::Error { name, message } => {
CallOutcome::Err(CallError {
name: name.clone(),
message: message.clone(),
})
}
};
CallAnswer {
origin: a.origin.clone(),
outcome,
attachment: att,
attachment_bytes: att_bytes,
}
})
.collect(),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct TraceSpec {
pub window: Duration,
}
const ORIGIN_CAPACITY: usize = 4096;
const FLEET_CAPACITY: usize = 8192;
pub async fn call_traced(
fleet: &crate::Fleet<'_>,
spec: CallSpec<'_>,
trace: TraceSpec,
) -> Result<TraceReport> {
use crate::bus::monitor::{FleetEvent, Monitor, MonitorSpec, StreamItem};
use crate::model::examples::Examples;
use crate::model::facts::describe_key;
use crate::model::timeline::TimelineRow;
use crate::model::trace::{TraceTarget, idiom_of, trace_row};
use zenkey::selector::{Scope, all_under};
let (scope, producer) = match spec.target {
CallTarget::Fleet => {
return Err(Error::unaskable(
"--trace",
"a trace attributes what it sees to one origin, and a fleet (`*`) call \
has none to attribute to — name one origin",
));
}
CallTarget::Host(id) => (
Scope::origin(&zenkey::origin::RemoteOrigin::from_host(id.clone())),
Some(spec.producer.to_string()),
),
CallTarget::Service(o) => (Scope::origin(o), None),
};
let origin = scope.chunk().to_string();
let slices = spec.slices;
let idiom = idiom_of(
slices
.and_then(|s| s.get(spec.producer))
.and_then(|s| s.procedures.iter().find(|p| p.path == spec.procedure)),
);
let target = TraceTarget {
origin: origin.clone(),
producer,
chain_chunk: spec
.procedure
.split('/')
.next()
.unwrap_or_default()
.to_string(),
registry_loaded: slices.is_some(),
};
let base = fleet.base();
let origin_scope = fleet.wire(all_under(scope));
let fleet_scope = fleet.wire(all_under(Scope::fleet()));
let session = fleet.session();
let origin_monitor = Monitor::start(
session,
MonitorSpec {
capacity: ORIGIN_CAPACITY,
..MonitorSpec::default()
},
)
.await?;
let mut origin_events = origin_monitor.events();
origin_monitor.watch(&origin_scope).await?;
let fleet_monitor = Monitor::start(
session,
MonitorSpec {
capacity: FLEET_CAPACITY,
..MonitorSpec::default()
},
)
.await?;
let mut fleet_events = fleet_monitor.events();
fleet_monitor.watch(&fleet_scope).await?;
let t0 = std::time::Instant::now();
let t0_unix_s = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs_f64())
.unwrap_or(0.0);
let (key, timeout, answers) = call_answers(fleet, spec).await?;
let call_returned_ms = t0.elapsed().as_secs_f64() * 1_000.0;
let reply_hlc = answers.iter().find_map(|a| a.timestamp);
let call = project_call(key, timeout, &answers);
let mut attributed = Vec::new();
let mut same_origin = Vec::new();
let mut pending_attributed = 0u64;
let mut pending_same_origin = 0u64;
let mut dropped = 0u64;
let mut concurrent_samples = 0u64;
let mut concurrent_dropped = 0u64;
let mut concurrent_keys = std::collections::HashSet::new();
let mut concurrent_examples = Examples::new(crate::judge::common::EXPANSION_CAP);
let reply_ntp64 = reply_hlc.map(|t| t.get_time().as_u64());
let deadline = tokio::time::sleep(trace.window);
tokio::pin!(deadline);
let mut origin_open = true;
let mut fleet_open = true;
while origin_open || fleet_open {
tokio::select! {
() = &mut deadline => break,
item = origin_events.recv(), if origin_open => match item {
None => origin_open = false,
Some(StreamItem::Dropped(n)) => {
dropped += n;
pending_attributed += n;
pending_same_origin += n;
}
Some(StreamItem::Event(FleetEvent::Sample(view))) => {
let desc = describe_key(base, &view.key, slices);
let Some(relation) = target.relation_of(&desc) else {
continue;
};
let row = TimelineRow::from_view(&view, t0, base);
let (lane, pending) = match relation {
crate::report::TraceRelation::DeclaredChain =>
(&mut attributed, &mut pending_attributed),
_ => (&mut same_origin, &mut pending_same_origin),
};
let break_before = (*pending > 0).then_some(*pending);
*pending = 0;
lane.push(trace_row(&row, relation, reply_ntp64, break_before));
}
Some(StreamItem::Event(_)) => {}
},
item = fleet_events.recv(), if fleet_open => match item {
None => fleet_open = false,
Some(StreamItem::Dropped(n)) => concurrent_dropped += n,
Some(StreamItem::Event(FleetEvent::Sample(view))) => {
let desc = describe_key(base, &view.key, None);
if target.relation_of(&desc).is_some() {
continue;
}
concurrent_samples += 1;
if concurrent_keys.insert(view.key.clone()) {
concurrent_examples.push_with(|| view.key.clone());
}
}
Some(StreamItem::Event(_)) => {}
},
}
}
let keys_evicted = origin_monitor.core().keys_evicted();
origin_monitor.stop();
fleet_monitor.stop();
Ok(TraceReport {
call,
scopes: vec![origin_scope, fleet_scope],
excluded: TRACE_EXCLUDED,
window_s: trace.window.as_secs_f64(),
subscribed_before_call: true,
t0_unix_s,
call_returned_ms,
hlc_reference: if reply_hlc.is_some() {
HlcReference::Reply
} else {
HlcReference::None
},
reply_hlc: reply_hlc.map(|t| t.to_string()),
chain_rule: TRACE_CHAIN_RULE,
registry_loaded: slices.is_some(),
idiom,
attributed,
same_origin,
concurrent: ConcurrentLane {
samples: concurrent_samples,
keys: concurrent_keys.len() as u64,
examples: concurrent_examples.into_vec(),
dropped: concurrent_dropped,
},
dropped,
keys_evicted,
})
}
#[cfg(test)]
mod tests {
use super::*;
use zenkey::slice::{ProcedureDecl, RegistrySlice, SubjectDecl};
fn slice_with_state_subject() -> SliceSet {
let mut health = SubjectDecl::new("health", zenkey::Class::State);
health.type_name = "Health".into();
health.ttl_s = Some(900);
let mut slice = RegistrySlice::new("1.0", "t", "sysinfo");
slice.subjects = vec![health];
SliceSet::from_slices(vec![slice])
}
#[test]
fn a_wildcard_retire_is_refused_unconditionally() {
for force in [false, true] {
let err = check_retire("", "v1/h-3fa9c2d41b7e/state/sysinfo/**", None, force)
.unwrap_err()
.to_string();
assert!(err.contains("blast radius"), "{err}");
}
}
#[test]
fn a_state_key_retires_without_a_registry() {
let got = check_retire("", "v1/h-3fa9c2d41b7e/state/sysinfo/health", None, false).unwrap();
assert_eq!(
got,
RetireClass::State {
registered: false,
ttl_s: None
}
);
let slices = slice_with_state_subject();
let got = check_retire(
"",
"v1/h-3fa9c2d41b7e/state/sysinfo/health",
Some(&slices),
false,
)
.unwrap();
assert_eq!(
got,
RetireClass::State {
registered: true,
ttl_s: Some(900)
}
);
}
#[test]
fn a_telemetry_retire_needs_i_know_and_cites_the_rfc() {
let key = "v1/h-3fa9c2d41b7e/telemetry/sysinfo/cpu/usage";
let err = check_retire("", key, None, false).unwrap_err().to_string();
assert!(err.contains("MUST NOT"), "{err}");
assert!(err.contains("v1.12"), "{err}");
assert!(err.contains("--i-know"), "{err}");
assert_eq!(
check_retire("", key, None, true).unwrap(),
RetireClass::NonState {
class: "telemetry".to_string()
}
);
}
#[test]
fn a_plane_retire_needs_i_know_too() {
let key = "v1/h-3fa9c2d41b7e/@rpc/sysinfo/introspect";
let err = check_retire("", key, None, false).unwrap_err().to_string();
assert!(err.contains("plane"), "{err}");
assert!(matches!(
check_retire("", key, None, true).unwrap(),
RetireClass::NonState { class } if class == "@rpc"
));
}
#[test]
fn an_unclassified_retire_needs_i_know_and_names_o4() {
let err = check_retire("", "some/foreign/key", None, false)
.unwrap_err()
.to_string();
assert!(err.contains("O4"), "{err}");
assert!(matches!(
check_retire("", "some/foreign/key", None, true).unwrap(),
RetireClass::Unclassified { .. }
));
let err = check_retire("acme", "other/v1/h-3fa9c2d41b7e/state/x/y", None, false)
.unwrap_err()
.to_string();
assert!(err.contains("cannot be classified"), "{err}");
}
fn slice_with_proc(kind: &str, fanout: Option<&str>) -> SliceSet {
let mut trigger = ProcedureDecl::new("capture/trigger");
trigger.kind = Some(Declared::parse(kind));
trigger.reply = Some("Ack".into());
trigger.fanout = fanout.map(Declared::parse);
trigger.idempotent = Some(false);
let mut slice = RegistrySlice::new("1.0", "t", "netring");
slice.procedures = vec![trigger];
SliceSet::from_slices(vec![slice])
}
#[test]
fn call_targets_parse_and_validate() {
assert_eq!(CallTarget::parse("*").unwrap(), CallTarget::Fleet);
assert!(matches!(
CallTarget::parse("@catalog").unwrap(),
CallTarget::Service(_)
));
assert!(matches!(
CallTarget::parse("h-3fa9c2d41b7e").unwrap(),
CallTarget::Host(_)
));
let err = CallTarget::parse("toolbx").unwrap_err().to_string();
assert!(err.contains("RFC 06 §6"), "{err}");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn fleet_calls_to_forbidden_fanout_are_refused() {
let session = crate::bus::session::open(&[], &[], false).await.unwrap();
let slices = slice_with_proc("write", Some("forbidden"));
let err = call(
&crate::Fleet::new(&session, ""),
CallSpec {
target: &CallTarget::Fleet,
producer: "netring",
procedure: "capture/trigger",
params: &[],
body: None,
attachment: None,
timeout: Duration::from_millis(100),
slices: Some(&slices),
},
)
.await
.unwrap_err()
.to_string();
assert!(err.contains("fanout"), "{err}");
assert!(err.contains("RFC 05 §2.1"), "{err}");
let err = call(
&crate::Fleet::new(&session, ""),
CallSpec {
target: &CallTarget::Fleet,
producer: "netring",
procedure: "capture/trigger",
params: &[],
body: None,
attachment: None,
timeout: Duration::from_millis(100),
slices: Some(&slice_with_proc("write", None)),
},
)
.await
.unwrap_err()
.to_string();
assert!(err.contains("defaults to forbidden"), "{err}");
assert!(err.contains("RFC 08 §2"), "{err}");
assert!(err.contains("RFC 05 §2.1"), "{err}");
let err = call(
&crate::Fleet::new(&session, ""),
CallSpec {
target: &CallTarget::Fleet,
producer: "netring",
procedure: "capture/trigger",
params: &[],
body: None,
attachment: None,
timeout: Duration::from_millis(100),
slices: Some(&slice_with_proc("write", Some("per-iface"))),
},
)
.await
.unwrap_err()
.to_string();
assert!(err.contains("per-iface"), "{err}");
assert!(err.contains("does not know"), "{err}");
assert!(
!err.contains("declares fanout = \"forbidden\""),
"the slice declared no such thing: {err}"
);
assert!(err.contains("RFC 05 §2.1"), "{err}");
for slices in [
slice_with_proc("write", Some("allowed")),
slice_with_proc("read", None),
] {
let report = call(
&crate::Fleet::new(&session, ""),
CallSpec {
target: &CallTarget::Fleet,
producer: "netring",
procedure: "capture/trigger",
params: &[],
body: None,
attachment: None,
timeout: Duration::from_millis(100),
slices: Some(&slices),
},
)
.await
.unwrap();
assert_eq!(report.exit_code(), 2, "silence stays exit 2");
}
}
}