use std::collections::{BTreeMap, BTreeSet};
use std::time::Duration;
use anyhow::Result;
use serde::Serialize;
use zenkey::grammar::{self, ClassOrPlane, SUBJECT_ALIVE, VERSION_CHUNK};
use zenoh::Session;
use crate::admin::StorageInfo;
const HOST_ALIVE_SWEEP: &str = "**/v1/*/state/*/alive";
const CATALOG_ALIVE_SWEEP: &str = "**/v1/@catalog/state/alive";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AliveToken {
pub base: String,
pub origin: String,
pub producer: Option<String>,
}
pub fn parse_alive_key(key: &str) -> Option<AliveToken> {
let chunks: Vec<&str> = key.split('/').collect();
for tail_len in [5usize, 4] {
let Some(split) = chunks.len().checked_sub(tail_len) else {
continue;
};
if chunks[split] != VERSION_CHUNK {
continue;
}
let tail = chunks[split..].join("/");
let Ok(parsed) = grammar::parse(&tail) else {
continue;
};
if !matches!(parsed.class, ClassOrPlane::Class(grammar::Class::State))
|| parsed.subject != [SUBJECT_ALIVE]
{
continue;
}
if (tail_len == 5) != parsed.producer.is_some() {
continue;
}
return Some(AliveToken {
base: chunks[..split].join("/"),
origin: parsed.origin.chunk().to_string(),
producer: parsed.producer.as_ref().map(|p| p.chunk()),
});
}
None
}
pub fn base_of_storage(storage: &StorageInfo) -> Option<String> {
if let Some(prefix) = storage.raw.get("strip_prefix").and_then(|v| v.as_str()) {
if prefix == VERSION_CHUNK {
return Some(String::new());
}
if let Some(base) = prefix.strip_suffix("/v1") {
return Some(base.to_string());
}
}
let key_expr = storage.key_expr.as_deref()?;
let mut base_chunks: Vec<&str> = Vec::new();
for chunk in key_expr.split('/') {
if chunk == VERSION_CHUNK {
return Some(base_chunks.join("/"));
}
if chunk.contains('*') || chunk.starts_with('@') {
return None;
}
base_chunks.push(chunk);
}
None
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
pub struct DiscoveredBase {
pub base: String,
pub origins: BTreeSet<String>,
pub producers: BTreeSet<String>,
pub storages: Vec<String>,
}
pub fn merge_signals(
tokens: impl IntoIterator<Item = AliveToken>,
storages: &[StorageInfo],
) -> Vec<DiscoveredBase> {
let mut bases: BTreeMap<String, DiscoveredBase> = BTreeMap::new();
fn entry<'m>(
bases: &'m mut BTreeMap<String, DiscoveredBase>,
base: &str,
) -> &'m mut DiscoveredBase {
bases
.entry(base.to_string())
.or_insert_with(|| DiscoveredBase {
base: base.to_string(),
..DiscoveredBase::default()
})
}
for token in tokens {
let row = entry(&mut bases, &token.base);
let producer = token
.producer
.unwrap_or_else(|| token.origin.trim_start_matches('@').to_string());
row.origins.insert(token.origin);
row.producers.insert(producer);
}
for storage in storages {
let Some(base) = base_of_storage(storage) else {
continue;
};
entry(&mut bases, &base)
.storages
.push(format!("{}@{}", storage.name, storage.zid));
}
for row in bases.values_mut() {
row.storages.sort();
row.storages.dedup();
}
bases.into_values().collect()
}
pub async fn discover_bases(session: &Session, timeout: Duration) -> Result<Vec<DiscoveredBase>> {
let mut tokens = Vec::new();
for sweep in [HOST_ALIVE_SWEEP, CATALOG_ALIVE_SWEEP] {
let Ok(replies) = session.liveliness().get(sweep).timeout(timeout).await else {
continue;
};
while let Ok(reply) = replies.recv_async().await {
let Ok(sample) = reply.result() else { continue };
if let Some(token) = parse_alive_key(sample.key_expr().as_str()) {
tokens.push(token);
}
}
}
let storages = crate::admin::storages(session, timeout)
.await
.unwrap_or_default();
Ok(merge_signals(tokens, &storages))
}
#[cfg(test)]
mod tests {
use super::*;
fn token(base: &str, origin: &str, producer: Option<&str>) -> AliveToken {
AliveToken {
base: base.into(),
origin: origin.into(),
producer: producer.map(str::to_string),
}
}
#[test]
fn alive_keys_attribute_by_fixed_arity() {
assert_eq!(
parse_alive_key("v1/h-3fa9c2d41b7e/state/sysinfo/alive"),
Some(token("", "h-3fa9c2d41b7e", Some("sysinfo")))
);
assert_eq!(
parse_alive_key("zensight/v1/h-3fa9c2d41b7e/state/sysinfo/alive"),
Some(token("zensight", "h-3fa9c2d41b7e", Some("sysinfo")))
);
assert_eq!(
parse_alive_key("acme/fleet-a/v1/h-aaaaaaaaaaaa/state/netring/alive"),
Some(token("acme/fleet-a", "h-aaaaaaaaaaaa", Some("netring")))
);
assert_eq!(
parse_alive_key("acme/v1/v1/h-3fa9c2d41b7e/state/tc/alive"),
Some(token("acme/v1", "h-3fa9c2d41b7e", Some("tc")))
);
assert_eq!(
parse_alive_key("v1/@catalog/state/alive"),
Some(token("", "@catalog", None))
);
assert_eq!(
parse_alive_key("acme/v1/v1/@catalog/state/alive"),
Some(token("acme/v1", "@catalog", None))
);
}
#[test]
fn alive_key_rejects_foreign_shapes() {
for key in [
"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", ] {
assert_eq!(parse_alive_key(key), None, "{key}");
}
}
#[test]
fn catalog_sweep_pins_to_the_typed_builder() {
assert_eq!(
CATALOG_ALIVE_SWEEP,
format!(
"**/{}",
zenkey::selector::service_alive(&zenkey::ServiceOrigin::catalog())
)
);
assert_eq!(
HOST_ALIVE_SWEEP,
format!(
"**/{}",
zenkey::selector::all_liveliness(zenkey::selector::Scope::fleet())
)
);
}
fn storage(strip_prefix: Option<&str>, key_expr: Option<&str>) -> StorageInfo {
StorageInfo {
zid: "z1".into(),
name: "latest".into(),
key_expr: key_expr.map(str::to_string),
strip_prefix: strip_prefix.map(str::to_string),
volume: None,
raw: match strip_prefix {
Some(p) => serde_json::json!({ "strip_prefix": p }),
None => serde_json::Value::Null,
},
}
}
#[test]
fn storage_bases_prefer_strip_prefix() {
assert_eq!(
base_of_storage(&storage(Some("zensight/v1"), None)),
Some("zensight".into())
);
assert_eq!(base_of_storage(&storage(Some("v1"), None)), Some("".into()));
assert_eq!(
base_of_storage(&storage(Some("acme/fleet-a/v1"), Some("other/v1/**"))),
Some("acme/fleet-a".into())
);
assert_eq!(
base_of_storage(&storage(Some("zensight"), Some("zensight/v1/*/state/**"))),
Some("zensight".into())
);
assert_eq!(
base_of_storage(&storage(None, Some("acme/fleet-a/v1/*/state/**"))),
Some("acme/fleet-a".into())
);
assert_eq!(
base_of_storage(&storage(None, Some("v1/*/state/**"))),
Some("".into())
);
assert_eq!(base_of_storage(&storage(None, Some("**"))), None);
assert_eq!(base_of_storage(&storage(None, Some("*/v1/**"))), None);
assert_eq!(base_of_storage(&storage(None, None)), None);
}
#[test]
fn signals_merge_sorted_and_deduped() {
let tokens = vec![
token("zensight", "h-3fa9c2d41b7e", Some("sysinfo")),
token("zensight", "h-aaaaaaaaaaaa", Some("sysinfo")), token("zensight", "@catalog", None), token("", "h-3fa9c2d41b7e", Some("tc")),
];
let storages = [
storage(Some("zensight/v1"), None),
storage(Some("zensight/v1"), None), storage(Some("acme/v1"), None), ];
let rows = merge_signals(tokens, &storages);
let bases: Vec<&str> = rows.iter().map(|r| r.base.as_str()).collect();
assert_eq!(bases, vec!["", "acme", "zensight"]);
let zs = &rows[2];
assert_eq!(zs.origins.len(), 3);
assert_eq!(
zs.producers.iter().collect::<Vec<_>>(),
vec!["catalog", "sysinfo"]
);
assert_eq!(zs.storages, vec!["latest@z1"]);
let acme = &rows[1];
assert!(acme.origins.is_empty(), "storage-only base has no origins");
assert_eq!(acme.storages, vec!["latest@z1"]);
}
}