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(feature = "decode")]
466pub(crate) fn rpc_key(base: &str, slice: &RegistrySlice, procedure: &str) -> Result<String> {
467 Ok(match &slice.service_origin {
468 Some(origin) => {
469 let o = origin.known().ok_or_else(|| {
475 Error::malformed(
476 format!("slice {}", slice.name),
477 format!("carries {:?} as a service origin", origin.token()),
478 )
479 })?;
480 zenkey::grammar::with_base(base, zenkey::selector::service_rpc(o, &[procedure]))
481 }
482 None => {
483 zenkey::grammar::with_base(base, zenkey::selector::fleet_rpc(&slice.name, &[procedure]))
484 }
485 })
486}
487
488#[cfg(test)]
489impl SliceSet {
490 pub(crate) fn from_toml_for_tests(toml: &str) -> SliceSet {
492 let mut set = SliceSet::default();
493 set.push(parse_slice(toml).unwrap(), toml.to_string());
494 set
495 }
496}
497
498#[cfg(test)]
499mod tests {
500 use super::*;
501
502 const A: &str = r#"
503 [registry]
504 version = "1.0"
505 app = "t"
506 convention = 1
507 [producer]
508 name = "alpha"
509 [[subject]]
510 path = "flow/{q}"
511 class = "telemetry"
512 type = "Point"
513 [[subject]]
514 path = "flow/special"
515 class = "telemetry"
516 type = "Special"
517 "#;
518
519 #[test]
520 fn refine_uses_shared_precedence() {
521 let mut set = SliceSet::default();
522 set.push(parse_slice(A).unwrap(), A.to_string());
523 let (s, binds) = set
525 .refine("alpha", "telemetry", &["flow", "special"])
526 .unwrap();
527 assert_eq!(s.type_name, "Special");
528 assert!(binds.is_empty());
529 let (s, binds) = set.refine("alpha", "telemetry", &["flow", "p95"]).unwrap();
530 assert_eq!(s.type_name, "Point");
531 assert_eq!(binds, vec![("q".to_string(), "p95".to_string())]);
532 assert!(set.refine("alpha", "state", &["flow", "p95"]).is_none());
533 }
534
535 #[test]
544 fn the_fold_to_one_slice_per_producer_records_what_it_discarded() {
545 let served = |origin: &str, raw: &str| crate::ServedSlice {
546 origin: origin.to_string(),
547 slice: parse_slice(raw).unwrap(),
548 raw: raw.to_string(),
549 };
550
551 let mut b_variant = A.to_string();
553 b_variant.push_str(
554 "\n[[subject]]\npath = \"extra\"\nclass = \"state\"\ntype = \"E\"\nttl_s = 1\n",
555 );
556 let set = SliceSet::from_served(vec![
557 served("h-aaaaaaaaaaaa", A),
558 served("h-bbbbbbbbbbbb", &b_variant),
559 ]);
560 assert_eq!(set.slices().len(), 1, "still one slice per producer");
561
562 let collapsed = set
563 .collapsed()
564 .as_option()
565 .copied()
566 .expect("a bus fold asked");
567 assert_eq!(collapsed.len(), 1, "{collapsed:?}");
568 assert_eq!(collapsed[0].producer, "alpha");
569 assert_eq!(
570 collapsed[0].origins,
571 vec!["h-aaaaaaaaaaaa", "h-bbbbbbbbbbbb"]
572 );
573 assert!(
574 !collapsed[0].agreed,
575 "the discarded answer differed — that is the finding"
576 );
577
578 let agreeing = SliceSet::from_served(vec![
580 served("h-aaaaaaaaaaaa", A),
581 served("h-bbbbbbbbbbbb", A),
582 ]);
583 assert!(agreeing.collapsed().as_option().copied().expect("asked")[0].agreed);
584
585 let one_host = SliceSet::from_served(vec![served("h-aaaaaaaaaaaa", A)]);
588 assert_eq!(
589 one_host.collapsed(),
590 Asked::Asked(&[][..]),
591 "asked, and nothing was collapsed"
592 );
593 assert_eq!(
595 SliceSet::from_slices(vec![parse_slice(A).unwrap()]).collapsed(),
596 Asked::NotAsked,
597 "not asked is not \"the fleet agrees\""
598 );
599 assert_eq!(SliceSet::default().collapsed(), Asked::NotAsked);
600 }
601
602 #[test]
603 fn cache_round_trips_and_last_slice_wins() {
604 let mut set = SliceSet::default();
605 set.push(parse_slice(A).unwrap(), A.to_string());
606 set.push(parse_slice(A).unwrap(), A.to_string());
608 assert_eq!(set.slices().len(), 1);
609
610 let dir = std::env::temp_dir().join(format!("zenkey-fleet-cache-{}", std::process::id()));
611 let _ = std::fs::remove_dir_all(&dir);
612 set.write_cache(&dir).unwrap();
613 let back = SliceSet::read_cache(&dir);
614 assert_eq!(back.slices().len(), 1);
615 assert_eq!(back.get("alpha").unwrap().subjects.len(), 2);
616 let _ = std::fs::remove_dir_all(&dir);
617 assert!(
619 SliceSet::read_cache(Path::new("/nonexistent-zkf"))
620 .slices()
621 .is_empty()
622 );
623 }
624
625 #[test]
633 fn a_re_pushed_producer_keeps_its_place_and_shadowing_is_first_wins() {
634 let newer = A.replace("version = \"1.0\"", "version = \"9.9\"");
635 let other = A.replace("name = \"alpha\"", "name = \"beta\"");
636
637 let mut set = SliceSet::default();
638 set.push(parse_slice(A).unwrap(), A.to_string());
639 set.push(parse_slice(&other).unwrap(), other.clone());
640 set.push(parse_slice(&newer).unwrap(), newer.clone());
641
642 assert_eq!(set.slices().len(), 2, "a re-push replaces, never appends");
643 assert_eq!(
644 set.slices()[0].name,
645 "alpha",
646 "the replacement keeps the producer's position"
647 );
648 assert_eq!(set.get("alpha").unwrap().version, "9.9", "last push wins");
649 assert_eq!(
650 set.entries().next().unwrap().1,
651 newer,
652 "the raw TOML rides with the slice it was parsed from"
653 );
654 assert!(set.get("gamma").is_none());
655 assert_eq!(
657 set.refine("alpha", "telemetry", &["flow", "special"])
658 .unwrap()
659 .0
660 .type_name,
661 "Special"
662 );
663 assert!(set.refine("gamma", "telemetry", &["flow"]).is_none());
664
665 let shadowed =
667 SliceSet::from_slices(vec![parse_slice(A).unwrap(), parse_slice(&newer).unwrap()]);
668 assert_eq!(
669 shadowed.get("alpha").unwrap().version,
670 "1.0",
671 "the earlier of two same-named slices answers"
672 );
673 assert_eq!(shadowed.slices().len(), 2, "neither is dropped");
674 }
675
676 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
679 async fn union_degrades_to_dirs_when_the_bus_is_silent() {
680 let session = crate::bus::session::open(&[], &[], false).await.unwrap();
681 let dir =
682 std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../fixture-tests/registry");
683 let out = SliceSet::from_union(
684 &crate::Fleet::new(&session, ""),
685 &[dir],
686 std::time::Duration::from_millis(200),
687 )
688 .await
689 .unwrap();
690 assert!(out.from_bus.is_empty(), "no bus answered");
691 assert!(!out.dirs_only.is_empty(), "dirs supplied the slices");
692 assert!(out.disagreements.is_empty());
693 assert_eq!(out.set.slices().len(), out.dirs_only.len());
694 }
695
696 fn set(toml: &str) -> SliceSet {
697 SliceSet::from_slices(vec![zenkey::parse_slice(toml).unwrap()])
698 }
699
700 const SERVED: &str = r#"
701[registry]
702version = "2.0"
703app = "t"
704convention = 1
705[producer]
706name = "netring"
707[[subject]]
708path = "flows"
709class = "telemetry"
710type = "TelemetryPoint"
711[[subject]]
712path = "brand/new"
713class = "telemetry"
714type = "TelemetryPoint"
715"#;
716
717 const LOCAL: &str = r#"
718[registry]
719version = "1.0"
720app = "t"
721convention = 1
722[producer]
723name = "netring"
724[[subject]]
725path = "flows"
726class = "telemetry"
727type = "TelemetryPoint"
728"#;
729
730 #[test]
733 fn the_diff_names_the_one_subject_that_moved() {
734 let report = set(SERVED).diff(&set(LOCAL));
735 assert_eq!(report.producers.len(), 1);
736 let p = &report.producers[0];
737 assert_eq!(p.served_version.as_deref(), Some("2.0"));
738 assert_eq!(p.local_version.as_deref(), Some("1.0"));
739 assert!(
740 p.findings.iter().any(|f| f.contains("brand/new")),
741 "{:?}",
742 p.findings
743 );
744 assert!(
745 p.findings.iter().any(|f| f.contains("2.0")),
746 "the version skew is a finding too: {:?}",
747 p.findings
748 );
749 }
750
751 #[test]
754 fn one_sided_producers_explain_themselves() {
755 let empty = SliceSet::from_slices(vec![]);
756 let served_only = set(SERVED).diff(&empty);
757 assert!(served_only.producers[0].findings[0].contains("absent from the local registry"));
758 assert!(served_only.producers[0].local_version.is_none());
759
760 let local_only = empty.diff(&set(LOCAL));
761 assert!(local_only.producers[0].findings[0].contains("silence is not a verdict"));
762 assert!(local_only.producers[0].served_version.is_none());
763 }
764}