1use std::path::{Path, PathBuf};
10use std::time::Duration;
11
12use crate::report::SliceDisagreement;
13use crate::report::{Asked, CollapsedProducer, ProducerDiff, RegistryDiff};
14use crate::{Error, Result};
15use zenkey::{Declared, RegistrySlice, parse_slice};
16
17#[derive(Debug, Clone, Default)]
25struct ParsedSubjects {
26 idx: Vec<usize>,
28 pats: Vec<zenkey::pattern::SubjectPattern>,
30}
31
32#[derive(Debug, Clone, Default)]
34pub struct SliceSet {
35 slices: Vec<RegistrySlice>,
36 raw: Vec<String>,
39 parsed: Vec<std::collections::BTreeMap<String, ParsedSubjects>>,
43 by_name: std::collections::BTreeMap<String, usize>,
57 collapsed: Asked<Vec<CollapsedProducer>>,
62}
63
64fn parse_subjects(slice: &RegistrySlice) -> std::collections::BTreeMap<String, ParsedSubjects> {
68 let mut out: std::collections::BTreeMap<String, ParsedSubjects> = Default::default();
69 for (i, s) in slice.subjects.iter().enumerate() {
70 if let Ok(p) = zenkey::pattern::SubjectPattern::parse(&s.path) {
71 let entry = out.entry(s.class.token().to_string()).or_default();
72 entry.idx.push(i);
73 entry.pats.push(p);
74 }
75 }
76 out
77}
78
79impl SliceSet {
80 pub fn from_dirs(dirs: &[PathBuf]) -> Result<SliceSet> {
84 let mut set = SliceSet::default();
85 for dir in dirs {
86 let mut paths: Vec<_> = std::fs::read_dir(dir)
87 .map_err(|e| Error::io(dir, e))?
88 .filter_map(|e| e.ok().map(|e| e.path()))
89 .filter(|p| p.extension().is_some_and(|e| e == "toml"))
90 .filter(|p| p.file_name().is_none_or(|n| n != "types.toml"))
91 .collect();
92 paths.sort();
93 for path in paths {
94 let text = std::fs::read_to_string(&path).map_err(|e| Error::io(&path, e))?;
95 let slice = parse_slice(&text)
96 .map_err(|e| Error::malformed_from(path.display().to_string(), e))?;
97 set.push(slice, text);
98 }
99 }
100 Ok(set)
101 }
102
103 pub async fn from_bus(fleet: &crate::Fleet<'_>, timeout: Duration) -> Result<SliceSet> {
112 Ok(SliceSet::from_served(
113 crate::bus::query::fleet_registry_by_origin(fleet, timeout).await?,
114 ))
115 }
116
117 pub fn from_served(served: Vec<crate::ServedSlice>) -> SliceSet {
123 let mut answers: std::collections::BTreeMap<
125 String,
126 (Vec<String>, Vec<String>, String, bool),
127 > = Default::default();
128 let mut set = SliceSet::default();
129 for s in served {
130 let entry = answers
131 .entry(s.slice.name.clone())
132 .or_insert_with(|| (Vec::new(), Vec::new(), s.raw.clone(), true));
133 entry.0.push(s.origin);
134 entry.1.push(s.slice.version.clone());
135 if s.raw != entry.2 {
139 entry.3 = false;
140 }
141 set.push(s.slice, s.raw);
142 }
143 set.collapsed = Asked::Asked(
147 answers
148 .into_iter()
149 .filter(|(_, (origins, ..))| origins.len() > 1)
151 .map(
152 |(producer, (origins, versions, _, agreed))| CollapsedProducer {
153 producer,
154 origins,
155 versions,
156 agreed,
157 },
158 )
159 .collect(),
160 );
161 set
162 }
163
164 pub fn collapsed(&self) -> Asked<&[CollapsedProducer]> {
175 match &self.collapsed {
176 Asked::NotAsked => Asked::NotAsked,
177 Asked::Asked(v) => Asked::Asked(v.as_slice()),
178 }
179 }
180
181 fn push(&mut self, slice: RegistrySlice, raw: String) {
182 let parsed = parse_subjects(&slice);
188 if let Some(&i) = self.by_name.get(&slice.name) {
189 self.slices[i] = slice;
190 self.raw[i] = raw;
191 self.parsed[i] = parsed;
192 } else {
193 self.by_name.insert(slice.name.clone(), self.slices.len());
194 self.slices.push(slice);
195 self.raw.push(raw);
196 self.parsed.push(parsed);
197 }
198 }
199
200 pub fn entries(&self) -> impl Iterator<Item = (&RegistrySlice, &str)> {
204 self.slices.iter().zip(self.raw.iter().map(String::as_str))
205 }
206
207 pub fn slices(&self) -> &[RegistrySlice] {
208 &self.slices
209 }
210
211 pub fn get(&self, name: &str) -> Option<&RegistrySlice> {
212 self.by_name.get(name).map(|&i| &self.slices[i])
213 }
214
215 pub fn by_service_origin(&self, origin: &str) -> Option<&RegistrySlice> {
218 self.slices
219 .iter()
220 .find(|s| s.service_origin.as_ref().map(Declared::token) == Some(origin))
221 }
222
223 pub fn refine<'s>(
226 &'s self,
227 producer: &str,
228 class: &str,
229 tail: &[&str],
230 ) -> Option<(&'s zenkey::slice::SubjectDecl, Vec<(String, String)>)> {
231 let i = *self.by_name.get(producer)?;
232 let slice = &self.slices[i];
233 let candidates = self.parsed[i].get(class)?;
237 let (winner, binds) = zenkey::pattern::best_match(&candidates.pats, tail)?;
238 let subject_idx = candidates.idx[winner];
239 Some((
240 &slice.subjects[subject_idx],
241 binds.into_iter().map(|(n, v)| (n.to_string(), v)).collect(),
242 ))
243 }
244
245 pub fn from_slices(slices: Vec<RegistrySlice>) -> SliceSet {
248 let raw = vec![String::new(); slices.len()];
249 let parsed = slices.iter().map(parse_subjects).collect();
250 let mut by_name = std::collections::BTreeMap::new();
253 for (i, s) in slices.iter().enumerate() {
254 by_name.entry(s.name.clone()).or_insert(i);
255 }
256 SliceSet {
257 slices,
258 raw,
259 parsed,
260 by_name,
261 collapsed: Asked::NotAsked,
266 }
267 }
268
269 pub fn write_cache(&self, dir: &Path) -> Result<()> {
273 std::fs::create_dir_all(dir).map_err(|e| Error::io(dir, e))?;
274 for (slice, raw) in self.slices.iter().zip(&self.raw) {
275 if raw.is_empty() {
276 continue; }
278 let path = dir.join(format!("{}.toml", slice.name));
279 std::fs::write(&path, raw).map_err(|e| Error::io(&path, e))?;
280 }
281 Ok(())
282 }
283
284 pub fn read_cache(dir: &Path) -> SliceSet {
287 if !dir.is_dir() {
288 return SliceSet::default();
289 }
290 SliceSet::from_dirs(&[dir.to_path_buf()]).unwrap_or_default()
291 }
292}
293
294#[derive(Debug, Clone, Copy, PartialEq, Eq)]
297pub enum SliceSource {
298 Bus,
299 Dirs,
300 Union,
301}
302
303#[derive(Debug, Clone)]
305pub struct UnionOutcome {
306 pub set: SliceSet,
307 pub from_bus: Vec<String>,
309 pub dirs_only: Vec<String>,
311 pub disagreements: Vec<SliceDisagreement>,
312}
313
314impl SliceSet {
315 pub async fn from_union(
322 fleet: &crate::Fleet<'_>,
323 dirs: &[std::path::PathBuf],
324 timeout: std::time::Duration,
325 ) -> Result<UnionOutcome> {
326 let bus = SliceSet::from_bus(fleet, timeout).await.unwrap_or_default();
327 let disk = if dirs.is_empty() {
328 SliceSet::default()
329 } else {
330 SliceSet::from_dirs(dirs)?
331 };
332
333 let mut merged = SliceSet::default();
338 let mut from_bus = Vec::new();
339 let mut dirs_only = Vec::new();
340 let mut disagreements = Vec::new();
341
342 for (served, raw) in bus.entries() {
343 from_bus.push(served.name.clone());
344 if let Some(local) = disk.get(&served.name)
345 && (local.version != served.version || local != served)
346 {
347 disagreements.push(SliceDisagreement {
348 producer: served.name.clone(),
349 bus_version: served.version.clone(),
350 dirs_version: local.version.clone(),
351 shape_differs: {
352 let mut a = served.clone();
354 let mut b = local.clone();
355 a.version = String::new();
356 b.version = String::new();
357 a != b
358 },
359 });
360 }
361 merged.push(served.clone(), raw.to_string());
362 }
363 for (local, raw) in disk.entries() {
364 if bus.get(&local.name).is_none() {
365 dirs_only.push(local.name.clone());
366 merged.push(local.clone(), raw.to_string());
367 }
368 }
369 merged.collapsed = bus.collapsed;
374
375 Ok(UnionOutcome {
376 set: merged,
377 from_bus,
378 dirs_only,
379 disagreements,
380 })
381 }
382}
383
384impl SliceSet {
385 pub fn diff(&self, local: &SliceSet) -> RegistryDiff {
393 let served = self;
394 let mut names: Vec<&str> = served
395 .slices()
396 .iter()
397 .chain(local.slices())
398 .map(|s| s.name.as_str())
399 .collect();
400 names.sort_unstable();
401 names.dedup();
402
403 let mut producers = Vec::new();
404 for name in names {
405 let s = served.get(name);
406 let l = local.get(name);
407 producers.push(match (s, l) {
408 (Some(s), Some(l)) => ProducerDiff {
409 producer: name.to_string(),
410 served_version: Some(s.version.clone()),
411 local_version: Some(l.version.clone()),
412 findings: zenkey::slice::diff(s, l)
413 .iter()
414 .map(|f| f.summary())
415 .collect(),
416 },
417 (Some(s), None) => ProducerDiff {
422 producer: name.to_string(),
423 served_version: Some(s.version.clone()),
424 local_version: None,
425 findings: vec!["served by the fleet, absent from the local registry".into()],
426 },
427 (None, Some(l)) => ProducerDiff {
428 producer: name.to_string(),
429 served_version: None,
430 local_version: Some(l.version.clone()),
431 findings: vec![
432 "declared locally, not served by any origin — down, or not deployed \
433 (silence is not a verdict, RFC 05 §3.1)"
434 .into(),
435 ],
436 },
437 (None, None) => unreachable!("name came from one of the two sets"),
438 });
439 }
440 RegistryDiff {
441 producers,
442 collapsed: match served.collapsed() {
446 Asked::NotAsked => Asked::NotAsked,
447 Asked::Asked(c) => Asked::Asked(c.to_vec()),
448 },
449 }
450 }
451}
452
453#[cfg(test)]
454impl SliceSet {
455 pub(crate) fn from_toml_for_tests(toml: &str) -> SliceSet {
457 let mut set = SliceSet::default();
458 set.push(parse_slice(toml).unwrap(), toml.to_string());
459 set
460 }
461}
462
463#[cfg(test)]
464mod tests {
465 use super::*;
466
467 const A: &str = r#"
468 [registry]
469 version = "1.0"
470 app = "t"
471 convention = 1
472 [producer]
473 name = "alpha"
474 [[subject]]
475 path = "flow/{q}"
476 class = "telemetry"
477 type = "Point"
478 [[subject]]
479 path = "flow/special"
480 class = "telemetry"
481 type = "Special"
482 "#;
483
484 #[test]
485 fn refine_uses_shared_precedence() {
486 let mut set = SliceSet::default();
487 set.push(parse_slice(A).unwrap(), A.to_string());
488 let (s, binds) = set
490 .refine("alpha", "telemetry", &["flow", "special"])
491 .unwrap();
492 assert_eq!(s.type_name, "Special");
493 assert!(binds.is_empty());
494 let (s, binds) = set.refine("alpha", "telemetry", &["flow", "p95"]).unwrap();
495 assert_eq!(s.type_name, "Point");
496 assert_eq!(binds, vec![("q".to_string(), "p95".to_string())]);
497 assert!(set.refine("alpha", "state", &["flow", "p95"]).is_none());
498 }
499
500 #[test]
509 fn the_fold_to_one_slice_per_producer_records_what_it_discarded() {
510 let served = |origin: &str, raw: &str| crate::ServedSlice {
511 origin: origin.to_string(),
512 slice: parse_slice(raw).unwrap(),
513 raw: raw.to_string(),
514 };
515
516 let mut b_variant = A.to_string();
518 b_variant.push_str(
519 "\n[[subject]]\npath = \"extra\"\nclass = \"state\"\ntype = \"E\"\nttl_s = 1\n",
520 );
521 let set = SliceSet::from_served(vec![
522 served("h-aaaaaaaaaaaa", A),
523 served("h-bbbbbbbbbbbb", &b_variant),
524 ]);
525 assert_eq!(set.slices().len(), 1, "still one slice per producer");
526
527 let collapsed = set
528 .collapsed()
529 .as_option()
530 .copied()
531 .expect("a bus fold asked");
532 assert_eq!(collapsed.len(), 1, "{collapsed:?}");
533 assert_eq!(collapsed[0].producer, "alpha");
534 assert_eq!(
535 collapsed[0].origins,
536 vec!["h-aaaaaaaaaaaa", "h-bbbbbbbbbbbb"]
537 );
538 assert!(
539 !collapsed[0].agreed,
540 "the discarded answer differed — that is the finding"
541 );
542
543 let agreeing = SliceSet::from_served(vec![
545 served("h-aaaaaaaaaaaa", A),
546 served("h-bbbbbbbbbbbb", A),
547 ]);
548 assert!(agreeing.collapsed().as_option().copied().expect("asked")[0].agreed);
549
550 let one_host = SliceSet::from_served(vec![served("h-aaaaaaaaaaaa", A)]);
553 assert_eq!(
554 one_host.collapsed(),
555 Asked::Asked(&[][..]),
556 "asked, and nothing was collapsed"
557 );
558 assert_eq!(
560 SliceSet::from_slices(vec![parse_slice(A).unwrap()]).collapsed(),
561 Asked::NotAsked,
562 "not asked is not \"the fleet agrees\""
563 );
564 assert_eq!(SliceSet::default().collapsed(), Asked::NotAsked);
565 }
566
567 #[test]
568 fn cache_round_trips_and_last_slice_wins() {
569 let mut set = SliceSet::default();
570 set.push(parse_slice(A).unwrap(), A.to_string());
571 set.push(parse_slice(A).unwrap(), A.to_string());
573 assert_eq!(set.slices().len(), 1);
574
575 let dir = std::env::temp_dir().join(format!("zenkey-fleet-cache-{}", std::process::id()));
576 let _ = std::fs::remove_dir_all(&dir);
577 set.write_cache(&dir).unwrap();
578 let back = SliceSet::read_cache(&dir);
579 assert_eq!(back.slices().len(), 1);
580 assert_eq!(back.get("alpha").unwrap().subjects.len(), 2);
581 let _ = std::fs::remove_dir_all(&dir);
582 assert!(
584 SliceSet::read_cache(Path::new("/nonexistent-zkf"))
585 .slices()
586 .is_empty()
587 );
588 }
589
590 #[test]
598 fn a_re_pushed_producer_keeps_its_place_and_shadowing_is_first_wins() {
599 let newer = A.replace("version = \"1.0\"", "version = \"9.9\"");
600 let other = A.replace("name = \"alpha\"", "name = \"beta\"");
601
602 let mut set = SliceSet::default();
603 set.push(parse_slice(A).unwrap(), A.to_string());
604 set.push(parse_slice(&other).unwrap(), other.clone());
605 set.push(parse_slice(&newer).unwrap(), newer.clone());
606
607 assert_eq!(set.slices().len(), 2, "a re-push replaces, never appends");
608 assert_eq!(
609 set.slices()[0].name,
610 "alpha",
611 "the replacement keeps the producer's position"
612 );
613 assert_eq!(set.get("alpha").unwrap().version, "9.9", "last push wins");
614 assert_eq!(
615 set.entries().next().unwrap().1,
616 newer,
617 "the raw TOML rides with the slice it was parsed from"
618 );
619 assert!(set.get("gamma").is_none());
620 assert_eq!(
622 set.refine("alpha", "telemetry", &["flow", "special"])
623 .unwrap()
624 .0
625 .type_name,
626 "Special"
627 );
628 assert!(set.refine("gamma", "telemetry", &["flow"]).is_none());
629
630 let shadowed =
632 SliceSet::from_slices(vec![parse_slice(A).unwrap(), parse_slice(&newer).unwrap()]);
633 assert_eq!(
634 shadowed.get("alpha").unwrap().version,
635 "1.0",
636 "the earlier of two same-named slices answers"
637 );
638 assert_eq!(shadowed.slices().len(), 2, "neither is dropped");
639 }
640
641 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
644 async fn union_degrades_to_dirs_when_the_bus_is_silent() {
645 let session = crate::bus::session::open(&[], &[], false).await.unwrap();
646 let dir =
647 std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../fixture-tests/registry");
648 let out = SliceSet::from_union(
649 &crate::Fleet::new(&session, ""),
650 &[dir],
651 std::time::Duration::from_millis(200),
652 )
653 .await
654 .unwrap();
655 assert!(out.from_bus.is_empty(), "no bus answered");
656 assert!(!out.dirs_only.is_empty(), "dirs supplied the slices");
657 assert!(out.disagreements.is_empty());
658 assert_eq!(out.set.slices().len(), out.dirs_only.len());
659 }
660
661 fn set(toml: &str) -> SliceSet {
662 SliceSet::from_slices(vec![zenkey::parse_slice(toml).unwrap()])
663 }
664
665 const SERVED: &str = r#"
666[registry]
667version = "2.0"
668app = "t"
669convention = 1
670[producer]
671name = "netring"
672[[subject]]
673path = "flows"
674class = "telemetry"
675type = "TelemetryPoint"
676[[subject]]
677path = "brand/new"
678class = "telemetry"
679type = "TelemetryPoint"
680"#;
681
682 const LOCAL: &str = r#"
683[registry]
684version = "1.0"
685app = "t"
686convention = 1
687[producer]
688name = "netring"
689[[subject]]
690path = "flows"
691class = "telemetry"
692type = "TelemetryPoint"
693"#;
694
695 #[test]
698 fn the_diff_names_the_one_subject_that_moved() {
699 let report = set(SERVED).diff(&set(LOCAL));
700 assert_eq!(report.producers.len(), 1);
701 let p = &report.producers[0];
702 assert_eq!(p.served_version.as_deref(), Some("2.0"));
703 assert_eq!(p.local_version.as_deref(), Some("1.0"));
704 assert!(
705 p.findings.iter().any(|f| f.contains("brand/new")),
706 "{:?}",
707 p.findings
708 );
709 assert!(
710 p.findings.iter().any(|f| f.contains("2.0")),
711 "the version skew is a finding too: {:?}",
712 p.findings
713 );
714 }
715
716 #[test]
719 fn one_sided_producers_explain_themselves() {
720 let empty = SliceSet::from_slices(vec![]);
721 let served_only = set(SERVED).diff(&empty);
722 assert!(served_only.producers[0].findings[0].contains("absent from the local registry"));
723 assert!(served_only.producers[0].local_version.is_none());
724
725 let local_only = empty.diff(&set(LOCAL));
726 assert!(local_only.producers[0].findings[0].contains("silence is not a verdict"));
727 assert!(local_only.producers[0].served_version.is_none());
728 }
729}