1use std::collections::{BTreeMap, BTreeSet};
22use std::time::Duration;
23
24use anyhow::Result;
25use serde::Serialize;
26use zenkey::grammar::{self, ClassOrPlane, SUBJECT_ALIVE, VERSION_CHUNK};
27use zenoh::Session;
28
29use crate::admin::StorageInfo;
30
31const HOST_ALIVE_SWEEP: &str = "**/v1/*/state/*/alive";
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.as_ref().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
124#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
126pub struct DiscoveredBase {
127 pub base: String,
129 pub origins: BTreeSet<String>,
131 pub producers: BTreeSet<String>,
133 pub storages: Vec<String>,
135}
136
137pub fn merge_signals(
140 tokens: impl IntoIterator<Item = AliveToken>,
141 storages: &[StorageInfo],
142) -> Vec<DiscoveredBase> {
143 let mut bases: BTreeMap<String, DiscoveredBase> = BTreeMap::new();
144 fn entry<'m>(
145 bases: &'m mut BTreeMap<String, DiscoveredBase>,
146 base: &str,
147 ) -> &'m mut DiscoveredBase {
148 bases
149 .entry(base.to_string())
150 .or_insert_with(|| DiscoveredBase {
151 base: base.to_string(),
152 ..DiscoveredBase::default()
153 })
154 }
155 for token in tokens {
156 let row = entry(&mut bases, &token.base);
157 let producer = token
159 .producer
160 .unwrap_or_else(|| token.origin.trim_start_matches('@').to_string());
161 row.origins.insert(token.origin);
162 row.producers.insert(producer);
163 }
164 for storage in storages {
165 let Some(base) = base_of_storage(storage) else {
166 continue;
167 };
168 entry(&mut bases, &base)
169 .storages
170 .push(format!("{}@{}", storage.name, storage.zid));
171 }
172 for row in bases.values_mut() {
173 row.storages.sort();
174 row.storages.dedup();
175 }
176 bases.into_values().collect()
177}
178
179pub async fn discover_bases(session: &Session, timeout: Duration) -> Result<Vec<DiscoveredBase>> {
185 let mut tokens = Vec::new();
186 for sweep in [HOST_ALIVE_SWEEP, CATALOG_ALIVE_SWEEP] {
187 let Ok(replies) = session.liveliness().get(sweep).timeout(timeout).await else {
188 continue;
189 };
190 while let Ok(reply) = replies.recv_async().await {
191 let Ok(sample) = reply.result() else { continue };
192 if let Some(token) = parse_alive_key(sample.key_expr().as_str()) {
193 tokens.push(token);
194 }
195 }
196 }
197 let storages = crate::admin::storages(session, timeout)
198 .await
199 .unwrap_or_default();
200 Ok(merge_signals(tokens, &storages))
201}
202
203#[cfg(test)]
204mod tests {
205 use super::*;
206
207 fn token(base: &str, origin: &str, producer: Option<&str>) -> AliveToken {
208 AliveToken {
209 base: base.into(),
210 origin: origin.into(),
211 producer: producer.map(str::to_string),
212 }
213 }
214
215 #[test]
216 fn alive_keys_attribute_by_fixed_arity() {
217 assert_eq!(
219 parse_alive_key("v1/h-3fa9c2d41b7e/state/sysinfo/alive"),
220 Some(token("", "h-3fa9c2d41b7e", Some("sysinfo")))
221 );
222 assert_eq!(
223 parse_alive_key("zensight/v1/h-3fa9c2d41b7e/state/sysinfo/alive"),
224 Some(token("zensight", "h-3fa9c2d41b7e", Some("sysinfo")))
225 );
226 assert_eq!(
227 parse_alive_key("acme/fleet-a/v1/h-aaaaaaaaaaaa/state/netring/alive"),
228 Some(token("acme/fleet-a", "h-aaaaaaaaaaaa", Some("netring")))
229 );
230 assert_eq!(
233 parse_alive_key("acme/v1/v1/h-3fa9c2d41b7e/state/tc/alive"),
234 Some(token("acme/v1", "h-3fa9c2d41b7e", Some("tc")))
235 );
236 assert_eq!(
238 parse_alive_key("v1/@catalog/state/alive"),
239 Some(token("", "@catalog", None))
240 );
241 assert_eq!(
242 parse_alive_key("acme/v1/v1/@catalog/state/alive"),
243 Some(token("acme/v1", "@catalog", None))
244 );
245 }
246
247 #[test]
248 fn alive_key_rejects_foreign_shapes() {
249 for key in [
250 "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", ] {
258 assert_eq!(parse_alive_key(key), None, "{key}");
259 }
260 }
261
262 #[test]
263 fn catalog_sweep_pins_to_the_typed_builder() {
264 assert_eq!(
265 CATALOG_ALIVE_SWEEP,
266 format!(
267 "**/{}",
268 zenkey::selector::service_alive(&zenkey::ServiceOrigin::catalog())
269 )
270 );
271 assert_eq!(
272 HOST_ALIVE_SWEEP,
273 format!(
274 "**/{}",
275 zenkey::selector::all_liveliness(zenkey::selector::Scope::fleet())
276 )
277 );
278 }
279
280 fn storage(strip_prefix: Option<&str>, key_expr: Option<&str>) -> StorageInfo {
281 StorageInfo {
282 zid: "z1".into(),
283 name: "latest".into(),
284 key_expr: key_expr.map(str::to_string),
285 raw: match strip_prefix {
286 Some(p) => serde_json::json!({ "strip_prefix": p }),
287 None => serde_json::Value::Null,
288 },
289 }
290 }
291
292 #[test]
293 fn storage_bases_prefer_strip_prefix() {
294 assert_eq!(
295 base_of_storage(&storage(Some("zensight/v1"), None)),
296 Some("zensight".into())
297 );
298 assert_eq!(base_of_storage(&storage(Some("v1"), None)), Some("".into()));
299 assert_eq!(
301 base_of_storage(&storage(Some("acme/fleet-a/v1"), Some("other/v1/**"))),
302 Some("acme/fleet-a".into())
303 );
304 assert_eq!(
306 base_of_storage(&storage(Some("zensight"), Some("zensight/v1/*/state/**"))),
307 Some("zensight".into())
308 );
309 assert_eq!(
311 base_of_storage(&storage(None, Some("acme/fleet-a/v1/*/state/**"))),
312 Some("acme/fleet-a".into())
313 );
314 assert_eq!(
315 base_of_storage(&storage(None, Some("v1/*/state/**"))),
316 Some("".into())
317 );
318 assert_eq!(base_of_storage(&storage(None, Some("**"))), None);
320 assert_eq!(base_of_storage(&storage(None, Some("*/v1/**"))), None);
321 assert_eq!(base_of_storage(&storage(None, None)), None);
322 }
323
324 #[test]
325 fn signals_merge_sorted_and_deduped() {
326 let tokens = vec![
327 token("zensight", "h-3fa9c2d41b7e", Some("sysinfo")),
328 token("zensight", "h-aaaaaaaaaaaa", Some("sysinfo")), token("zensight", "@catalog", None), token("", "h-3fa9c2d41b7e", Some("tc")),
331 ];
332 let storages = [
333 storage(Some("zensight/v1"), None),
334 storage(Some("zensight/v1"), None), storage(Some("acme/v1"), None), ];
337 let rows = merge_signals(tokens, &storages);
338 let bases: Vec<&str> = rows.iter().map(|r| r.base.as_str()).collect();
340 assert_eq!(bases, vec!["", "acme", "zensight"]);
341 let zs = &rows[2];
342 assert_eq!(zs.origins.len(), 3);
343 assert_eq!(
344 zs.producers.iter().collect::<Vec<_>>(),
345 vec!["catalog", "sysinfo"]
346 );
347 assert_eq!(zs.storages, vec!["latest@z1"]);
348 let acme = &rows[1];
349 assert!(acme.origins.is_empty(), "storage-only base has no origins");
350 assert_eq!(acme.storages, vec!["latest@z1"]);
351 }
352}