use std::time::Duration;
use crate::Result;
use zenkey::slice::{DeprecationDecl, RegistrySlice};
use crate::judge::common::new_prefix;
use crate::report::{Asked, CutoverVerdict, RetiredEntry, RetiredReport};
pub fn scope_note(entries: usize, new_prefix: &str, window: Duration) -> String {
let window = window.as_secs_f64();
format!(
"retired check: {window}s window over {entries} ledger entr(y|ies) — \
watching the retired families and their replacements, with {new_prefix}** \
as the fleet's proof of life. `**` cannot cross `@`-chunks: verbatim \
planes and the admin space are outside this watch by construction (O5)."
)
}
pub fn retired_selector(slice: &RegistrySlice, path: &str) -> String {
let tail = zenkey::pattern::SubjectPattern::parse(path)
.map(|p| p.selector_tail())
.unwrap_or_else(|_| path.to_string());
match &slice.service_origin {
Some(origin) => format!("v1/{origin}/*/{tail}"),
None => format!("v1/*/*/{}/{tail}", slice.name),
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct EntryEvidence {
pub wire_samples: Asked<u64>,
pub still_declared: Asked<bool>,
pub subscribers: Asked<usize>,
pub life_samples: Asked<u64>,
}
pub fn entry_verdict(ev: EntryEvidence) -> CutoverVerdict {
let EntryEvidence {
wire_samples,
still_declared,
subscribers,
life_samples,
} = ev;
if wire_samples.as_option().is_some_and(|n| *n > 0)
|| still_declared.as_option() == Some(&true)
|| subscribers.as_option().is_some_and(|n| *n > 0)
{
return CutoverVerdict::OldStillSpeaks;
}
match (wire_samples.as_option(), life_samples.as_option()) {
(Some(0), Some(n)) if *n > 0 => CutoverVerdict::Pass,
_ => CutoverVerdict::Unproven,
}
}
pub fn overall(entries: &[RetiredEntry]) -> CutoverVerdict {
if entries
.iter()
.any(|e| e.verdict == CutoverVerdict::OldStillSpeaks)
{
CutoverVerdict::OldStillSpeaks
} else if entries
.iter()
.any(|e| e.verdict == CutoverVerdict::Unproven)
{
CutoverVerdict::Unproven
} else {
CutoverVerdict::Pass
}
}
enum Identity {
Host(String),
Service(String),
}
struct Matcher {
identity: Identity,
old: Option<zenkey::pattern::SubjectPattern>,
replacement: Option<zenkey::pattern::SubjectPattern>,
}
impl Matcher {
fn covers(&self, parsed: &zenkey::grammar::StructuralKey<'_>) -> bool {
match &self.identity {
Identity::Host(name) => parsed.producer().is_some_and(|p| p.name() == name.as_str()),
Identity::Service(origin) => {
parsed.producer().is_none() && parsed.origin.chunk() == origin.as_str()
}
}
}
}
pub async fn run_retired(
fleet: &crate::Fleet<'_>,
local: &crate::SliceSet,
registries: Vec<String>,
listen: Option<Duration>,
timeout: Duration,
) -> Result<RetiredReport> {
let (session, base) = (fleet.session(), fleet.base());
let mut ledger: Vec<(&RegistrySlice, &DeprecationDecl)> = local
.slices()
.iter()
.flat_map(|s| s.deprecated.iter().map(move |d| (s, d)))
.collect();
ledger.sort_by(|(sa, da), (sb, db)| {
(sa.name.as_str(), da.path.as_str()).cmp(&(sb.name.as_str(), db.path.as_str()))
});
let served = crate::SliceSet::from_bus(fleet, timeout).await?;
let admin = crate::bus::admin::declared_entities(session, timeout).await?;
let own_zid = session.zid().to_string();
let matchers: Vec<Matcher> = ledger
.iter()
.map(|(slice, decl)| Matcher {
identity: match &slice.service_origin {
Some(origin) => Identity::Service(origin.token().to_string()),
None => Identity::Host(slice.name.clone()),
},
old: zenkey::pattern::SubjectPattern::parse(&decl.path).ok(),
replacement: decl
.replaced_by
.as_deref()
.and_then(|p| zenkey::pattern::SubjectPattern::parse(p).ok()),
})
.collect();
let new_prefix = new_prefix(base);
let mut old_counts = vec![0u64; ledger.len()];
let mut repl_counts = vec![0u64; ledger.len()];
let (mut plane_samples, mut dropped) = (0u64, 0u64);
if let Some(window) = listen {
let monitor = crate::Monitor::start(session, crate::MonitorSpec::default()).await?;
let mut events = monitor.events();
let monitor = monitor.watching(["**"]).await?;
let deadline = tokio::time::Instant::now() + window;
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))) => {
if s.key.starts_with(&new_prefix) {
plane_samples += 1;
}
let Some(parsed) = zenkey::grammar::parse_full(base, &s.key) else {
continue;
};
if !matches!(parsed.class, zenkey::grammar::ClassOrPlane::Class(_)) {
continue;
}
for (i, m) in matchers.iter().enumerate() {
if !m.covers(&parsed) {
continue;
}
if m.old
.as_ref()
.is_some_and(|p| p.matches(&parsed.subject).is_some())
{
old_counts[i] += 1;
}
if m.replacement
.as_ref()
.is_some_and(|p| p.matches(&parsed.subject).is_some())
{
repl_counts[i] += 1;
}
}
}
Some(crate::StreamItem::Dropped(n)) => dropped += n,
Some(_) => continue,
None => break,
}
}
monitor.shutdown().await?;
}
let entries: Vec<RetiredEntry> = ledger
.iter()
.enumerate()
.map(|(i, (slice, decl))| {
let selector = retired_selector(slice, &decl.path);
let wire_samples = listen.map(|_| old_counts[i]);
let still_declared = served
.get(&slice.name)
.map(|served| served.serves_subject(&decl.path));
let subscribers = admin.as_ref().map(|entities| {
let family = zenkey::grammar::with_base(base, &selector);
let Ok(family) = zenoh::key_expr::KeyExpr::try_from(family) else {
return 0;
};
entities
.entities
.iter()
.filter(|e| e.kind == crate::EntityKind::Subscriber)
.filter(|e| e.node_zid != own_zid)
.filter(|e| {
zenoh::key_expr::KeyExpr::try_from(e.keyexpr.as_str())
.map(|k| k.intersects(&family))
.unwrap_or(false)
})
.count()
});
let replacement_samples = match (&decl.replaced_by, listen) {
(Some(_), Some(_)) => Some(repl_counts[i]),
_ => None,
};
let life = match &decl.replaced_by {
Some(_) => replacement_samples,
None => listen.map(|_| plane_samples),
};
RetiredEntry {
producer: slice.name.clone(),
path: decl.path.clone(),
since: decl.since.clone(),
replaced_by: decl.replaced_by.clone(),
selector,
wire_samples: wire_samples.into(),
still_declared,
subscribers,
replacement_samples: replacement_samples.into(),
verdict: entry_verdict(EntryEvidence {
wire_samples: wire_samples.into(),
still_declared: still_declared.into(),
subscribers: subscribers.into(),
life_samples: life.into(),
}),
}
})
.collect();
let verdict = overall(&entries);
Ok(RetiredReport {
registries,
entries,
window_s: listen.map(|d| d.as_secs_f64()).into(),
plane_samples: listen.map(|_| plane_samples).into(),
dropped: listen.map(|_| dropped).into(),
introspect_answered: served.slices().len(),
admin_entities: admin.as_ref().map(|e| e.entities.len()),
verdict,
})
}
#[cfg(test)]
mod tests {
use super::*;
fn verdict(
wire_samples: Option<u64>,
still_declared: Option<bool>,
subscribers: Option<usize>,
life_samples: Option<u64>,
) -> CutoverVerdict {
entry_verdict(EntryEvidence {
wire_samples: wire_samples.into(),
still_declared: still_declared.into(),
subscribers: subscribers.into(),
life_samples: life_samples.into(),
})
}
fn slice(toml: &str) -> RegistrySlice {
zenkey::parse_slice(toml).expect("fixture slice parses")
}
#[test]
fn a_silent_replacement_is_unproven_not_a_pass() {
assert_eq!(
verdict(Some(0), Some(false), Some(0), Some(0)),
CutoverVerdict::Unproven
);
assert_eq!(
verdict(Some(0), Some(false), Some(0), Some(12)),
CutoverVerdict::Pass
);
}
#[test]
fn any_sign_of_life_beats_everything_else() {
assert_eq!(
verdict(Some(3), Some(false), Some(0), Some(10_000)),
CutoverVerdict::OldStillSpeaks
);
assert_eq!(
verdict(None, Some(true), None, None),
CutoverVerdict::OldStillSpeaks
);
assert_eq!(
verdict(Some(0), Some(false), Some(1), Some(12)),
CutoverVerdict::OldStillSpeaks
);
}
#[test]
fn an_unlistened_entry_cannot_pass() {
assert_eq!(
verdict(None, Some(false), Some(0), None),
CutoverVerdict::Unproven
);
assert_eq!(verdict(None, None, None, None), CutoverVerdict::Unproven);
}
#[test]
fn the_selector_states_the_family_shape() {
let host = slice(
"[registry]\nversion = \"2.0\"\napp = \"demo\"\nconvention = 1\n\
[producer]\nname = \"logs\"\n",
);
assert_eq!(
retired_selector(&host, "logs/errors_total"),
"v1/*/*/logs/logs/errors_total"
);
assert_eq!(
retired_selector(&host, "logs/by_unit/{unit}/messages_total"),
"v1/*/*/logs/logs/by_unit/*/messages_total"
);
let mut svc = slice(
"[registry]\nversion = \"2.0\"\napp = \"demo\"\nconvention = 1\n\
[producer]\nname = \"catalog\"\n",
);
svc.service_origin = Some(zenkey::Declared::parse("@catalog"));
assert_eq!(
retired_selector(&svc, "entity/{id}"),
"v1/@catalog/*/entity/*"
);
}
#[test]
fn the_overall_verdict_is_worst_of() {
let entry = |verdict| RetiredEntry {
producer: "logs".into(),
path: "logs/errors_total".into(),
since: None,
replaced_by: None,
selector: "v1/*/*/logs/logs/errors_total".into(),
wire_samples: crate::report::Asked::NotAsked,
still_declared: None,
subscribers: None,
replacement_samples: crate::report::Asked::NotAsked,
verdict,
};
assert_eq!(overall(&[]), CutoverVerdict::Pass);
assert_eq!(
overall(&[entry(CutoverVerdict::Pass), entry(CutoverVerdict::Pass)]),
CutoverVerdict::Pass
);
assert_eq!(
overall(&[entry(CutoverVerdict::Pass), entry(CutoverVerdict::Unproven)]),
CutoverVerdict::Unproven
);
assert_eq!(
overall(&[
entry(CutoverVerdict::Unproven),
entry(CutoverVerdict::OldStillSpeaks),
entry(CutoverVerdict::Pass),
]),
CutoverVerdict::OldStillSpeaks
);
}
#[test]
fn the_scope_note_states_what_it_cannot_see() {
let note = scope_note(16, "acme/v1/", Duration::from_secs(30));
assert!(note.contains("30s window"));
assert!(
note.contains("cannot cross"),
"a wildcard scope must not be presented as total coverage (O5): {note}"
);
}
}