use std::time::Duration;
use anyhow::{Result, anyhow, bail};
use zenkey::origin::{HostId, ServiceOrigin};
use zenkey::qos::QosProfile;
use zenoh::Session;
use crate::registry::SliceSet;
use crate::report::{CallAnswer, CallError, CallReport};
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| anyhow!("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<()> {
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,
};
put.await
.map_err(|e| anyhow!("put {}: {e}", self.publisher.key_expr()))
}
pub async fn retire(&self) -> Result<()> {
self.publisher
.delete()
.await
.map_err(|e| anyhow!("delete {}: {e}", self.publisher.key_expr()))
}
pub async fn undeclare(self) -> Result<()> {
self.publisher
.undeclare()
.await
.map_err(|e| anyhow!("undeclare publisher: {e}"))
}
pub async fn matching_status(&self) -> Result<bool> {
self.publisher
.matching_status()
.await
.map(|s| s.matching())
.map_err(|e| anyhow!("matching status: {e}"))
}
pub async fn matching_events(&self) -> Result<MatchingEvents> {
let listener = self
.publisher
.matching_listener()
.await
.map_err(|e| anyhow!("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('$') {
bail!(
"{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::facts::describe_key(base, key, slices).facts;
use crate::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(),
});
}
bail!(
"{key} is {class}-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.",
class = v.class
);
}
KeyShape::V1(v) => {
if force {
return Ok(RetireClass::NonState {
class: v.class.clone(),
});
}
bail!(
"{key} sits on the {class} 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.",
class = 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 });
}
bail!(
"cannot classify {key} 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| anyhow!("matching listener: {e}"))?;
Ok(MatchingEvents { listener })
}
pub async fn recv(&self) -> Option<bool> {
self.listener.recv_async().await.ok().map(|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).map_err(|e| anyhow!("{e}"))?,
));
}
HostId::parse(s)
.map(CallTarget::Host)
.map_err(|e| anyhow!("{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()))
}
}
#[allow(clippy::too_many_arguments)]
pub async fn call(
session: &Session,
base: &str,
target: &CallTarget,
producer: &str,
procedure: &str,
params: &[String],
body: Option<Vec<u8>>,
attachment: Option<Vec<u8>>,
timeout: Duration,
slices: Option<&SliceSet>,
) -> Result<CallReport> {
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)
&& proc_decl.fanout.as_deref() == Some("forbidden")
{
bail!(
"procedure {producer}/{procedure} declares fanout = \"forbidden\" — 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 = zenkey::grammar::with_base(base, relative);
if !params.is_empty() {
key.push('?');
key.push_str(¶ms.join(";"));
}
let answers =
crate::query::fleet_get_call(session, base, &key, body, attachment, timeout).await?;
Ok(CallReport {
key: key.clone(),
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),
};
match &a.answer {
crate::query::Answer::Value(bytes) => {
let bytes = bytes.to_bytes();
match serde_json::from_slice::<serde_json::Value>(&bytes) {
Ok(v) => CallAnswer {
origin: a.origin.clone(),
ok: true,
value: Some(v),
text: None,
attachment: att,
attachment_bytes: att_bytes,
error: None,
},
Err(_) => CallAnswer {
origin: a.origin.clone(),
ok: true,
value: None,
text: Some(String::from_utf8_lossy(&bytes).to_string()),
attachment: att,
attachment_bytes: att_bytes,
error: None,
},
}
}
crate::query::Answer::Error { name, message } => CallAnswer {
origin: a.origin.clone(),
ok: false,
value: None,
text: None,
attachment: att,
attachment_bytes: att_bytes,
error: Some(CallError {
name: name.clone(),
message: message.clone(),
}),
},
}
})
.collect(),
})
}
#[cfg(test)]
mod tests {
use super::*;
use zenkey::slice::{ProcedureDecl, RegistrySlice, SubjectDecl};
fn slice_with_state_subject() -> SliceSet {
SliceSet::from_slices(vec![RegistrySlice {
version: "1.0".into(),
app: "t".into(),
convention: 1,
name: "sysinfo".into(),
service_origin: None,
description: None,
subjects: vec![SubjectDecl {
path: "health".into(),
class: "state".into(),
type_name: "Health".into(),
common: None,
since: None,
description: None,
qos: None,
ttl_s: Some(900),
unit: None,
rate: None,
cardinality: None,
encoding: None,
}],
procedures: vec![],
blob: vec![],
media: vec![],
deprecated: vec![],
}])
}
#[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".into()
}
);
}
#[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 classify"), "{err}");
}
fn slice_with_proc(fanout: Option<&str>) -> SliceSet {
SliceSet::from_slices(vec![RegistrySlice {
version: "1.0".into(),
app: "t".into(),
convention: 1,
name: "netring".into(),
service_origin: None,
description: None,
subjects: vec![],
procedures: vec![ProcedureDecl {
path: "capture/trigger".into(),
kind: "write".into(),
reply: Some("Ack".into()),
request: None,
encoding: None,
fanout: fanout.map(str::to_string),
idempotent: Some(false),
since: None,
description: None,
}],
blob: vec![],
media: vec![],
deprecated: vec![],
}])
}
#[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::session::open(&[], &[], false).await.unwrap();
let slices = slice_with_proc(Some("forbidden"));
let err = call(
&session,
"",
&CallTarget::Fleet,
"netring",
"capture/trigger",
&[],
None,
None,
Duration::from_millis(100),
Some(&slices),
)
.await
.unwrap_err()
.to_string();
assert!(err.contains("fanout"), "{err}");
assert!(err.contains("RFC 05 §2.1"), "{err}");
let report = call(
&session,
"",
&CallTarget::Fleet,
"netring",
"capture/trigger",
&[],
None,
None,
Duration::from_millis(100),
Some(&slice_with_proc(None)),
)
.await
.unwrap();
assert_eq!(report.exit_code(), 2, "silence stays exit 2");
}
}