use std::time::Duration;
use crate::{Error, Result};
use zenkey::grammar::with_base;
use zenkey::{Declared, RegistrySlice};
use crate::bus::query::{Answer, GetOpts, RepeatingRegistry, fleet_get, state_snapshot};
use crate::judge::common::{FINDING_CAP, is_synthetic_marker};
use crate::model::examples::Examples;
use crate::report::{CheckId, DoctorFinding, DoctorReport, DoctorSeverity, DriftVerdict};
#[derive(Debug, Clone)]
pub struct DoctorSpec {
pub deep: bool,
pub sample: Option<usize>,
pub timeout: Duration,
pub listen: Option<Duration>,
}
fn finding(
severity: DoctorSeverity,
check: CheckId,
subject: impl Into<String>,
evidence: impl Into<String>,
citation: Option<&str>,
) -> DoctorFinding {
DoctorFinding {
severity,
check,
subject: subject.into(),
evidence: evidence.into(),
citation: citation.map(str::to_string),
}
}
fn rpc_key(base: &str, slice: &RegistrySlice, procedure: &str) -> Result<String> {
Ok(match &slice.service_origin {
Some(origin) => {
let o = origin.known().ok_or_else(|| {
Error::malformed(
format!("slice {}", slice.name),
format!("carries {:?} as a service origin", origin.token()),
)
})?;
with_base(base, zenkey::selector::service_rpc(o, &[procedure]))
}
None => with_base(base, zenkey::selector::fleet_rpc(&slice.name, &[procedure])),
})
}
pub async fn run_doctor(
fleet: &crate::Fleet<'_>,
locals: Option<&crate::model::registry::SliceSet>,
spec: &DoctorSpec,
) -> Result<DoctorReport> {
let (session, base) = (fleet.session(), fleet.base());
let locals = locals.filter(|set| !set.slices().is_empty());
let roster = crate::bus::roster::roster(fleet, spec.timeout).await?;
let mut findings: Vec<DoctorFinding> = Vec::new();
let mut synced: Vec<String> = Vec::new();
let mut answered = 0usize;
for local in locals.iter().flat_map(|set| set.slices()) {
let key = rpc_key(base, local, "introspect")?;
let answers = fleet_get(fleet, &key, &GetOpts::new(spec.timeout)).await?;
for answer in &answers {
let Answer::Value(bytes) = &answer.answer else {
continue;
};
answered += 1;
let served_toml = bytes.to_bytes();
let served_toml = String::from_utf8_lossy(&served_toml);
let served = match zenkey::parse_slice(&served_toml) {
Ok(s) => s,
Err(e) => {
findings.push(finding(
DoctorSeverity::Error,
CheckId::SliceParse,
format!("{}/{}", answer.origin, local.name),
format!("served slice does not parse: {e}"),
Some("RFC 08 §6"),
));
continue;
}
};
let diff = zenkey::slice::diff(&served, local);
if diff.is_empty() {
synced.push(format!(
"{}/{} (registry {})",
answer.origin, local.name, served.version
));
} else {
for f in &diff {
findings.push(finding(
DoctorSeverity::Error,
CheckId::SliceSync,
format!("{}/{}", answer.origin, local.name),
f.summary(),
Some("RFC 08 §6"),
));
}
}
}
}
let sweep = if locals.is_none() {
let repeating = RepeatingRegistry::declare(fleet, spec.timeout).await?;
let slices: Vec<RegistrySlice> = repeating
.fetch()
.await?
.into_iter()
.map(|(s, _)| s)
.collect();
repeating.undeclare().await?;
Some(slices)
} else {
None
};
if let Some(slices) = &sweep {
answered = slices.len();
}
let live: usize = roster.values().map(Vec::len).sum();
findings.extend(judge_introspect_coverage(
&roster,
locals.map(crate::model::registry::SliceSet::slices),
answered,
));
let routers = crate::routers(session, spec.timeout)
.await
.unwrap_or_default();
let mut router_version = None;
if routers.is_empty() {
findings.push(finding(
DoctorSeverity::Info,
CheckId::AdminUnreachable,
"mesh",
"no routers answered @/*/router (peer-only mesh, or the admin space is \
disabled) — storage/version checks skipped",
None,
));
} else {
let versions: std::collections::BTreeSet<&str> = routers
.iter()
.filter_map(|r| r.version.as_deref())
.collect();
if versions.len() > 1 {
findings.push(finding(
DoctorSeverity::Error,
CheckId::RouterVersionSkew,
"mesh",
format!("router version skew across the mesh: {versions:?}"),
None,
));
} else {
router_version = versions.iter().next().map(|v| v.to_string());
}
}
let slice_set: std::borrow::Cow<'_, crate::model::registry::SliceSet> = match sweep {
Some(slices) => {
std::borrow::Cow::Owned(crate::model::registry::SliceSet::from_slices(slices))
}
None => match locals {
Some(set) => std::borrow::Cow::Borrowed(set),
None => std::borrow::Cow::Owned(crate::model::registry::SliceSet::default()),
},
};
let mut described: Vec<(String, zenkey::schema::SchemaSet)> = Vec::new();
let mut described_by_origin: Vec<crate::model::decode::DescribedSchema> = Vec::new();
let mut undescribed = 0usize;
let describe_opts = GetOpts::new(spec.timeout);
for slice in slice_set.slices() {
let key = rpc_key(base, slice, "describe")?;
let answers = fleet_get(fleet, &key, &describe_opts).await?;
let before = described_by_origin.len();
for a in answers {
let origin = a.origin;
let Answer::Value(bytes) = a.answer else {
continue;
};
let cow = bytes.to_bytes();
let Some(set) = std::str::from_utf8(&cow)
.ok()
.and_then(|t| zenkey::schema::SchemaSet::parse(t).ok())
else {
continue;
};
if described_by_origin.len() == before {
described.push((slice.name.clone(), set.clone()));
}
described_by_origin.push(crate::model::decode::DescribedSchema {
origin,
producer: slice.name.clone(),
set,
});
}
if described_by_origin.len() == before {
undescribed += 1;
}
}
for gap in crate::model::decode::totality_gaps(&described, &slice_set) {
findings.push(finding(
DoctorSeverity::Error,
CheckId::DescribeTotality,
gap.producer.clone(),
format!(
"describe is not total — missing: {}",
gap.missing.join(", ")
),
Some("RFC 08 §7"),
));
}
for drift in crate::model::decode::schema_drift(&described_by_origin) {
let servers: Vec<String> = drift
.servers
.iter()
.map(|s| match s.hash.as_option() {
Some(h) => format!("{}@{} ({h})", s.producer, s.origin),
None => format!("{}@{} (no identity served)", s.producer, s.origin),
})
.collect();
let (severity, evidence) = match drift.verdict {
DriftVerdict::Disagree => (
DoctorSeverity::Error,
format!("served with different schemas by {}", servers.join(", ")),
),
DriftVerdict::Unjudgeable => (
DoctorSeverity::Warning,
format!(
"agreement cannot be judged — {} served no schema identity: {} \
(RFC 09 §5.1 O4; the hash exists for exactly this, RFC 08 §7)",
drift
.servers
.iter()
.filter(|s| s.hash.is_not_asked())
.count(),
servers.join(", ")
),
),
};
findings.push(finding(
severity,
CheckId::SchemaDrift,
drift.type_name.clone(),
evidence,
Some("RFC 08 §7"),
));
}
if undescribed > 0 {
findings.push(finding(
DoctorSeverity::Info,
CheckId::DescribeMissing,
"fleet",
format!(
"{undescribed} producer(s) serve no describe (a SHOULD; generic tools \
render their payloads structurally)"
),
Some("RFC 08 §7"),
));
}
if spec.deep {
let now = std::time::SystemTime::now();
let mut unstamped = 0usize;
for slice in slice_set.slices() {
for subject in &slice.subjects {
let (Some(ttl), true) = (subject.ttl_s, subject.class.is(&zenkey::Class::State))
else {
continue;
};
let Ok(pattern) = zenkey::pattern::SubjectPattern::parse(&subject.path) else {
continue;
};
let selector = match &slice.service_origin {
Some(origin) => with_base(
base,
format!("v1/{origin}/state/{}", pattern.selector_tail()),
),
None => with_base(
base,
format!("v1/*/state/{}/{}", slice.name, pattern.selector_tail()),
),
};
let samples = state_snapshot(session, &selector, spec.timeout, spec.sample).await?;
let (family_findings, family_unstamped) = judge_state_samples(&samples, ttl, now);
findings.extend(family_findings);
unstamped += family_unstamped;
}
}
if unstamped > 0 {
findings.push(finding(
DoctorSeverity::Warning,
CheckId::UnstampedState,
"fleet",
format!(
"{unstamped} state sample(s) carry no HLC timestamp — the deployment \
lacks timestamping, which LWW requires; freshness is unjudgeable \
for them"
),
Some("RFC 04 §4"),
));
}
let storages = crate::storages(session, spec.timeout)
.await
.unwrap_or_default();
let coverage = crate::state_coverage(&slice_set, base, &storages);
let uncovered: Vec<&crate::CoverageRow> = coverage
.iter()
.filter(|r| r.coverage == crate::Coverage::Uncovered)
.collect();
if !uncovered.is_empty() {
findings.push(finding(
DoctorSeverity::Info,
CheckId::StorageCoverage,
"fleet",
format!(
"{} state famil(y|ies) have no storage coverage (volatile seeding \
may ride the advanced-pub/sub cache): {}",
uncovered.len(),
uncovered
.iter()
.map(|r| format!("{}/{}", r.producer, r.path))
.collect::<Vec<_>>()
.join(", ")
),
Some("RFC 04 §3.5"),
));
}
}
let observation = match spec.listen {
Some(window) => {
let store = crate::model::decode::SchemaStore::new(base, spec.timeout);
for (producer, set) in &described {
store.insert(producer, set.clone());
}
let _sealed = store.seal();
let (listen_findings, summary) =
observe_traffic(fleet, &slice_set, &store, &described, window).await?;
findings.extend(listen_findings);
Some(summary)
}
None => None,
};
Ok(DoctorReport {
findings,
synced: locals.is_some().then_some(synced).into(),
introspect_answered: answered,
live_producers: live,
describe_served: described.len(),
describe_missing: undescribed,
routers: routers.len(),
router_version,
deep: spec.deep,
observation,
})
}
const DECODE_BUDGET: u8 = 2;
const SAME_FINDING: &str = "more key(s) with the same finding";
fn emit_capped(
findings: &mut Vec<DoctorFinding>,
ex: Examples<DoctorFinding>,
check: CheckId,
tail: &str,
) {
let more = ex.more(tail);
findings.extend(ex.into_vec());
if let Some(evidence) = more {
findings.push(finding(
DoctorSeverity::Info,
check,
"fleet",
evidence,
None,
));
}
}
struct RateWindow {
cap: Option<u64>,
seen: u64,
}
async fn observe_traffic(
fleet: &crate::Fleet<'_>,
slices: &crate::model::registry::SliceSet,
store: &crate::model::decode::SchemaStore,
described: &[(String, zenkey::schema::SchemaSet)],
window: Duration,
) -> Result<(Vec<DoctorFinding>, crate::report::ObservationSummary)> {
use std::collections::BTreeMap;
let (session, base) = (fleet.session(), fleet.base());
let scopes = crate::judge::common::data_plane_scopes(base, slices);
let monitor = crate::Monitor::start(session, crate::MonitorSpec::default()).await?;
let mut events = monitor.events();
let monitor = monitor.watching(&scopes).await?;
let started = tokio::time::Instant::now();
let deadline = started + window;
let mut fields =
crate::judge::field::FieldObservation::new(crate::judge::field::DEFAULT_MAX_PATHS);
let mut samples: u64 = 0;
let mut dropped: u64 = 0;
let mut synthetic: u64 = 0;
let mut facts_cache = crate::model::facts::FactsCache::default();
let mut decode_budget: BTreeMap<String, u8> = BTreeMap::new();
let mut unregistered: BTreeMap<String, u64> = BTreeMap::new();
let mut qos_bad: BTreeMap<String, (String, u64, u64)> = BTreeMap::new();
let mut undecodable: BTreeMap<String, (String, u64)> = BTreeMap::new();
let mut invalid: BTreeMap<String, (String, u64)> = BTreeMap::new();
let mut event_counts: BTreeMap<(String, String), RateWindow> = BTreeMap::new();
let mut foreign_stampers: BTreeMap<String, u64> = BTreeMap::new();
let window_over = tokio::time::sleep_until(deadline);
tokio::pin!(window_over);
loop {
let item = tokio::select! {
item = events.recv() => item,
() = &mut window_over => break,
};
match item {
Some(crate::StreamItem::Event(crate::FleetEvent::Sample(s))) => {
samples += 1;
if let Some(att) = &s.attachment
&& is_synthetic_marker(&att.to_bytes())
{
synthetic += 1;
}
if let Some(crate::StampProvenance::Foreign { stamper }) = s.stamped_by {
*foreign_stampers.entry(stamper.to_string()).or_default() += 1;
}
let is_put = s.kind == zenoh::sample::SampleKind::Put;
let bytes = s.payload.to_bytes();
if is_put {
if bytes.len() > crate::model::decode::OBSERVE_LIMIT {
fields.observe_unread(&s.key);
} else {
let doc = crate::model::decode::structural_value(&bytes);
fields.observe(&s.key, started.elapsed().as_secs_f64(), doc.as_ref());
}
}
facts_cache.ensure(base, &s.key, Some(slices));
let facts = facts_cache.get(&s.key).expect("just ensured this key");
match &facts.registration {
crate::model::facts::Registration::Unregistered => {
*unregistered.entry(s.key.clone()).or_default() += 1;
}
crate::model::facts::Registration::Registered(sf) => {
if let (Some(profile), Some(declared)) = (sf.declared_qos(), &sf.qos) {
let entry = qos_bad
.entry(s.key.clone())
.or_insert_with(|| (declared.token().to_string(), 0, 0));
entry.2 += 1;
if !s.qos_matches(profile) {
entry.1 += 1;
}
}
if let (Some(rate), crate::model::facts::KeyShape::V1(v)) =
(&sf.rate, &facts.shape)
&& v.class == "events"
{
let family = match &v.producer {
Some(p) => format!("{p}/{}", sf.path),
None => format!("{}/{}", v.origin, sf.path),
};
event_counts
.entry((family, rate.token()))
.or_insert_with(|| RateWindow {
cap: rate.cap_per_hour(),
seen: 0,
})
.seen += 1;
}
let budget = decode_budget.entry(s.key.clone()).or_default();
if is_put && *budget < DECODE_BUDGET {
*budget += 1;
let d = crate::model::decode::decode_sample(
fleet,
store,
Some(slices),
&s.key,
Some(&s.encoding),
&s.payload.to_bytes(),
)
.await;
match d.verdict {
crate::Verdict::NotValidated(
zenkey::schema::validate::NotValidated::Undecodable,
) => {
let e = undecodable.entry(s.key.clone()).or_insert_with(|| {
(
d.decode_error
.unwrap_or_else(|| "does not decode".into()),
0,
)
});
e.1 += 1;
}
crate::Verdict::Invalid(errors) => {
let e = invalid
.entry(s.key.clone())
.or_insert_with(|| (errors.join("; "), 0));
e.1 += 1;
}
_ => {}
}
}
}
_ => {}
}
}
Some(crate::StreamItem::Dropped(n)) => dropped += n,
Some(_) => continue,
None => break,
}
}
let keys_seen = facts_cache.len();
monitor.shutdown().await?;
let budgets =
crate::judge::budget::BudgetObservation::observe(base, slices, facts_cache.keys());
let window_s = window.as_secs_f64();
let mut findings = Vec::new();
let mut ex = Examples::new(FINDING_CAP);
for (key, (error, n)) in &undecodable {
ex.push_with(|| {
finding(
DoctorSeverity::Error,
CheckId::PayloadUndecodable,
key.clone(),
format!(
"payload does not decode as its declared type: {error} ({n} sample(s) tried)"
),
Some("RFC 08 §7"),
)
});
}
emit_capped(&mut findings, ex, CheckId::PayloadUndecodable, SAME_FINDING);
let mut ex = Examples::new(FINDING_CAP);
for (key, (violations, n)) in &invalid {
ex.push_with(|| {
finding(
DoctorSeverity::Error,
CheckId::PayloadInvalid,
key.clone(),
format!("payload violates the served schema: {violations} ({n} sample(s) tried)"),
Some("RFC 08 §7"),
)
});
}
emit_capped(&mut findings, ex, CheckId::PayloadInvalid, SAME_FINDING);
findings.extend(judge_qos_observed(&qos_bad));
if !foreign_stampers.is_empty() {
let mut named: Vec<String> = foreign_stampers
.iter()
.map(|(zid, n)| format!("{zid} ({n} sample(s))"))
.collect();
named.sort();
findings.push(finding(
DoctorSeverity::Info,
CheckId::TimestampStampedElsewhere,
"fleet".to_string(),
format!(
"HLCs on this bus are stamped by {} node(s) that are not the publishing \
session — a deployment with router-side timestamping, which is legal and \
common. Latency measured from these stamps is stamper→observer, not \
publisher→observer: {}",
foreign_stampers.len(),
named.join(", ")
),
Some("RFC 09 §5.1 O7"),
));
}
let mut ex = Examples::new(FINDING_CAP);
for (key, n) in &unregistered {
ex.push_with(|| {
finding(
DoctorSeverity::Warning,
CheckId::UnregisteredTraffic,
key.clone(),
format!(
"{n} sample(s) on a subject the producer's slice does not declare — \
for a conforming producer, a subject that is not registered does not exist"
),
Some("RFC 08 §2"),
)
});
}
emit_capped(
&mut findings,
ex,
CheckId::UnregisteredTraffic,
SAME_FINDING,
);
if window <= Duration::from_secs(3600) {
for ((family, rate), RateWindow { cap, seen: count }) in &event_counts {
let Some(cap) = cap else {
continue;
};
if count > cap {
findings.push(finding(
DoctorSeverity::Warning,
CheckId::RateOverDeclared,
family.clone(),
format!(
"{count} event(s) in {window_s:.0}s exceeds the declared \
`{rate}` cap ({cap}/h)"
),
Some("RFC 04 §1.3"),
));
}
}
}
findings.extend(judge_cardinality(slices, &budgets, window_s));
let field_ctx = field_context_from(slices, described, &facts_cache);
findings.extend(crate::judge::field::judge_fields(
&fields, window_s, &field_ctx,
));
Ok((
findings,
crate::report::ObservationSummary {
window_s,
scopes,
samples,
keys_seen,
dropped,
synthetic_marked: synthetic,
field_paths_dropped: fields.dropped_paths(),
facts_evicted: facts_cache.evicted(),
},
))
}
fn field_context_from(
slices: &crate::model::registry::SliceSet,
described: &[(String, zenkey::schema::SchemaSet)],
facts: &crate::model::facts::FactsCache,
) -> std::collections::BTreeMap<String, crate::judge::field::KeyFieldContext> {
use std::collections::BTreeMap;
let mut declared_cache: BTreeMap<(String, String), Option<crate::judge::field::DeclaredPaths>> =
BTreeMap::new();
let mut ctx = BTreeMap::new();
for (key, f) in facts.iter() {
let mut c = crate::judge::field::KeyFieldContext::default();
if let crate::model::facts::Registration::Registered(sf) = &f.registration {
c.ttl_s = sf.ttl_s;
c.type_name = Some(sf.type_name.clone());
if let Some(producer) = crate::judge::common::producer_of(f, Some(slices))
&& !sf.type_name.is_empty()
{
let declared = declared_cache
.entry((producer.clone(), sf.type_name.clone()))
.or_insert_with(|| {
described
.iter()
.find(|(name, _)| *name == producer)
.and_then(|(_, set)| set.get(&sf.type_name))
.and_then(|schema| schema.json_document())
.and_then(crate::judge::field::DeclaredPaths::from_json_schema)
});
c.declared = declared.clone();
}
}
ctx.insert(key.to_string(), c);
}
ctx
}
fn judge_qos_observed(
qos_bad: &std::collections::BTreeMap<String, (String, u64, u64)>,
) -> Vec<DoctorFinding> {
let mut findings = Vec::new();
let mut ex = Examples::new(FINDING_CAP);
for (key, (declared, bad, total)) in qos_bad.iter().filter(|(_, (_, bad, _))| *bad > 0) {
ex.push_with(|| {
finding(
DoctorSeverity::Warning,
CheckId::QosObservedMismatch,
key.clone(),
format!(
"{bad} of {total} sample(s) did not ride the declared {declared} — this \
is what actually rode: an interceptor MAY rewrite QoS, so it is a \
deviation, not proof of the publisher"
),
Some("RFC 04 §3"),
)
});
}
emit_capped(
&mut findings,
ex,
CheckId::QosObservedMismatch,
SAME_FINDING,
);
findings
}
fn judge_introspect_coverage(
roster: &std::collections::BTreeMap<String, Vec<String>>,
locals: Option<&[RegistrySlice]>,
answered: usize,
) -> Option<DoctorFinding> {
let live: usize = roster.values().map(Vec::len).sum();
let (in_scope, scope) = match locals {
None => (
live,
"scope: the whole roster (fleet-wide wildcard sweep)".to_string(),
),
Some(locals) => {
let named = |origin: &str, producer: &str| {
let base_name = zenkey::grammar::Producer::parse_chunk(producer)
.map(|p| p.name().to_string())
.unwrap_or_else(|_| producer.to_string());
locals.iter().any(|l| {
l.name == base_name
|| l.service_origin.as_ref().map(Declared::token) == Some(origin)
})
};
let in_scope: usize = roster
.iter()
.map(|(origin, producers)| producers.iter().filter(|p| named(origin, p)).count())
.sum();
let mut names: Vec<&str> = locals.iter().map(|l| l.name.as_str()).collect();
names.sort_unstable();
names.dedup();
let not_asked = live - in_scope;
(
in_scope,
format!(
"scope: the producer(s) the local registry names ({}); {} other \
live producer(s) were not asked and are not counted (O4)",
names.join(", "),
not_asked
),
)
}
};
(answered < in_scope).then(|| {
finding(
DoctorSeverity::Error,
CheckId::IntrospectCoverage,
"fleet",
format!(
"{} of {} live producer(s) in scope did not answer introspect — \
alive ⇒ callable, so this is a finding, not a boot race; {scope}",
in_scope - answered,
in_scope
),
Some("RFC 04 §5"),
)
})
}
fn judge_state_samples(
samples: &[crate::StateSample],
ttl: i64,
now: std::time::SystemTime,
) -> (Vec<DoctorFinding>, usize) {
let mut findings = Vec::new();
let mut unstamped = 0usize;
for sample in samples {
match sample.timestamp {
Some(ts) => {
let stamped = ts.get_time().to_system_time();
if let Ok(age) = now.duration_since(stamped)
&& age.as_secs() as i64 > ttl
{
findings.push(finding(
DoctorSeverity::Error,
CheckId::StaleState,
sample.key.clone(),
format!(
"{}s old against ttl {ttl}s (refresh <= ttl/2)",
age.as_secs()
),
Some("RFC 04 §1.2"),
));
}
}
None => unstamped += 1,
}
}
(findings, unstamped)
}
fn judge_cardinality(
slices: &crate::model::registry::SliceSet,
observed: &crate::judge::budget::BudgetObservation,
window_s: f64,
) -> Vec<DoctorFinding> {
let mut findings = Vec::new();
let mut over: Examples<DoctorFinding> = Examples::new(FINDING_CAP);
for slice in slices.slices() {
for s in &slice.subjects {
if !s.path.contains('{') {
continue; }
if s.path.contains("...") {
let seen: usize = observed
.family(&slice.name, &s.path)
.map(|origins| origins.values().map(|keys| keys.len()).sum())
.unwrap_or(0);
findings.push(finding(
DoctorSeverity::Info,
CheckId::CardinalityOverDeclared,
format!("{}/{}", slice.name, s.path),
format!(
"exempt: rest-variable — a `{{var...}}` family is unbounded by \
construction, so its declared cardinality is not a bound this \
check can pass or fail; {seen} distinct key(s) observed in \
{window_s:.0}s"
),
Some("RFC 08 §6.1"),
));
continue;
}
let Some(declared) = s.cardinality else {
continue; };
let Some(origins) = observed.family(&slice.name, &s.path) else {
continue; };
for (origin, keys) in origins {
if keys.len() as i64 <= declared {
continue; }
let examples = Examples::collect(
crate::judge::common::EXPANSION_CAP,
keys.iter().map(String::as_str),
);
let subject = if origin.starts_with('@') {
format!("{origin}/{}", s.path)
} else {
format!("{origin}/{}/{}", slice.name, s.path)
};
over.push_with(|| {
finding(
DoctorSeverity::Warning,
CheckId::CardinalityOverDeclared,
subject,
format!(
"{} distinct key(s) observed in {window_s:.0}s exceed the \
declared cardinality {declared} — e.g. {}. A bounded window \
observes a lower bound: the population is at least this",
keys.len(),
examples.as_slice().join(", ")
),
Some("RFC 04 §1.2"),
)
});
}
}
}
emit_capped(
&mut findings,
over,
CheckId::CardinalityOverDeclared,
"more origin famil(y|ies) over their declared cardinality",
);
findings
}
#[cfg(test)]
mod tests {
use super::*;
const BOUNDED: &str = r#"
[registry]
version = "1.0"
app = "t"
convention = 1
[producer]
name = "sysinfo"
[[subject]]
path = "disk/{mount}/used"
class = "telemetry"
type = "Point"
cardinality = 16
"#;
#[test]
fn cardinality_over_declared_fires_with_count_and_examples() {
let slices = crate::model::registry::SliceSet::from_toml_for_tests(BOUNDED);
let keys: Vec<String> = (0..40)
.map(|i| format!("v1/h-aaaaaaaaaaaa/telemetry/sysinfo/disk/m{i:02}/used"))
.collect();
let obs = crate::judge::budget::BudgetObservation::observe(
"",
&slices,
keys.iter().map(String::as_str),
);
let findings = judge_cardinality(&slices, &obs, 10.0);
assert_eq!(findings.len(), 1, "{findings:?}");
let f = &findings[0];
assert_eq!(f.check, CheckId::CardinalityOverDeclared);
assert_eq!(f.severity, DoctorSeverity::Warning);
assert_eq!(f.subject, "h-aaaaaaaaaaaa/sysinfo/disk/{mount}/used");
assert!(f.evidence.contains("40 distinct key(s)"), "{}", f.evidence);
assert!(f.evidence.contains("cardinality 16"), "{}", f.evidence);
assert!(f.evidence.contains("10s"), "the window is stated");
assert!(
f.evidence.contains("disk/m00/used"),
"examples are named: {}",
f.evidence
);
}
#[test]
fn cardinality_under_declared_is_not_a_finding() {
let slices = crate::model::registry::SliceSet::from_toml_for_tests(BOUNDED);
let keys = [
"v1/h-aaaaaaaaaaaa/telemetry/sysinfo/disk/root/used",
"v1/h-aaaaaaaaaaaa/telemetry/sysinfo/disk/var/used",
"v1/h-bbbbbbbbbbbb/telemetry/sysinfo/disk/root/used",
];
let obs = crate::judge::budget::BudgetObservation::observe("", &slices, keys);
assert!(judge_cardinality(&slices, &obs, 5.0).is_empty());
}
#[test]
fn rest_variable_families_are_exempt_and_say_so() {
let toml = r#"
[registry]
version = "1.0"
app = "t"
convention = 1
[producer]
name = "gnmi"
[[subject]]
path = "{device}/{path...}"
class = "telemetry"
type = "Point"
cardinality = 2
"#;
let slices = crate::model::registry::SliceSet::from_toml_for_tests(toml);
let keys: Vec<String> = (0..5)
.map(|i| format!("v1/h-aaaaaaaaaaaa/telemetry/gnmi/sw1/if/eth{i}/rx"))
.collect();
let obs = crate::judge::budget::BudgetObservation::observe(
"",
&slices,
keys.iter().map(String::as_str),
);
let findings = judge_cardinality(&slices, &obs, 5.0);
assert_eq!(findings.len(), 1, "{findings:?}");
let f = &findings[0];
assert_eq!(f.severity, DoctorSeverity::Info, "an exemption, not a pass");
assert!(
f.evidence.starts_with("exempt: rest-variable"),
"{}",
f.evidence
);
assert!(f.evidence.contains("5 distinct key(s)"), "{}", f.evidence);
let quiet = judge_cardinality(&slices, &Default::default(), 5.0);
assert_eq!(quiet.len(), 1);
assert!(quiet[0].evidence.starts_with("exempt: rest-variable"));
}
#[test]
fn qos_mismatch_cap_bounds_violators_not_map_entries() {
let mut qos_bad: std::collections::BTreeMap<String, (String, u64, u64)> =
std::collections::BTreeMap::new();
for i in 0..30u32 {
let bad = if i < 5 { 0 } else { 1 };
qos_bad.insert(
format!("v1/h-a/telemetry/x/k{i:02}"),
("tel".into(), bad, 10),
);
}
let findings = judge_qos_observed(&qos_bad);
let per_key: Vec<&DoctorFinding> = findings
.iter()
.filter(|f| f.severity == DoctorSeverity::Warning)
.collect();
assert_eq!(per_key.len(), FINDING_CAP, "the cap bounds the findings");
assert!(
per_key.iter().all(|f| f.evidence.starts_with("1 of 10")),
"only violators become findings: {findings:#?}"
);
assert!(
per_key.iter().any(|f| f.subject.ends_with("k24")),
"violators past the first {FINDING_CAP} map entries (k20..k24) are \
not dropped: {findings:#?}"
);
let note = findings
.iter()
.find(|f| f.severity == DoctorSeverity::Info)
.expect("a remainder note");
assert_eq!(
note.evidence, "… and 5 more key(s) with the same finding",
"the note counts violators (25 − 20), not map entries"
);
let few: std::collections::BTreeMap<String, (String, u64, u64)> = qos_bad
.iter()
.take(10)
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
let findings = judge_qos_observed(&few);
assert_eq!(findings.len(), 5, "{findings:#?}");
assert!(
findings
.iter()
.all(|f| f.severity == DoctorSeverity::Warning)
);
}
fn roster_of(entries: &[(&str, &[&str])]) -> std::collections::BTreeMap<String, Vec<String>> {
entries
.iter()
.map(|(origin, producers)| {
(
origin.to_string(),
producers.iter().map(|p| p.to_string()).collect(),
)
})
.collect()
}
fn slice_named(name: &str) -> RegistrySlice {
zenkey::parse_slice(&format!(
"[registry]\nversion = \"1.0\"\napp = \"t\"\nconvention = 1\n\
[producer]\nname = \"{name}\"\n"
))
.expect("fixture slice parses")
}
#[test]
fn a_partial_registry_does_not_count_unasked_producers_against_coverage() {
let roster = roster_of(&[("h-aaaaaaaaaaaa", &["sysinfo", "extra"])]);
let locals = [slice_named("sysinfo")];
assert_eq!(
judge_introspect_coverage(&roster, Some(&locals), 1),
None,
"the un-asked producer is out of scope, not silent"
);
}
#[test]
fn introspect_coverage_evidence_states_its_scope() {
let roster = roster_of(&[
("h-aaaaaaaaaaaa", &["sysinfo", "extra"]),
("h-bbbbbbbbbbbb", &["sysinfo-2"]),
]);
let locals = [slice_named("sysinfo")];
let f = judge_introspect_coverage(&roster, Some(&locals), 1).expect("a finding");
assert_eq!(f.check, CheckId::IntrospectCoverage);
assert!(f.evidence.contains("1 of 2"), "{}", f.evidence);
assert!(
f.evidence.contains("the local registry names (sysinfo)"),
"{}",
f.evidence
);
assert!(
f.evidence
.contains("1 other live producer(s) were not asked"),
"{}",
f.evidence
);
}
#[test]
fn a_service_slice_scopes_its_origin_into_coverage() {
let roster = roster_of(&[("@catalog", &["catalog"]), ("h-aaaaaaaaaaaa", &["extra"])]);
let locals = [zenkey::parse_slice(
"[registry]\nversion = \"1.0\"\napp = \"t\"\nconvention = 1\n\
[service]\nname = \"catalog\"\norigin = \"@catalog\"\n",
)
.expect("service slice parses")];
assert_eq!(judge_introspect_coverage(&roster, Some(&locals), 1), None);
let f = judge_introspect_coverage(&roster, Some(&locals), 0).expect("a finding");
assert!(f.evidence.contains("1 of 1"), "{}", f.evidence);
}
#[test]
fn the_wildcard_sweep_judges_the_whole_roster() {
let roster = roster_of(&[("h-aaaaaaaaaaaa", &["sysinfo", "extra"])]);
let f = judge_introspect_coverage(&roster, None, 1).expect("a finding");
assert!(f.evidence.contains("1 of 2"), "{}", f.evidence);
assert!(f.evidence.contains("whole roster"), "{}", f.evidence);
assert_eq!(judge_introspect_coverage(&roster, None, 2), None);
}
#[test]
fn freshness_judgement_is_pure_and_ttl_bound() {
let now = std::time::SystemTime::now();
let fresh_ts = zenoh::time::Timestamp::new(
zenoh::time::NTP64::from(now.duration_since(std::time::UNIX_EPOCH).unwrap()),
zenoh::time::TimestampId::rand(),
);
let stale_ts = zenoh::time::Timestamp::new(
zenoh::time::NTP64::from(
now.duration_since(std::time::UNIX_EPOCH).unwrap() - Duration::from_secs(120),
),
zenoh::time::TimestampId::rand(),
);
let samples = vec![
crate::StateSample {
key: "b/v1/h-aaaaaaaaaaaa/state/p/health".into(),
timestamp: Some(fresh_ts),
payload_len: 2,
},
crate::StateSample {
key: "b/v1/h-bbbbbbbbbbbb/state/p/health".into(),
timestamp: Some(stale_ts),
payload_len: 2,
},
crate::StateSample {
key: "b/v1/h-cccccccccccc/state/p/health".into(),
timestamp: None,
payload_len: 2,
},
];
let (findings, unstamped) = judge_state_samples(&samples, 30, now);
assert_eq!(
findings.len(),
1,
"only the stale stamped sample is a finding"
);
assert_eq!(findings[0].check, CheckId::StaleState);
assert!(findings[0].subject.contains("h-bbbbbbbbbbbb"));
assert_eq!(unstamped, 1, "the unstamped sample is counted, not judged");
}
}