use std::collections::BTreeMap;
use std::time::Duration;
use anyhow::Result;
use zenkey::grammar::with_base;
use zenoh::Session;
pub async fn roster(
session: &Session,
base: &str,
timeout: Duration,
) -> Result<BTreeMap<String, Vec<String>>> {
let mut out: BTreeMap<String, Vec<String>> = BTreeMap::new();
let catalog_alive = zenkey::selector::service_alive(&zenkey::ServiceOrigin::catalog());
for expr in [
with_base(
base,
zenkey::selector::all_liveliness(zenkey::selector::Scope::fleet()),
),
with_base(base, catalog_alive),
] {
let Ok(replies) = session.liveliness().get(&expr).timeout(timeout).await else {
continue;
};
while let Ok(reply) = replies.recv_async().await {
let Ok(sample) = reply.result() else { continue };
let key = sample.key_expr().as_str();
let Some(parsed) = zenkey::grammar::parse_full(base, key) else {
continue;
};
let origin = parsed.origin.chunk().to_string();
let producer = parsed
.producer
.as_ref()
.map(|p| p.chunk())
.unwrap_or_else(|| origin.trim_start_matches('@').to_string());
out.entry(origin).or_default().push(producer);
}
}
for producers in out.values_mut() {
producers.sort();
producers.dedup();
}
Ok(out)
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct ProducerInfo {
pub name: String,
pub alive: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub app: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub registry_version: Option<String>,
pub subjects: usize,
pub procedures: usize,
#[serde(skip_serializing_if = "Vec::is_empty", default)]
pub blob_tiers: Vec<String>,
#[serde(skip_serializing_if = "Vec::is_empty", default)]
pub media: Vec<MediaStreamInfo>,
pub deprecated_served: usize,
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct MediaStreamInfo {
pub path: String,
pub encoding: String,
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct Freshness {
pub producer: String,
pub path: String,
pub ttl_s: i64,
#[serde(skip_serializing_if = "Option::is_none")]
pub age_s: Option<i64>,
pub stale: bool,
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct NodeInfo {
pub origin: String,
pub producers: Vec<ProducerInfo>,
#[serde(skip_serializing_if = "Vec::is_empty", default)]
pub freshness: Vec<Freshness>,
}
enum Node {
Host(zenkey::origin::RemoteOrigin),
Service(zenkey::ServiceOrigin),
}
impl Node {
fn parse(origin: &str) -> Result<Node> {
if origin.starts_with('@') {
zenkey::ServiceOrigin::new(origin)
.map(Node::Service)
.map_err(|e| anyhow::anyhow!("{e}"))
} else {
zenkey::origin::RemoteOrigin::parse(origin)
.map(Node::Host)
.map_err(|e| anyhow::anyhow!("{e} — a hostname is not an origin (RFC 06 §6)"))
}
}
fn alive_selector(&self) -> String {
match self {
Node::Host(o) => {
zenkey::selector::all_liveliness(zenkey::selector::Scope::origin(o)).to_string()
}
Node::Service(o) => zenkey::selector::service_alive(o).to_string(),
}
}
fn introspect_selector(&self) -> String {
match self {
Node::Host(o) => {
zenkey::selector::rpc(zenkey::selector::Scope::origin(o), "*", &["introspect"])
.to_string()
}
Node::Service(o) => zenkey::selector::service_rpc(o, &["introspect"]).to_string(),
}
}
fn state_selector(&self) -> String {
let scope = match self {
Node::Host(o) => zenkey::selector::Scope::origin(o),
Node::Service(o) => zenkey::selector::Scope::origin(o),
};
zenkey::selector::all_state(scope).to_string()
}
}
pub async fn node_info(
session: &Session,
base: &str,
origin: &str,
timeout: Duration,
with_freshness: bool,
) -> Result<NodeInfo> {
let node = Node::parse(origin)?;
let mut alive: Vec<String> = Vec::new();
let alive_expr = with_base(base, node.alive_selector());
if let Ok(replies) = session.liveliness().get(&alive_expr).timeout(timeout).await {
while let Ok(reply) = replies.recv_async().await {
let Ok(sample) = reply.result() else { continue };
let Some(parsed) = zenkey::grammar::parse_full(base, sample.key_expr().as_str()) else {
continue;
};
alive.push(
parsed
.producer
.as_ref()
.map(|p| p.chunk())
.unwrap_or_else(|| parsed.origin.chunk().trim_start_matches('@').to_string()),
);
}
}
alive.sort();
alive.dedup();
let introspect = with_base(base, node.introspect_selector());
let answers = crate::query::fleet_get(session, base, &introspect, None, timeout)
.await
.unwrap_or_default();
let served: Vec<zenkey::slice::RegistrySlice> = answers
.into_iter()
.filter(|a| a.origin == origin)
.filter_map(|a| {
let crate::query::Answer::Value(bytes) = a.answer else {
return None;
};
let toml = String::from_utf8_lossy(&bytes.to_bytes()).to_string();
match zenkey::parse_slice(&toml) {
Ok(slice) => Some(slice),
Err(e) => {
tracing::warn!(origin, "introspect reply did not parse, skipping: {e}");
None
}
}
})
.collect();
let mine: Vec<&zenkey::slice::RegistrySlice> = served.iter().collect();
let mut names: Vec<String> = alive.clone();
names.extend(mine.iter().map(|s| s.name.clone()));
names.sort();
names.dedup();
let producers: Vec<ProducerInfo> = names
.iter()
.map(|name| {
let slice = mine.iter().find(|s| &s.name == name);
ProducerInfo {
name: name.clone(),
alive: alive.iter().any(|a| a == name),
app: slice.map(|s| s.app.clone()),
registry_version: slice.map(|s| s.version.clone()),
subjects: slice.map(|s| s.subjects.len()).unwrap_or(0),
procedures: slice.map(|s| s.procedures.len()).unwrap_or(0),
blob_tiers: slice
.map(|s| s.blob.iter().map(|b| b.tier.clone()).collect())
.unwrap_or_default(),
media: slice
.map(|s| {
s.media
.iter()
.map(|m| MediaStreamInfo {
path: m.path.clone(),
encoding: m.encoding.clone(),
})
.collect()
})
.unwrap_or_default(),
deprecated_served: slice.map(|s| s.deprecated.len()).unwrap_or(0),
}
})
.collect();
let mut freshness = Vec::new();
if with_freshness && !mine.is_empty() {
let selector = with_base(base, node.state_selector());
let samples = crate::query::state_snapshot(session, &selector, timeout, None)
.await
.unwrap_or_default();
let now = std::time::SystemTime::now();
for slice in &mine {
for subject in &slice.subjects {
let Some(ttl) = subject.ttl_s else { continue };
if subject.class != "state" {
continue;
}
let age = samples
.iter()
.filter_map(|s| {
let parsed = zenkey::grammar::parse_full(base, &s.key)?;
let p = parsed.producer.as_ref()?.name().to_string();
if p != slice.name {
return None;
}
let tail: Vec<&str> = parsed.subject.clone();
let pattern = zenkey::pattern::SubjectPattern::parse(&subject.path).ok()?;
pattern.matches(&tail)?;
s.timestamp.map(|t| {
now.duration_since(t.get_time().to_system_time())
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
})
})
.min();
freshness.push(Freshness {
producer: slice.name.clone(),
path: subject.path.clone(),
ttl_s: ttl,
age_s: age,
stale: match age {
Some(a) => a > ttl,
None => true,
},
});
}
}
}
Ok(NodeInfo {
origin: origin.to_string(),
producers,
freshness,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BridgeMatch {
pub host_id: zenkey::origin::HostId,
pub source: String,
pub key: String,
}
pub async fn bridge_resolve(
session: &zenoh::Session,
base: &str,
producer: &str,
label: &str,
timeout: std::time::Duration,
) -> Result<(Vec<BridgeMatch>, usize)> {
let relative =
zenkey::selector::producer_state(zenkey::selector::Scope::fleet(), producer, &["health"])
.to_string();
let key = zenkey::grammar::with_base(base, relative);
let answers = crate::query::fleet_get(session, base, &key, None, timeout).await?;
let mut matches = Vec::new();
let seen = answers.len();
for a in &answers {
let crate::query::Answer::Value(bytes) = &a.answer else {
continue;
};
let Ok(doc) = serde_json::from_slice::<serde_json::Value>(&bytes.to_bytes()) else {
continue;
};
let (Some(host_id), Some(source)) = (
doc.get("host_id").and_then(|v| v.as_str()),
doc.get("source").and_then(|v| v.as_str()),
) else {
continue;
};
if source == label
&& let Ok(id) = zenkey::origin::HostId::parse(host_id)
{
matches.push(BridgeMatch {
host_id: id,
source: source.to_string(),
key: a.key.clone(),
});
}
}
matches.dedup_by(|a, b| a.host_id == b.host_id);
Ok((matches, seen))
}