1use std::collections::BTreeMap;
22use std::time::Duration;
23
24use crate::Result;
25use zenkey::grammar::{self, ClassOrPlane, SUBJECT_ALIVE, VERSION_CHUNK};
26use zenoh::Session;
27
28use crate::report::{DiscoveredBase, StorageInfo};
29
30const HOST_ALIVE_SWEEP: &str = "**/v1/*/state/*/alive";
34
35const CATALOG_ALIVE_SWEEP: &str = "**/v1/@catalog/state/alive";
38
39#[derive(Debug, Clone, PartialEq, Eq)]
42pub struct AliveToken {
43 pub base: String,
45 pub origin: String,
47 pub producer: Option<String>,
49}
50
51pub fn parse_alive_key(key: &str) -> Option<AliveToken> {
60 let chunks: Vec<&str> = key.split('/').collect();
61 for tail_len in [5usize, 4] {
63 let Some(split) = chunks.len().checked_sub(tail_len) else {
64 continue;
65 };
66 if chunks[split] != VERSION_CHUNK {
67 continue;
68 }
69 let tail = chunks[split..].join("/");
70 let Ok(parsed) = grammar::parse(&tail) else {
71 continue;
72 };
73 if !matches!(parsed.class, ClassOrPlane::Class(grammar::Class::State))
74 || parsed.subject != [SUBJECT_ALIVE]
75 {
76 continue;
77 }
78 if (tail_len == 5) != parsed.producer().is_some() {
81 continue;
82 }
83 return Some(AliveToken {
84 base: chunks[..split].join("/"),
85 origin: parsed.origin.chunk().to_string(),
86 producer: parsed.producer().map(|p| p.chunk()),
87 });
88 }
89 None
90}
91
92pub fn base_of_storage(storage: &StorageInfo) -> Option<String> {
100 if let Some(prefix) = storage.raw.get("strip_prefix").and_then(|v| v.as_str()) {
101 if prefix == VERSION_CHUNK {
102 return Some(String::new());
103 }
104 if let Some(base) = prefix.strip_suffix("/v1") {
105 return Some(base.to_string());
106 }
107 }
110 let key_expr = storage.key_expr.as_deref()?;
111 let mut base_chunks: Vec<&str> = Vec::new();
112 for chunk in key_expr.split('/') {
113 if chunk == VERSION_CHUNK {
114 return Some(base_chunks.join("/"));
115 }
116 if chunk.contains('*') || chunk.starts_with('@') {
117 return None;
118 }
119 base_chunks.push(chunk);
120 }
121 None
122}
123
124pub fn merge_signals(
127 tokens: impl IntoIterator<Item = AliveToken>,
128 storages: &[StorageInfo],
129) -> Vec<DiscoveredBase> {
130 let mut bases: BTreeMap<String, DiscoveredBase> = BTreeMap::new();
131
132 fn entry<'m>(
133 bases: &'m mut BTreeMap<String, DiscoveredBase>,
134 base: &str,
135 ) -> &'m mut DiscoveredBase {
136 bases
137 .entry(base.to_string())
138 .or_insert_with(|| DiscoveredBase {
139 base: base.to_string(),
140 ..DiscoveredBase::default()
141 })
142 }
143 for token in tokens {
144 let row = entry(&mut bases, &token.base);
145 let producer = token
147 .producer
148 .unwrap_or_else(|| token.origin.trim_start_matches('@').to_string());
149 row.origins.insert(token.origin);
150 row.producers.insert(producer);
151 }
152 for storage in storages {
153 let Some(base) = base_of_storage(storage) else {
154 continue;
155 };
156 entry(&mut bases, &base)
157 .storages
158 .push(format!("{}@{}", storage.name, storage.zid));
159 }
160 for row in bases.values_mut() {
161 row.storages.sort();
162 row.storages.dedup();
163 }
164 bases.into_values().collect()
165}
166
167pub async fn discover_bases(session: &Session, timeout: Duration) -> Result<Vec<DiscoveredBase>> {
177 let mut tokens = Vec::new();
178 for sweep in [HOST_ALIVE_SWEEP, CATALOG_ALIVE_SWEEP] {
179 let Ok(replies) = session.liveliness().get(sweep).timeout(timeout).await else {
180 continue;
181 };
182 while let Ok(reply) = replies.recv_async().await {
183 let Ok(sample) = reply.result() else { continue };
184 if let Some(token) = parse_alive_key(sample.key_expr().as_str()) {
185 tokens.push(token);
186 }
187 }
188 }
189 let storages = crate::bus::admin::storages(session, timeout)
190 .await
191 .unwrap_or_default();
192 Ok(merge_signals(tokens, &storages))
193}
194
195#[cfg(test)]
196mod tests {
197 use super::*;
198
199 fn token(base: &str, origin: &str, producer: Option<&str>) -> AliveToken {
200 AliveToken {
201 base: base.into(),
202 origin: origin.into(),
203 producer: producer.map(str::to_string),
204 }
205 }
206
207 #[test]
208 fn alive_keys_attribute_by_fixed_arity() {
209 assert_eq!(
211 parse_alive_key("v1/h-3fa9c2d41b7e/state/sysinfo/alive"),
212 Some(token("", "h-3fa9c2d41b7e", Some("sysinfo")))
213 );
214 assert_eq!(
215 parse_alive_key("zensight/v1/h-3fa9c2d41b7e/state/sysinfo/alive"),
216 Some(token("zensight", "h-3fa9c2d41b7e", Some("sysinfo")))
217 );
218 assert_eq!(
219 parse_alive_key("acme/fleet-a/v1/h-aaaaaaaaaaaa/state/netring/alive"),
220 Some(token("acme/fleet-a", "h-aaaaaaaaaaaa", Some("netring")))
221 );
222 assert_eq!(
225 parse_alive_key("acme/v1/v1/h-3fa9c2d41b7e/state/tc/alive"),
226 Some(token("acme/v1", "h-3fa9c2d41b7e", Some("tc")))
227 );
228 assert_eq!(
230 parse_alive_key("v1/@catalog/state/alive"),
231 Some(token("", "@catalog", None))
232 );
233 assert_eq!(
234 parse_alive_key("acme/v1/v1/@catalog/state/alive"),
235 Some(token("acme/v1", "@catalog", None))
236 );
237 }
238
239 #[test]
240 fn alive_key_rejects_foreign_shapes() {
241 for key in [
242 "other/junk/alive", "alive", "zensight/v1/notanorigin/state/p/alive", "zensight/v1/h-3fa9c2d41b7e/telemetry/p/alive", "zensight/v2/h-3fa9c2d41b7e/state/p/alive", "zensight/v1/h-3fa9c2d41b7e/state/p/health", "zensight/v1/h-3fa9c2d41b7e/state/p/device/d0/alive", ] {
250 assert_eq!(parse_alive_key(key), None, "{key}");
251 }
252 }
253
254 #[test]
255 fn catalog_sweep_pins_to_the_typed_builder() {
256 assert_eq!(
257 CATALOG_ALIVE_SWEEP,
258 format!(
259 "**/{}",
260 zenkey::selector::service_alive(&zenkey::ServiceOrigin::catalog())
261 )
262 );
263 assert_eq!(
264 HOST_ALIVE_SWEEP,
265 format!(
266 "**/{}",
267 zenkey::selector::all_liveliness(zenkey::selector::Scope::fleet())
268 )
269 );
270 }
271
272 fn storage(strip_prefix: Option<&str>, key_expr: Option<&str>) -> StorageInfo {
273 StorageInfo {
274 zid: "z1".into(),
275 name: "latest".into(),
276 key_expr: key_expr.map(str::to_string),
277 strip_prefix: strip_prefix.map(str::to_string),
278 volume: None,
279 raw: match strip_prefix {
280 Some(p) => serde_json::json!({ "strip_prefix": p }),
281 None => serde_json::Value::Null,
282 },
283 }
284 }
285
286 #[test]
287 fn storage_bases_prefer_strip_prefix() {
288 assert_eq!(
289 base_of_storage(&storage(Some("zensight/v1"), None)),
290 Some("zensight".into())
291 );
292 assert_eq!(base_of_storage(&storage(Some("v1"), None)), Some("".into()));
293 assert_eq!(
295 base_of_storage(&storage(Some("acme/fleet-a/v1"), Some("other/v1/**"))),
296 Some("acme/fleet-a".into())
297 );
298 assert_eq!(
300 base_of_storage(&storage(Some("zensight"), Some("zensight/v1/*/state/**"))),
301 Some("zensight".into())
302 );
303 assert_eq!(
305 base_of_storage(&storage(None, Some("acme/fleet-a/v1/*/state/**"))),
306 Some("acme/fleet-a".into())
307 );
308 assert_eq!(
309 base_of_storage(&storage(None, Some("v1/*/state/**"))),
310 Some("".into())
311 );
312 assert_eq!(base_of_storage(&storage(None, Some("**"))), None);
314 assert_eq!(base_of_storage(&storage(None, Some("*/v1/**"))), None);
315 assert_eq!(base_of_storage(&storage(None, None)), None);
316 }
317
318 #[test]
319 fn signals_merge_sorted_and_deduped() {
320 let tokens = vec![
321 token("zensight", "h-3fa9c2d41b7e", Some("sysinfo")),
322 token("zensight", "h-aaaaaaaaaaaa", Some("sysinfo")), token("zensight", "@catalog", None), token("", "h-3fa9c2d41b7e", Some("tc")),
325 ];
326 let storages = [
327 storage(Some("zensight/v1"), None),
328 storage(Some("zensight/v1"), None), storage(Some("acme/v1"), None), ];
331 let rows = merge_signals(tokens, &storages);
332 let bases: Vec<&str> = rows.iter().map(|r| r.base.as_str()).collect();
334 assert_eq!(bases, vec!["", "acme", "zensight"]);
335 let zs = &rows[2];
336 assert_eq!(zs.origins.len(), 3);
337 assert_eq!(
338 zs.producers.iter().collect::<Vec<_>>(),
339 vec!["catalog", "sysinfo"]
340 );
341 assert_eq!(zs.storages, vec!["latest@z1"]);
342 let acme = &rows[1];
343 assert!(acme.origins.is_empty(), "storage-only base has no origins");
344 assert_eq!(acme.storages, vec!["latest@z1"]);
345 }
346}