use zenoh::key_expr::keyexpr;
use crate::report::{
Attribution, ConsumerRow, DeclaredEntities, DeprecationFact, EntityKind, OriginAttachment,
Relation, TopologyReport,
};
pub fn join_consumers(
target: &str,
declared: &DeclaredEntities,
attachments: &[OriginAttachment],
topology: Option<&TopologyReport>,
self_zid: &str,
) -> Vec<ConsumerRow> {
let Ok(target_ke) = keyexpr::new(target) else {
return Vec::new();
};
let base_tree = base_tree_of(target);
let base_ke = keyexpr::new(base_tree.as_str()).ok();
let mut rows: Vec<ConsumerRow> = Vec::new();
for entity in &declared.entities {
if !matches!(entity.kind, EntityKind::Subscriber | EntityKind::Querier) {
continue;
}
let Ok(declared_ke) = keyexpr::new(entity.keyexpr.as_str()) else {
continue;
};
let Some(relation) = relate(target_ke, declared_ke, base_ke) else {
continue;
};
let zids = crate::bus::admin::source_zids(&entity.sources);
let attributed: Vec<(String, Attribution)> = if zids.is_empty() {
vec![(entity.node_zid.clone(), Attribution::ReportedOnly)]
} else {
zids.into_iter()
.map(|z| (z, Attribution::Session))
.collect()
};
for (zid, attribution) in attributed {
if rows
.iter()
.any(|r| r.zid == zid && r.kind == entity.kind && r.keyexpr == entity.keyexpr)
{
continue;
}
let mut origins: Vec<String> = attachments
.iter()
.filter(|a| a.session_zid.as_deref() == Some(zid.as_str()))
.map(|a| a.origin.clone())
.collect();
origins.sort();
origins.dedup();
let whatami = topology
.and_then(|t| t.nodes.iter().find(|n| n.zid == zid))
.map(|n| n.whatami.clone());
rows.push(ConsumerRow {
is_self: zid == self_zid,
zid,
whatami,
origins,
attribution,
kind: entity.kind,
keyexpr: entity.keyexpr.clone(),
relation,
total_wildcard: relation == Relation::Total,
});
}
}
rows.sort_by(|a, b| {
a.relation
.cmp(&b.relation)
.then_with(|| a.keyexpr.cmp(&b.keyexpr))
.then_with(|| a.zid.cmp(&b.zid))
});
rows
}
fn relate(target: &keyexpr, declared: &keyexpr, base: Option<&keyexpr>) -> Option<Relation> {
if !declared.intersects(target) {
return None;
}
if base.is_some_and(|b| declared.includes(b)) {
return Some(Relation::Total);
}
let wider = declared.includes(target);
let narrower = target.includes(declared);
Some(match (wider, narrower) {
(true, true) => Relation::Exact,
(false, true) => Relation::Narrower,
(true, false) => Relation::Wider,
(false, false) => Relation::Intersects,
})
}
fn base_tree_of(target: &str) -> String {
let chunks: Vec<&str> = target.split('/').collect();
match chunks.iter().position(|c| *c == "v1") {
Some(i) => {
let mut prefix = chunks[..=i].join("/");
prefix.push_str("/**");
prefix
}
None => "**".to_string(),
}
}
pub fn declaring_sessions(target: &str, declared: &DeclaredEntities, kind: EntityKind) -> usize {
let Ok(target_ke) = keyexpr::new(target) else {
return 0;
};
let mut zids: Vec<String> = Vec::new();
for entity in declared.entities.iter().filter(|e| e.kind == kind) {
let Ok(declared_ke) = keyexpr::new(entity.keyexpr.as_str()) else {
continue;
};
if !declared_ke.intersects(target_ke) {
continue;
}
let sources = crate::bus::admin::source_zids(&entity.sources);
if sources.is_empty() {
zids.push(entity.node_zid.clone());
} else {
zids.extend(sources);
}
}
zids.sort();
zids.dedup();
zids.len()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubjectTarget {
pub class: String,
pub selector: String,
pub deprecated: Option<DeprecationFact>,
}
pub fn subject_target(
slices: &crate::SliceSet,
base: &str,
producer: &str,
path: &str,
) -> Option<SubjectTarget> {
let slice = slices.get(producer)?;
let deprecated = slice
.deprecated
.iter()
.find(|d| d.path == path && d.kind == zenkey::slice::DeprecatedKind::Subject)
.map(|d| DeprecationFact {
since: d.since.clone(),
replaced_by: d.replaced_by.clone(),
});
let subject = slice.subjects.iter().find(|s| s.path == path);
let class = match subject {
Some(s) => s.class.token().to_string(),
None => {
deprecated.as_ref()?;
"*".to_string()
}
};
let tail = zenkey::pattern::SubjectPattern::parse(path)
.ok()?
.selector_tail();
let relative = match &slice.service_origin {
Some(origin) => format!("v1/{origin}/{class}/{tail}"),
None => format!("v1/*/{class}/{}/{tail}", slice.name),
};
Some(SubjectTarget {
class,
selector: zenkey::grammar::with_base(base, relative),
deprecated,
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::report::{DeclaredEntity, TopologyNode};
use serde_json::json;
fn entity(
kind: EntityKind,
keyexpr: &str,
node: &str,
sources: serde_json::Value,
) -> DeclaredEntity {
DeclaredEntity {
kind,
keyexpr: keyexpr.into(),
node_zid: node.into(),
sources,
}
}
fn peers(zids: &[&str]) -> serde_json::Value {
json!({ "routers": [], "peers": zids, "clients": [] })
}
const TARGET: &str = "v1/h-cccccccccccc/state/demo/**";
#[test]
fn a_narrow_and_a_total_subscriber_rank_by_relation() {
let declared = DeclaredEntities {
entities: vec![
entity(
EntityKind::Subscriber,
"**",
"router1",
peers(&["zid-wide"]),
),
entity(
EntityKind::Subscriber,
"v1/h-cccccccccccc/state/demo/health",
"router1",
peers(&["zid-narrow"]),
),
entity(
EntityKind::Publisher,
"v1/h-cccccccccccc/state/demo/health",
"router1",
peers(&["zid-pub"]),
),
entity(
EntityKind::Subscriber,
"v1/h-dddddddddddd/state/demo/health",
"router1",
peers(&["zid-other"]),
),
],
};
let attachments = vec![OriginAttachment {
origin: "h-cccccccccccc".into(),
session_zid: Some("zid-narrow".into()),
reporter_zid: "router1".into(),
token_key: "v1/h-cccccccccccc/state/demo/alive".into(),
}];
let rows = join_consumers(TARGET, &declared, &attachments, None, "me");
assert_eq!(rows.len(), 2, "{rows:?}");
assert_eq!(rows[0].relation, Relation::Narrower);
assert_eq!(rows[0].zid, "zid-narrow");
assert_eq!(rows[0].origins, ["h-cccccccccccc"]);
assert_eq!(rows[0].attribution, Attribution::Session);
assert!(!rows[0].total_wildcard);
assert_eq!(rows[1].relation, Relation::Total);
assert!(rows[1].total_wildcard);
assert!(rows[1].origins.is_empty(), "session only, unattributed");
}
#[test]
fn every_relation_is_reachable() {
let declared = DeclaredEntities {
entities: vec![
entity(EntityKind::Subscriber, TARGET, "n", peers(&["exact"])),
entity(
EntityKind::Querier,
"v1/*/state/demo/**",
"n",
peers(&["wider"]),
),
entity(
EntityKind::Subscriber,
"v1/h-cccccccccccc/*/demo/health",
"n",
peers(&["intersects"]),
),
entity(EntityKind::Subscriber, "v1/**", "n", peers(&["total"])),
],
};
let rows = join_consumers(TARGET, &declared, &[], None, "me");
let got: Vec<(&str, Relation)> =
rows.iter().map(|r| (r.zid.as_str(), r.relation)).collect();
assert_eq!(
got,
[
("exact", Relation::Exact),
("wider", Relation::Wider),
("intersects", Relation::Intersects),
("total", Relation::Total),
]
);
assert_eq!(rows[1].kind, EntityKind::Querier);
}
#[test]
fn attribution_and_self_are_stated_not_guessed() {
let declared = DeclaredEntities {
entities: vec![
entity(EntityKind::Subscriber, TARGET, "reporter", json!({})),
entity(EntityKind::Subscriber, TARGET, "router1", peers(&["me"])),
entity(EntityKind::Subscriber, TARGET, "router2", peers(&["me"])),
],
};
let topology = TopologyReport {
nodes: vec![TopologyNode {
zid: "me".into(),
whatami: "peer".into(),
version: None,
locators: vec![],
locators_via_links: vec![],
answered: true,
}],
edges: vec![],
asked: "@/*/*".into(),
answered: 1,
self_zid: "me".into(),
};
let rows = join_consumers(TARGET, &declared, &[], Some(&topology), "me");
assert_eq!(rows.len(), 2, "{rows:?}");
let me = rows.iter().find(|r| r.zid == "me").unwrap();
assert!(me.is_self);
assert_eq!(me.whatami.as_deref(), Some("peer"));
assert_eq!(me.attribution, Attribution::Session);
let reported = rows.iter().find(|r| r.zid == "reporter").unwrap();
assert_eq!(reported.attribution, Attribution::ReportedOnly);
assert!(reported.whatami.is_none(), "not heard of is not a kind");
}
#[test]
fn a_total_wildcard_does_not_reach_a_verbatim_plane() {
let declared = DeclaredEntities {
entities: vec![entity(EntityKind::Subscriber, "**", "n", peers(&["wide"]))],
};
let rows = join_consumers("v1/@catalog/state/entity/**", &declared, &[], None, "me");
assert!(rows.is_empty(), "{rows:?}");
assert_eq!(base_tree_of("acme/v1/@catalog/state/**"), "acme/v1/**");
assert_eq!(base_tree_of("@/*/*"), "**");
}
#[test]
fn declaring_sessions_count_distinct_zids() {
let declared = DeclaredEntities {
entities: vec![
entity(EntityKind::Publisher, TARGET, "n", peers(&["a", "b"])),
entity(EntityKind::Publisher, TARGET, "n", peers(&["a"])),
entity(EntityKind::Publisher, TARGET, "reporter", json!({})),
entity(EntityKind::Queryable, "v1/**", "n", peers(&["q"])),
entity(EntityKind::Queryable, "other/**", "n", peers(&["r"])),
],
};
assert_eq!(
declaring_sessions(TARGET, &declared, EntityKind::Publisher),
3
);
assert_eq!(
declaring_sessions(TARGET, &declared, EntityKind::Queryable),
1
);
}
fn slices() -> crate::SliceSet {
let toml = r#"
[registry]
version = "1.0"
app = "acme"
convention = 1
[producer]
name = "sysinfo"
[[subject]]
path = "disk/{mount}/used"
class = "telemetry"
type = "TelemetryPoint"
[[subject]]
path = "health"
class = "state"
type = "HealthSnapshot"
[[deprecated]]
path = "ingest/legacy_total"
since = "2.0"
replaced_by = "disk/{mount}/bytes_used"
"#;
let slice = zenkey::parse_slice(toml).expect("a slice");
crate::SliceSet::from_slices(vec![slice])
}
#[test]
fn a_subject_resolves_to_its_family_selector() {
let set = slices();
let t = subject_target(&set, "acme", "sysinfo", "disk/{mount}/used").unwrap();
assert_eq!(t.class, "telemetry");
assert_eq!(t.selector, "acme/v1/*/telemetry/sysinfo/disk/*/used");
assert!(t.deprecated.is_none());
let t = subject_target(&set, "", "sysinfo", "health").unwrap();
assert_eq!(t.selector, "v1/*/state/sysinfo/health");
let t = subject_target(&set, "", "sysinfo", "ingest/legacy_total").unwrap();
assert_eq!(t.class, "*");
assert_eq!(t.selector, "v1/*/*/sysinfo/ingest/legacy_total");
assert_eq!(
t.deprecated,
Some(DeprecationFact {
since: Some("2.0".into()),
replaced_by: Some("disk/{mount}/bytes_used".into()),
})
);
assert!(subject_target(&set, "", "sysinfo", "nope").is_none());
assert!(subject_target(&set, "", "logs", "health").is_none());
}
}