1use std::path::{Path, PathBuf};
10use std::time::Duration;
11
12use anyhow::{Result, anyhow};
13use zenkey::{RegistrySlice, parse_slice};
14use zenoh::Session;
15
16#[derive(Debug, Clone, Default)]
18pub struct SliceSet {
19 slices: Vec<RegistrySlice>,
20 raw: Vec<String>,
23}
24
25impl SliceSet {
26 pub fn from_dirs(dirs: &[PathBuf]) -> Result<SliceSet> {
30 let mut set = SliceSet::default();
31 for dir in dirs {
32 let mut paths: Vec<_> = std::fs::read_dir(dir)
33 .map_err(|e| anyhow!("--registry {}: {e}", dir.display()))?
34 .filter_map(|e| e.ok().map(|e| e.path()))
35 .filter(|p| p.extension().is_some_and(|e| e == "toml"))
36 .filter(|p| p.file_name().is_none_or(|n| n != "types.toml"))
37 .collect();
38 paths.sort();
39 for path in paths {
40 let text = std::fs::read_to_string(&path)
41 .map_err(|e| anyhow!("{}: {e}", path.display()))?;
42 let slice = parse_slice(&text).map_err(|e| {
43 anyhow!(
44 "{}: does not parse as a registry slice: {e}",
45 path.display()
46 )
47 })?;
48 set.push(slice, text);
49 }
50 }
51 Ok(set)
52 }
53
54 pub async fn from_bus(session: &Session, base: &str, timeout: Duration) -> Result<SliceSet> {
57 let pairs = crate::query::fleet_registry_raw(session, base, timeout).await?;
58 let mut set = SliceSet::default();
59 for (slice, raw) in pairs {
60 set.push(slice, raw);
61 }
62 Ok(set)
63 }
64
65 fn push(&mut self, slice: RegistrySlice, raw: String) {
66 if let Some(i) = self.slices.iter().position(|s| s.name == slice.name) {
70 self.slices[i] = slice;
71 self.raw[i] = raw;
72 } else {
73 self.slices.push(slice);
74 self.raw.push(raw);
75 }
76 }
77
78 pub fn slices(&self) -> &[RegistrySlice] {
79 &self.slices
80 }
81
82 pub fn get(&self, name: &str) -> Option<&RegistrySlice> {
83 self.slices.iter().find(|s| s.name == name)
84 }
85
86 pub fn by_service_origin(&self, origin: &str) -> Option<&RegistrySlice> {
89 self.slices
90 .iter()
91 .find(|s| s.service_origin.as_deref() == Some(origin))
92 }
93
94 pub fn refine<'s>(
97 &'s self,
98 producer: &str,
99 class: &str,
100 tail: &[&str],
101 ) -> Option<(&'s zenkey::slice::SubjectDecl, Vec<(String, String)>)> {
102 let slice = self.get(producer)?;
103 let candidates: Vec<(usize, zenkey::pattern::SubjectPattern)> = slice
107 .subjects
108 .iter()
109 .enumerate()
110 .filter(|(_, s)| s.class == class)
111 .filter_map(|(i, s)| {
112 zenkey::pattern::SubjectPattern::parse(&s.path)
113 .ok()
114 .map(|p| (i, p))
115 })
116 .collect();
117 let patterns: Vec<zenkey::pattern::SubjectPattern> =
118 candidates.iter().map(|(_, p)| p.clone()).collect();
119 let (winner, binds) = zenkey::pattern::best_match(&patterns, tail)?;
120 let (subject_idx, _) = candidates[winner];
121 Some((
122 &slice.subjects[subject_idx],
123 binds.into_iter().map(|(n, v)| (n.to_string(), v)).collect(),
124 ))
125 }
126
127 pub fn write_cache(&self, dir: &Path) -> Result<()> {
131 std::fs::create_dir_all(dir)?;
132 for (slice, raw) in self.slices.iter().zip(&self.raw) {
133 std::fs::write(dir.join(format!("{}.toml", slice.name)), raw)?;
134 }
135 Ok(())
136 }
137
138 pub fn read_cache(dir: &Path) -> SliceSet {
141 if !dir.is_dir() {
142 return SliceSet::default();
143 }
144 SliceSet::from_dirs(&[dir.to_path_buf()]).unwrap_or_default()
145 }
146}
147
148#[cfg(test)]
149mod tests {
150 use super::*;
151
152 const A: &str = r#"
153 [registry]
154 version = "1.0"
155 app = "t"
156 convention = 1
157 [producer]
158 name = "alpha"
159 [[subject]]
160 path = "flow/{q}"
161 class = "telemetry"
162 type = "Point"
163 [[subject]]
164 path = "flow/special"
165 class = "telemetry"
166 type = "Special"
167 "#;
168
169 #[test]
170 fn refine_uses_shared_precedence() {
171 let mut set = SliceSet::default();
172 set.push(parse_slice(A).unwrap(), A.to_string());
173 let (s, binds) = set
175 .refine("alpha", "telemetry", &["flow", "special"])
176 .unwrap();
177 assert_eq!(s.type_name, "Special");
178 assert!(binds.is_empty());
179 let (s, binds) = set.refine("alpha", "telemetry", &["flow", "p95"]).unwrap();
180 assert_eq!(s.type_name, "Point");
181 assert_eq!(binds, vec![("q".to_string(), "p95".to_string())]);
182 assert!(set.refine("alpha", "state", &["flow", "p95"]).is_none());
183 }
184
185 #[test]
186 fn cache_round_trips_and_last_slice_wins() {
187 let mut set = SliceSet::default();
188 set.push(parse_slice(A).unwrap(), A.to_string());
189 set.push(parse_slice(A).unwrap(), A.to_string());
191 assert_eq!(set.slices().len(), 1);
192
193 let dir = std::env::temp_dir().join(format!("zenkey-fleet-cache-{}", std::process::id()));
194 let _ = std::fs::remove_dir_all(&dir);
195 set.write_cache(&dir).unwrap();
196 let back = SliceSet::read_cache(&dir);
197 assert_eq!(back.slices().len(), 1);
198 assert_eq!(back.get("alpha").unwrap().subjects.len(), 2);
199 let _ = std::fs::remove_dir_all(&dir);
200 assert!(
202 SliceSet::read_cache(Path::new("/nonexistent-zkf"))
203 .slices()
204 .is_empty()
205 );
206 }
207}