use std::path::{Path, PathBuf};
use std::time::Duration;
use crate::report::SliceDisagreement;
use crate::report::{ProducerDiff, RegistryDiff};
use crate::{Error, Result};
use zenkey::{Declared, RegistrySlice, parse_slice};
#[derive(Debug, Clone, Default)]
struct ParsedSubjects {
idx: Vec<usize>,
pats: Vec<zenkey::pattern::SubjectPattern>,
}
#[derive(Debug, Clone, Default)]
pub struct SliceSet {
slices: Vec<RegistrySlice>,
raw: Vec<String>,
parsed: Vec<std::collections::BTreeMap<String, ParsedSubjects>>,
by_name: std::collections::BTreeMap<String, usize>,
}
fn parse_subjects(slice: &RegistrySlice) -> std::collections::BTreeMap<String, ParsedSubjects> {
let mut out: std::collections::BTreeMap<String, ParsedSubjects> = Default::default();
for (i, s) in slice.subjects.iter().enumerate() {
if let Ok(p) = zenkey::pattern::SubjectPattern::parse(&s.path) {
let entry = out.entry(s.class.token().to_string()).or_default();
entry.idx.push(i);
entry.pats.push(p);
}
}
out
}
impl SliceSet {
pub fn from_dirs(dirs: &[PathBuf]) -> Result<SliceSet> {
let mut set = SliceSet::default();
for dir in dirs {
let mut paths: Vec<_> = std::fs::read_dir(dir)
.map_err(|e| Error::io(dir, e))?
.filter_map(|e| e.ok().map(|e| e.path()))
.filter(|p| p.extension().is_some_and(|e| e == "toml"))
.filter(|p| p.file_name().is_none_or(|n| n != "types.toml"))
.collect();
paths.sort();
for path in paths {
let text = std::fs::read_to_string(&path).map_err(|e| Error::io(&path, e))?;
let slice = parse_slice(&text)
.map_err(|e| Error::malformed_from(path.display().to_string(), e))?;
set.push(slice, text);
}
}
Ok(set)
}
pub async fn from_bus(fleet: &crate::Fleet<'_>, timeout: Duration) -> Result<SliceSet> {
let pairs = crate::bus::query::fleet_registry_raw(fleet, timeout).await?;
let mut set = SliceSet::default();
for (slice, raw) in pairs {
set.push(slice, raw);
}
Ok(set)
}
fn push(&mut self, slice: RegistrySlice, raw: String) {
let parsed = parse_subjects(&slice);
if let Some(&i) = self.by_name.get(&slice.name) {
self.slices[i] = slice;
self.raw[i] = raw;
self.parsed[i] = parsed;
} else {
self.by_name.insert(slice.name.clone(), self.slices.len());
self.slices.push(slice);
self.raw.push(raw);
self.parsed.push(parsed);
}
}
pub fn entries(&self) -> impl Iterator<Item = (&RegistrySlice, &str)> {
self.slices.iter().zip(self.raw.iter().map(String::as_str))
}
pub fn slices(&self) -> &[RegistrySlice] {
&self.slices
}
pub fn get(&self, name: &str) -> Option<&RegistrySlice> {
self.by_name.get(name).map(|&i| &self.slices[i])
}
pub fn by_service_origin(&self, origin: &str) -> Option<&RegistrySlice> {
self.slices
.iter()
.find(|s| s.service_origin.as_ref().map(Declared::token) == Some(origin))
}
pub fn refine<'s>(
&'s self,
producer: &str,
class: &str,
tail: &[&str],
) -> Option<(&'s zenkey::slice::SubjectDecl, Vec<(String, String)>)> {
let i = *self.by_name.get(producer)?;
let slice = &self.slices[i];
let candidates = self.parsed[i].get(class)?;
let (winner, binds) = zenkey::pattern::best_match(&candidates.pats, tail)?;
let subject_idx = candidates.idx[winner];
Some((
&slice.subjects[subject_idx],
binds.into_iter().map(|(n, v)| (n.to_string(), v)).collect(),
))
}
pub fn from_slices(slices: Vec<RegistrySlice>) -> SliceSet {
let raw = vec![String::new(); slices.len()];
let parsed = slices.iter().map(parse_subjects).collect();
let mut by_name = std::collections::BTreeMap::new();
for (i, s) in slices.iter().enumerate() {
by_name.entry(s.name.clone()).or_insert(i);
}
SliceSet {
slices,
raw,
parsed,
by_name,
}
}
pub fn write_cache(&self, dir: &Path) -> Result<()> {
std::fs::create_dir_all(dir).map_err(|e| Error::io(dir, e))?;
for (slice, raw) in self.slices.iter().zip(&self.raw) {
if raw.is_empty() {
continue; }
let path = dir.join(format!("{}.toml", slice.name));
std::fs::write(&path, raw).map_err(|e| Error::io(&path, e))?;
}
Ok(())
}
pub fn read_cache(dir: &Path) -> SliceSet {
if !dir.is_dir() {
return SliceSet::default();
}
SliceSet::from_dirs(&[dir.to_path_buf()]).unwrap_or_default()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SliceSource {
Bus,
Dirs,
Union,
}
#[derive(Debug, Clone)]
pub struct UnionOutcome {
pub set: SliceSet,
pub from_bus: Vec<String>,
pub dirs_only: Vec<String>,
pub disagreements: Vec<SliceDisagreement>,
}
impl SliceSet {
pub async fn from_union(
fleet: &crate::Fleet<'_>,
dirs: &[std::path::PathBuf],
timeout: std::time::Duration,
) -> Result<UnionOutcome> {
let bus = SliceSet::from_bus(fleet, timeout).await.unwrap_or_default();
let disk = if dirs.is_empty() {
SliceSet::default()
} else {
SliceSet::from_dirs(dirs)?
};
let mut merged = SliceSet::default();
let mut from_bus = Vec::new();
let mut dirs_only = Vec::new();
let mut disagreements = Vec::new();
for (served, raw) in bus.entries() {
from_bus.push(served.name.clone());
if let Some(local) = disk.get(&served.name)
&& (local.version != served.version || local != served)
{
disagreements.push(SliceDisagreement {
producer: served.name.clone(),
bus_version: served.version.clone(),
dirs_version: local.version.clone(),
shape_differs: {
let mut a = served.clone();
let mut b = local.clone();
a.version = String::new();
b.version = String::new();
a != b
},
});
}
merged.push(served.clone(), raw.to_string());
}
for (local, raw) in disk.entries() {
if bus.get(&local.name).is_none() {
dirs_only.push(local.name.clone());
merged.push(local.clone(), raw.to_string());
}
}
Ok(UnionOutcome {
set: merged,
from_bus,
dirs_only,
disagreements,
})
}
}
impl SliceSet {
pub fn diff(&self, local: &SliceSet) -> RegistryDiff {
let served = self;
let mut names: Vec<&str> = served
.slices()
.iter()
.chain(local.slices())
.map(|s| s.name.as_str())
.collect();
names.sort_unstable();
names.dedup();
let mut producers = Vec::new();
for name in names {
let s = served.get(name);
let l = local.get(name);
producers.push(match (s, l) {
(Some(s), Some(l)) => ProducerDiff {
producer: name.to_string(),
served_version: Some(s.version.clone()),
local_version: Some(l.version.clone()),
findings: zenkey::slice::diff(s, l)
.iter()
.map(|f| f.summary())
.collect(),
},
(Some(s), None) => ProducerDiff {
producer: name.to_string(),
served_version: Some(s.version.clone()),
local_version: None,
findings: vec!["served by the fleet, absent from the local registry".into()],
},
(None, Some(l)) => ProducerDiff {
producer: name.to_string(),
served_version: None,
local_version: Some(l.version.clone()),
findings: vec![
"declared locally, not served by any origin — down, or not deployed \
(silence is not a verdict, RFC 05 §3.1)"
.into(),
],
},
(None, None) => unreachable!("name came from one of the two sets"),
});
}
RegistryDiff { producers }
}
}
#[cfg(test)]
impl SliceSet {
pub(crate) fn from_toml_for_tests(toml: &str) -> SliceSet {
let mut set = SliceSet::default();
set.push(parse_slice(toml).unwrap(), toml.to_string());
set
}
}
#[cfg(test)]
mod tests {
use super::*;
const A: &str = r#"
[registry]
version = "1.0"
app = "t"
convention = 1
[producer]
name = "alpha"
[[subject]]
path = "flow/{q}"
class = "telemetry"
type = "Point"
[[subject]]
path = "flow/special"
class = "telemetry"
type = "Special"
"#;
#[test]
fn refine_uses_shared_precedence() {
let mut set = SliceSet::default();
set.push(parse_slice(A).unwrap(), A.to_string());
let (s, binds) = set
.refine("alpha", "telemetry", &["flow", "special"])
.unwrap();
assert_eq!(s.type_name, "Special");
assert!(binds.is_empty());
let (s, binds) = set.refine("alpha", "telemetry", &["flow", "p95"]).unwrap();
assert_eq!(s.type_name, "Point");
assert_eq!(binds, vec![("q".to_string(), "p95".to_string())]);
assert!(set.refine("alpha", "state", &["flow", "p95"]).is_none());
}
#[test]
fn cache_round_trips_and_last_slice_wins() {
let mut set = SliceSet::default();
set.push(parse_slice(A).unwrap(), A.to_string());
set.push(parse_slice(A).unwrap(), A.to_string());
assert_eq!(set.slices().len(), 1);
let dir = std::env::temp_dir().join(format!("zenkey-fleet-cache-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
set.write_cache(&dir).unwrap();
let back = SliceSet::read_cache(&dir);
assert_eq!(back.slices().len(), 1);
assert_eq!(back.get("alpha").unwrap().subjects.len(), 2);
let _ = std::fs::remove_dir_all(&dir);
assert!(
SliceSet::read_cache(Path::new("/nonexistent-zkf"))
.slices()
.is_empty()
);
}
#[test]
fn a_re_pushed_producer_keeps_its_place_and_shadowing_is_first_wins() {
let newer = A.replace("version = \"1.0\"", "version = \"9.9\"");
let other = A.replace("name = \"alpha\"", "name = \"beta\"");
let mut set = SliceSet::default();
set.push(parse_slice(A).unwrap(), A.to_string());
set.push(parse_slice(&other).unwrap(), other.clone());
set.push(parse_slice(&newer).unwrap(), newer.clone());
assert_eq!(set.slices().len(), 2, "a re-push replaces, never appends");
assert_eq!(
set.slices()[0].name,
"alpha",
"the replacement keeps the producer's position"
);
assert_eq!(set.get("alpha").unwrap().version, "9.9", "last push wins");
assert_eq!(
set.entries().next().unwrap().1,
newer,
"the raw TOML rides with the slice it was parsed from"
);
assert!(set.get("gamma").is_none());
assert_eq!(
set.refine("alpha", "telemetry", &["flow", "special"])
.unwrap()
.0
.type_name,
"Special"
);
assert!(set.refine("gamma", "telemetry", &["flow"]).is_none());
let shadowed =
SliceSet::from_slices(vec![parse_slice(A).unwrap(), parse_slice(&newer).unwrap()]);
assert_eq!(
shadowed.get("alpha").unwrap().version,
"1.0",
"the earlier of two same-named slices answers"
);
assert_eq!(shadowed.slices().len(), 2, "neither is dropped");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn union_degrades_to_dirs_when_the_bus_is_silent() {
let session = crate::bus::session::open(&[], &[], false).await.unwrap();
let dir =
std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../fixture-tests/registry");
let out = SliceSet::from_union(
&crate::Fleet::new(&session, ""),
&[dir],
std::time::Duration::from_millis(200),
)
.await
.unwrap();
assert!(out.from_bus.is_empty(), "no bus answered");
assert!(!out.dirs_only.is_empty(), "dirs supplied the slices");
assert!(out.disagreements.is_empty());
assert_eq!(out.set.slices().len(), out.dirs_only.len());
}
fn set(toml: &str) -> SliceSet {
SliceSet::from_slices(vec![zenkey::parse_slice(toml).unwrap()])
}
const SERVED: &str = r#"
[registry]
version = "2.0"
app = "t"
convention = 1
[producer]
name = "netring"
[[subject]]
path = "flows"
class = "telemetry"
type = "TelemetryPoint"
[[subject]]
path = "brand/new"
class = "telemetry"
type = "TelemetryPoint"
"#;
const LOCAL: &str = r#"
[registry]
version = "1.0"
app = "t"
convention = 1
[producer]
name = "netring"
[[subject]]
path = "flows"
class = "telemetry"
type = "TelemetryPoint"
"#;
#[test]
fn the_diff_names_the_one_subject_that_moved() {
let report = set(SERVED).diff(&set(LOCAL));
assert_eq!(report.producers.len(), 1);
let p = &report.producers[0];
assert_eq!(p.served_version.as_deref(), Some("2.0"));
assert_eq!(p.local_version.as_deref(), Some("1.0"));
assert!(
p.findings.iter().any(|f| f.contains("brand/new")),
"{:?}",
p.findings
);
assert!(
p.findings.iter().any(|f| f.contains("2.0")),
"the version skew is a finding too: {:?}",
p.findings
);
}
#[test]
fn one_sided_producers_explain_themselves() {
let empty = SliceSet::from_slices(vec![]);
let served_only = set(SERVED).diff(&empty);
assert!(served_only.producers[0].findings[0].contains("absent from the local registry"));
assert!(served_only.producers[0].local_version.is_none());
let local_only = empty.diff(&set(LOCAL));
assert!(local_only.producers[0].findings[0].contains("silence is not a verdict"));
assert!(local_only.producers[0].served_version.is_none());
}
}