use std::time::Duration;
use anyhow::{Result, anyhow};
use zenoh::Session;
use zenoh::query::{ConsolidationMode, QueryTarget};
#[derive(Debug, Clone)]
pub struct AdminEntry {
pub key: String,
pub value: serde_json::Value,
}
pub async fn admin_get(
session: &Session,
selector: &str,
timeout: Duration,
) -> Result<Vec<AdminEntry>> {
let replies = session
.get(selector)
.target(QueryTarget::All)
.consolidation(ConsolidationMode::None)
.timeout(timeout)
.await
.map_err(|e| anyhow!("admin get {selector}: {e}"))?;
let mut out = Vec::new();
while let Ok(reply) = replies.recv_async().await {
let Ok(sample) = reply.result() else { continue };
let bytes = sample.payload().to_bytes();
let value = serde_json::from_slice(&bytes).unwrap_or_else(|_| {
serde_json::Value::String(String::from_utf8_lossy(&bytes).to_string())
});
out.push(AdminEntry {
key: sample.key_expr().as_str().to_string(),
value,
});
}
out.sort_by(|a, b| a.key.cmp(&b.key));
Ok(out)
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct RouterInfo {
pub zid: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub version: Option<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub locators: Vec<String>,
pub raw: serde_json::Value,
}
pub async fn routers(session: &Session, timeout: Duration) -> Result<Vec<RouterInfo>> {
let entries = admin_get(session, "@/*/router", timeout).await?;
Ok(entries
.into_iter()
.map(|e| {
let zid = e
.value
.get("zid")
.and_then(|v| v.as_str())
.map(str::to_string)
.unwrap_or_else(|| {
e.key.split('/').nth(1).unwrap_or("?").to_string()
});
let version = e
.value
.get("version")
.and_then(|v| v.as_str())
.map(str::to_string);
let locators = e
.value
.get("locators")
.and_then(|v| v.as_array())
.map(|a| {
a.iter()
.filter_map(|l| l.as_str().map(str::to_string))
.collect()
})
.unwrap_or_default();
RouterInfo {
zid,
version,
locators,
raw: e.value,
}
})
.collect())
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct StorageInfo {
pub zid: String,
pub name: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub key_expr: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub strip_prefix: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub volume: Option<String>,
pub raw: serde_json::Value,
}
pub fn storage_from_admin_entry(key: &str, value: &serde_json::Value) -> Option<StorageInfo> {
let chunks: Vec<&str> = key.split('/').collect();
let storages_pos = chunks.iter().position(|c| *c == "storages")?;
if chunks.get(storages_pos.checked_sub(1)?) != Some(&"storage_manager") {
return None;
}
let name = chunks.get(storages_pos + 1)?;
let zid = chunks.get(1).unwrap_or(&"?");
let text = |field: &str| {
value
.get(field)
.and_then(|v| v.as_str())
.map(str::to_string)
};
let volume = text("volume").or_else(|| {
value
.get("volume")
.and_then(|v| v.get("id"))
.and_then(|v| v.as_str())
.map(str::to_string)
});
Some(StorageInfo {
zid: (*zid).to_string(),
name: (*name).to_string(),
key_expr: text("key_expr"),
strip_prefix: text("strip_prefix"),
volume,
raw: value.clone(),
})
}
pub fn merge_storage_rows(mut rows: Vec<StorageInfo>) -> Vec<StorageInfo> {
rows.sort_by(|a, b| (&a.zid, &a.name).cmp(&(&b.zid, &b.name)));
let mut out: Vec<StorageInfo> = Vec::with_capacity(rows.len());
for row in rows {
match out.last_mut() {
Some(prev) if prev.zid == row.zid && prev.name == row.name => {
prev.key_expr = prev.key_expr.take().or(row.key_expr);
prev.strip_prefix = prev.strip_prefix.take().or(row.strip_prefix);
prev.volume = prev.volume.take().or(row.volume);
if prev.raw.as_object().map(|o| o.len()).unwrap_or(0)
< row.raw.as_object().map(|o| o.len()).unwrap_or(0)
{
prev.raw = row.raw;
}
}
_ => out.push(row),
}
}
out
}
pub async fn storages(session: &Session, timeout: Duration) -> Result<Vec<StorageInfo>> {
let entries = admin_get(
session,
"@/*/router/**/storage_manager/storages/**",
timeout,
)
.await?;
let rows: Vec<StorageInfo> = entries
.iter()
.filter_map(|e| storage_from_admin_entry(&e.key, &e.value))
.collect();
Ok(merge_storage_rows(rows))
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
#[serde(tag = "coverage", content = "storage")]
pub enum Coverage {
Covered(String),
Partial(String),
Uncovered,
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct CoverageRow {
pub producer: String,
pub path: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub ttl_s: Option<i64>,
#[serde(flatten)]
pub coverage: Coverage,
}
pub fn state_coverage(
slices: &crate::registry::SliceSet,
base: &str,
storages: &[StorageInfo],
) -> Vec<CoverageRow> {
use zenoh::key_expr::keyexpr;
let storage_kes: Vec<(&StorageInfo, &keyexpr)> = storages
.iter()
.filter_map(|s| {
let ke = s.key_expr.as_deref()?;
keyexpr::new(ke).ok().map(|ke| (s, ke))
})
.collect();
let mut rows = Vec::new();
for slice in slices.slices() {
for subject in &slice.subjects {
if subject.class != "state" {
continue;
}
let Ok(pattern) = zenkey::pattern::SubjectPattern::parse(&subject.path) else {
continue;
};
let selector = match &slice.service_origin {
Some(origin) => zenkey::grammar::with_base(
base,
format!("v1/{origin}/state/{}", pattern.selector_tail()),
),
None => zenkey::grammar::with_base(
base,
format!("v1/*/state/{}/{}", slice.name, pattern.selector_tail()),
),
};
let Ok(family) = keyexpr::new(selector.as_str()) else {
continue;
};
let mut coverage = Coverage::Uncovered;
for (info, ke) in &storage_kes {
if ke.includes(family) {
coverage = Coverage::Covered(format!("{}@{}", info.name, info.zid));
break;
}
if ke.intersects(family) && coverage == Coverage::Uncovered {
coverage = Coverage::Partial(format!("{}@{}", info.name, info.zid));
}
}
rows.push(CoverageRow {
producer: slice.name.clone(),
path: subject.path.clone(),
ttl_s: subject.ttl_s,
coverage,
});
}
}
rows
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
#[serde(rename_all = "snake_case")]
pub enum EntityKind {
Subscriber,
Publisher,
Queryable,
Querier,
Token,
}
impl EntityKind {
fn from_chunk(chunk: &str) -> Option<EntityKind> {
Some(match chunk {
"subscriber" => EntityKind::Subscriber,
"publisher" => EntityKind::Publisher,
"queryable" => EntityKind::Queryable,
"querier" => EntityKind::Querier,
"token" => EntityKind::Token,
_ => return None,
})
}
fn chunk(self) -> &'static str {
match self {
EntityKind::Subscriber => "subscriber",
EntityKind::Publisher => "publisher",
EntityKind::Queryable => "queryable",
EntityKind::Querier => "querier",
EntityKind::Token => "token",
}
}
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct DeclaredEntity {
pub kind: EntityKind,
pub keyexpr: String,
pub node_zid: String,
pub sources: serde_json::Value,
}
#[derive(Debug, Clone, Default, serde::Serialize)]
pub struct DeclaredEntities {
pub entities: Vec<DeclaredEntity>,
}
pub fn declared_from_admin_entry(key: &str, value: &serde_json::Value) -> Option<DeclaredEntity> {
let mut chunks = key.split('/');
if chunks.next()? != "@" {
return None;
}
let zid = chunks.next()?;
let _whatami = chunks.next()?;
let kind = EntityKind::from_chunk(chunks.next()?)?;
let keyexpr: Vec<&str> = chunks.collect();
if keyexpr.is_empty() {
return None;
}
Some(DeclaredEntity {
kind,
keyexpr: keyexpr.join("/"),
node_zid: zid.to_string(),
sources: value.clone(),
})
}
pub async fn declared_entities(
session: &Session,
timeout: Duration,
) -> Result<Option<DeclaredEntities>> {
let mut entities = Vec::new();
let mut any_reply = false;
for kind in [
EntityKind::Subscriber,
EntityKind::Publisher,
EntityKind::Queryable,
EntityKind::Querier,
EntityKind::Token,
] {
let selector = format!("@/*/*/{}/**", kind.chunk());
let entries = admin_get(session, &selector, timeout).await?;
any_reply |= !entries.is_empty();
entities.extend(
entries
.iter()
.filter_map(|e| declared_from_admin_entry(&e.key, &e.value)),
);
}
if !any_reply {
return Ok(None);
}
Ok(Some(DeclaredEntities { entities }))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn storage_extraction_tolerates_layouts() {
let v = serde_json::json!({"key_expr": "zs/v1/*/state/**", "volume": "fs"});
let s = storage_from_admin_entry(
"@/abc123/router/config/plugins/storage_manager/storages/latest",
&v,
)
.unwrap();
assert_eq!(s.zid, "abc123");
assert_eq!(s.name, "latest");
assert_eq!(s.key_expr.as_deref(), Some("zs/v1/*/state/**"));
assert_eq!(s.volume.as_deref(), Some("fs"), "a bare-string volume");
let v = serde_json::json!({
"key_expr": "zs/v1/*/state/**",
"strip_prefix": "zs/v1",
"volume": {"id": "rocksdb", "dir": "latest"},
});
let s = storage_from_admin_entry(
"@/abc123/router/config/plugins/storage_manager/storages/durable",
&v,
)
.unwrap();
assert_eq!(s.volume.as_deref(), Some("rocksdb"), "an object volume");
assert_eq!(s.strip_prefix.as_deref(), Some("zs/v1"));
let s = storage_from_admin_entry(
"@/abc123/router/status/plugins/storage_manager/storages/latest/info",
&serde_json::json!("ok"),
)
.unwrap();
assert_eq!(s.name, "latest");
assert!(s.key_expr.is_none());
assert!(
storage_from_admin_entry(
"@/abc123/router/config/plugins/storage_manager/volumes/fs",
&serde_json::json!({}),
)
.is_none()
);
}
#[test]
fn merging_a_storage_keeps_every_field_either_side_named() {
let config = StorageInfo {
zid: "z1".into(),
name: "latest".into(),
key_expr: Some("zs/**".into()),
strip_prefix: Some("zs".into()),
volume: None,
raw: serde_json::json!({"key_expr": "zs/**", "strip_prefix": "zs"}),
};
let status = StorageInfo {
zid: "z1".into(),
name: "latest".into(),
key_expr: None,
strip_prefix: None,
volume: Some("fs".into()),
raw: serde_json::json!({"volume": "fs"}),
};
for rows in [
vec![config.clone(), status.clone()],
vec![status, config.clone()],
] {
let merged = merge_storage_rows(rows);
assert_eq!(merged.len(), 1, "one row per (zid, name)");
assert_eq!(merged[0].key_expr.as_deref(), Some("zs/**"));
assert_eq!(merged[0].strip_prefix.as_deref(), Some("zs"));
assert_eq!(merged[0].volume.as_deref(), Some("fs"));
}
let other = StorageInfo {
zid: "z1".into(),
name: "history".into(),
key_expr: None,
strip_prefix: None,
volume: None,
raw: serde_json::Value::Null,
};
assert_eq!(merge_storage_rows(vec![config, other]).len(), 2);
}
fn slices_with_state() -> crate::registry::SliceSet {
let toml = r#"
[registry]
version = "1.0"
app = "t"
convention = 1
[producer]
name = "tc"
[[subject]]
path = "health"
class = "state"
type = "Health"
ttl_s = 60
[[subject]]
path = "config/{iface}"
class = "state"
type = "Config"
ttl_s = 120
[[subject]]
path = "bandwidth"
class = "telemetry"
type = "Point"
"#;
crate::registry::SliceSet::from_toml_for_tests(toml)
}
fn storage(name: &str, key_expr: &str) -> StorageInfo {
StorageInfo {
zid: "z1".into(),
name: name.into(),
key_expr: Some(key_expr.into()),
strip_prefix: None,
volume: None,
raw: serde_json::Value::Null,
}
}
#[test]
fn coverage_judges_covered_partial_uncovered() {
let slices = slices_with_state();
let rows = state_coverage(&slices, "zs", &[storage("latest", "zs/v1/*/state/**")]);
assert_eq!(rows.len(), 2);
assert!(
rows.iter()
.all(|r| matches!(r.coverage, Coverage::Covered(_)))
);
let rows = state_coverage(
&slices,
"zs",
&[storage("one", "zs/v1/*/state/tc/config/eth0")],
);
let health = rows.iter().find(|r| r.path == "health").unwrap();
assert_eq!(health.coverage, Coverage::Uncovered);
let config = rows.iter().find(|r| r.path == "config/{iface}").unwrap();
assert!(matches!(config.coverage, Coverage::Partial(_)));
let rows = state_coverage(&slices, "zs", &[]);
assert!(rows.iter().all(|r| r.coverage == Coverage::Uncovered));
let rows = state_coverage(&slices, "", &[storage("latest", "v1/*/state/**")]);
assert_eq!(rows.len(), 2);
assert!(
rows.iter()
.all(|r| matches!(r.coverage, Coverage::Covered(_)))
);
}
#[test]
fn declared_entities_parse_the_admin_key_shape() {
let v = serde_json::json!({"routers": [], "peers": ["p1"], "clients": []});
let e = declared_from_admin_entry("@/a1b2c3/router/subscriber/zensight/v1/*/state/**", &v)
.unwrap();
assert_eq!(e.kind, EntityKind::Subscriber);
assert_eq!(e.keyexpr, "zensight/v1/*/state/**");
assert_eq!(e.node_zid, "a1b2c3");
let e = declared_from_admin_entry("@/z/peer/publisher/v1/h-a/telemetry/x/m", &v).unwrap();
assert_eq!(e.kind, EntityKind::Publisher);
assert_eq!(e.keyexpr, "v1/h-a/telemetry/x/m");
assert!(declared_from_admin_entry("@/z/router/config/x", &v).is_none());
assert!(declared_from_admin_entry("@/z/router/subscriber", &v).is_none());
assert!(declared_from_admin_entry("not/admin/at/all", &v).is_none());
}
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct OriginAttachment {
pub origin: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub session_zid: Option<String>,
pub reporter_zid: String,
pub token_key: String,
}
fn source_zids(sources: &serde_json::Value) -> Vec<String> {
let mut out = Vec::new();
for kind in ["routers", "peers", "clients"] {
if let Some(list) = sources.get(kind).and_then(|v| v.as_array()) {
out.extend(list.iter().filter_map(|z| z.as_str().map(str::to_string)));
}
}
out.sort();
out.dedup();
out
}
pub async fn origin_attachments(
session: &Session,
base: &str,
timeout: Duration,
) -> Result<Vec<OriginAttachment>> {
let entries = admin_get(session, "@/*/*/token/**", timeout).await?;
let mut out: Vec<OriginAttachment> = Vec::new();
for e in &entries {
let Some(decl) = declared_from_admin_entry(&e.key, &e.value) else {
continue;
};
if decl.kind != EntityKind::Token {
continue;
}
let Some(parsed) = zenkey::grammar::parse_full(base, &decl.keyexpr) else {
continue;
};
if parsed.subject.last().copied() != Some("alive") {
continue;
}
let origin = parsed.origin.chunk().to_string();
let zids = source_zids(&decl.sources);
let session_zid = match zids.as_slice() {
[only] => Some(only.clone()),
_ => None,
};
let attachment = OriginAttachment {
origin,
session_zid,
reporter_zid: decl.node_zid.clone(),
token_key: decl.keyexpr.clone(),
};
if !out.iter().any(|a| {
a.origin == attachment.origin
&& a.session_zid == attachment.session_zid
&& a.reporter_zid == attachment.reporter_zid
}) {
out.push(attachment);
}
}
Ok(out)
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct TopologyNode {
pub zid: String,
pub whatami: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub version: Option<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub locators: Vec<String>,
pub answered: bool,
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct TopologyEdge {
pub reporter: String,
pub peer: String,
pub whatami: String,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub links: Vec<String>,
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct TopologyReport {
pub nodes: Vec<TopologyNode>,
pub edges: Vec<TopologyEdge>,
pub asked: String,
pub answered: usize,
pub self_zid: String,
}
pub async fn topology(session: &Session, timeout: Duration) -> Result<TopologyReport> {
const ASKED: &str = "@/*/*";
let entries = admin_get(session, ASKED, timeout).await?;
let mut nodes: Vec<TopologyNode> = Vec::new();
let mut edges: Vec<TopologyEdge> = Vec::new();
for e in &entries {
let mut chunks = e.key.split('/');
let (Some("@"), Some(zid), Some(whatami), None) =
(chunks.next(), chunks.next(), chunks.next(), chunks.next())
else {
continue;
};
let doc = &e.value;
nodes.push(TopologyNode {
zid: doc
.get("zid")
.and_then(|v| v.as_str())
.unwrap_or(zid)
.to_string(),
whatami: whatami.to_string(),
version: doc
.get("version")
.and_then(|v| v.as_str())
.map(str::to_string),
locators: doc
.get("locators")
.and_then(|v| v.as_array())
.map(|a| {
a.iter()
.filter_map(|l| l.as_str().map(str::to_string))
.collect()
})
.unwrap_or_default(),
answered: true,
});
for s in doc
.get("sessions")
.and_then(|v| v.as_array())
.map(|a| a.as_slice())
.unwrap_or_default()
{
let Some(peer) = s.get("peer").and_then(|v| v.as_str()) else {
continue;
};
edges.push(TopologyEdge {
reporter: zid.to_string(),
peer: peer.to_string(),
whatami: s
.get("whatami")
.and_then(|v| v.as_str())
.unwrap_or("unknown")
.to_string(),
links: s
.get("links")
.and_then(|v| v.as_array())
.map(|a| {
a.iter()
.filter_map(|l| {
Some(format!(
"{} -> {}",
l.get("src")?.as_str()?,
l.get("dst")?.as_str()?
))
})
.collect()
})
.unwrap_or_default(),
});
}
}
let answered = nodes.len();
for e in &edges {
if !nodes.iter().any(|n| n.zid == e.peer) {
nodes.push(TopologyNode {
zid: e.peer.clone(),
whatami: e.whatami.clone(),
version: None,
locators: Vec::new(),
answered: false,
});
}
}
nodes.sort_by(|a, b| a.zid.cmp(&b.zid));
nodes.dedup_by(|a, b| a.zid == b.zid);
Ok(TopologyReport {
nodes,
edges,
asked: ASKED.to_string(),
answered,
self_zid: session.zid().to_string(),
})
}