use std::time::Duration;
use anyhow::{Context, Result};
use zenkey::grammar::with_base;
use zenkey::{RegistrySlice, parse_slice};
use zenoh::Session;
use zenoh::query::{ConsolidationMode, QueryTarget};
pub enum Answer {
Value(zenoh::bytes::ZBytes),
Error { name: String, message: String },
}
pub struct FleetAnswer {
pub origin: String,
pub answer: Answer,
}
pub async fn fleet_get(
session: &Session,
base: &str,
key: &str,
payload: Option<Vec<u8>>,
timeout: Duration,
) -> Result<Vec<FleetAnswer>> {
let mut builder = session
.get(key)
.target(QueryTarget::All)
.consolidation(ConsolidationMode::None)
.timeout(timeout);
if let Some(body) = payload {
builder = builder.payload(body);
}
let replies = builder
.await
.map_err(|e| anyhow::anyhow!("{e}"))
.with_context(|| format!("query failed: {key}"))?;
let mut out = Vec::new();
while let Ok(reply) = replies.recv_async().await {
match reply.result() {
Ok(sample) => {
let origin = origin_of(base, sample.key_expr().as_str());
out.push(FleetAnswer {
origin,
answer: Answer::Value(sample.payload().clone()),
});
}
Err(err) => {
let bytes = err.payload().to_bytes();
let (name, message) = match serde_json::from_slice::<serde_json::Value>(&bytes) {
Ok(v) => (
v.get("error")
.and_then(|e| e.as_str())
.unwrap_or("error/unparsed")
.to_string(),
v.get("message")
.and_then(|m| m.as_str())
.unwrap_or_default()
.to_string(),
),
Err(_) => (
"error/unparsed".to_string(),
String::from_utf8_lossy(&bytes).to_string(),
),
};
out.push(FleetAnswer {
origin: "?".to_string(),
answer: Answer::Error { name, message },
});
}
}
}
Ok(out)
}
fn origin_of(base: &str, key: &str) -> String {
zenkey::grammar::parse_full(base, key)
.map(|k| k.origin.chunk().to_string())
.unwrap_or_else(|| "?".to_string())
}
pub async fn fleet_registry(
session: &Session,
base: &str,
timeout: Duration,
) -> Result<Vec<(String, RegistrySlice)>> {
Ok(fleet_registry_raw(session, base, timeout)
.await?
.into_iter()
.map(|(slice, _)| (slice.name.clone(), slice))
.collect())
}
pub async fn fleet_registry_raw(
session: &Session,
base: &str,
timeout: Duration,
) -> Result<Vec<(RegistrySlice, String)>> {
let keys = [
with_base(base, zenkey::selector::fleet_rpc("*", &["introspect"])),
with_base(
base,
zenkey::selector::service_rpc(&zenkey::ServiceOrigin::catalog(), &["introspect"]),
),
];
let mut slices = Vec::new();
for key in keys {
let answers = fleet_get(session, base, &key, None, timeout).await?;
for answer in answers {
let Answer::Value(bytes) = answer.answer else {
continue;
};
let served_toml = String::from_utf8_lossy(&bytes.to_bytes()).to_string();
match parse_slice(&served_toml) {
Ok(slice) => slices.push((slice, served_toml)),
Err(e) => tracing::warn!(
origin = %answer.origin,
"introspect reply did not parse, skipping: {e}"
),
}
}
}
Ok(slices)
}
#[derive(Debug, Clone)]
pub struct StateSample {
pub key: String,
pub timestamp: Option<zenoh::time::Timestamp>,
pub payload_len: usize,
}
pub async fn state_snapshot(
session: &Session,
selector: &str,
timeout: Duration,
) -> Result<Vec<StateSample>> {
let replies = session
.get(selector)
.target(QueryTarget::All)
.consolidation(ConsolidationMode::None)
.timeout(timeout)
.await
.map_err(|e| anyhow::anyhow!("{e}"))
.with_context(|| format!("state snapshot failed: {selector}"))?;
let mut out = Vec::new();
while let Ok(reply) = replies.recv_async().await {
let Ok(sample) = reply.result() else { continue };
out.push(StateSample {
key: sample.key_expr().as_str().to_string(),
timestamp: sample.timestamp().copied(),
payload_len: sample.payload().len(),
});
}
Ok(out)
}