use std::collections::BTreeMap;
use std::time::Duration;
use crate::report::{Freshness, MediaStreamInfo, NodeInfo, ProducerInfo};
use crate::{Error, Result};
use zenkey::grammar::with_base;
pub async fn roster(
fleet: &crate::Fleet<'_>,
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 [
fleet.wire(zenkey::selector::all_liveliness(
zenkey::selector::Scope::fleet(),
)),
fleet.wire(catalog_alive),
] {
let Ok(replies) = fleet
.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((origin, producer)) = token_identity(fleet.base(), key) else {
continue;
};
out.entry(origin).or_default().push(producer);
}
}
for producers in out.values_mut() {
producers.sort();
producers.dedup();
}
Ok(out)
}
pub struct RosterWatch {
monitor: crate::Monitor,
events: crate::EventStream,
roster: BTreeMap<String, Vec<String>>,
base: String,
pending: RosterChange,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct RosterChange {
pub node_up: bool,
pub node_down: bool,
}
const BURST_QUIET: Duration = Duration::from_millis(50);
impl RosterWatch {
pub async fn start(fleet: &crate::Fleet<'_>, timeout: Duration) -> Result<RosterWatch> {
let liveliness = vec![
fleet.wire(zenkey::selector::all_liveliness(
zenkey::selector::Scope::fleet(),
)),
fleet.wire(zenkey::selector::service_alive(
&zenkey::ServiceOrigin::catalog(),
)),
];
let monitor = crate::Monitor::start(
fleet.session(),
crate::MonitorSpec {
selectors: vec![],
liveliness,
..Default::default()
},
)
.await?;
let events = monitor.events();
let roster = roster(fleet, timeout).await?;
Ok(RosterWatch {
monitor,
events,
roster,
base: fleet.base().to_string(),
pending: RosterChange::default(),
})
}
pub fn roster(&self) -> &BTreeMap<String, Vec<String>> {
&self.roster
}
pub async fn next_change(&mut self) -> Option<RosterChange> {
next_change_in(
&mut self.events,
&mut self.roster,
&self.base,
&mut self.pending,
)
.await
}
pub async fn stop(self) -> Result<()> {
self.monitor.shutdown().await
}
}
async fn next_change_in(
events: &mut crate::EventStream,
roster: &mut BTreeMap<String, Vec<String>>,
base: &str,
pending: &mut RosterChange,
) -> Option<RosterChange> {
loop {
if let Some(change) = take_change(pending) {
return Some(change);
}
let mut item = events.recv().await;
loop {
match item {
None => return take_change(pending),
Some(crate::StreamItem::Dropped(_)) => {}
Some(crate::StreamItem::Event(ev)) => {
let transition = match ev {
crate::FleetEvent::NodeUp(key) => Some((key, true)),
crate::FleetEvent::NodeDown(key) => Some((key, false)),
_ => None,
};
if let Some((key, up)) = transition
&& apply_token(roster, base, &key, up)
{
if up {
pending.node_up = true;
} else {
pending.node_down = true;
}
}
}
}
match tokio::time::timeout(BURST_QUIET, events.recv()).await {
Ok(next) => item = next,
Err(_) => break,
}
}
}
}
fn take_change(pending: &mut RosterChange) -> Option<RosterChange> {
if !pending.node_up && !pending.node_down {
return None;
}
Some(std::mem::take(pending))
}
pub fn token_identity(base: &str, key: &str) -> Option<(String, String)> {
let parsed = zenkey::grammar::parse_full(base, key)?;
let origin = parsed.origin.chunk().to_string();
let producer = parsed
.producer()
.map(|p| p.chunk())
.unwrap_or_else(|| origin.trim_start_matches('@').to_string());
Some((origin, producer))
}
pub fn apply_token(
roster: &mut BTreeMap<String, Vec<String>>,
base: &str,
key: &str,
up: bool,
) -> bool {
let Some((origin, producer)) = token_identity(base, key) else {
return false;
};
if up {
let entry = roster.entry(origin).or_default();
if entry.contains(&producer) {
return false;
}
entry.push(producer);
entry.sort();
return true;
}
let Some(entry) = roster.get_mut(&origin) else {
return false;
};
let before = entry.len();
entry.retain(|p| p != &producer);
let changed = entry.len() != before;
if entry.is_empty() {
roster.remove(&origin);
}
changed
}
pub fn node_rows(
roster: &BTreeMap<String, Vec<String>>,
slices: Option<&crate::SliceSet>,
) -> crate::report::NodeList {
let mut nodes = Vec::new();
for (origin, producers) in roster {
for producer in producers {
let joined = slices.and_then(|s| {
let base_name = zenkey::grammar::Producer::parse_chunk(producer)
.map(|pr| pr.name().to_string())
.unwrap_or_else(|_| producer.clone());
s.get(&base_name)
});
nodes.push(crate::report::NodeRow {
origin: origin.clone(),
producer: producer.clone(),
app: joined.map(|s| s.app.clone()),
registry_version: joined.map(|s| s.version.clone()),
});
}
}
crate::report::NodeList {
nodes,
slices_joined: slices.is_some(),
}
}
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(Error::from)
} else {
zenkey::origin::RemoteOrigin::parse(origin)
.map(Node::Host)
.map_err(|e| {
Error::unaskable(
"origin",
format!("{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),
zenkey::selector::Producers::all(),
&["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(
fleet: &crate::Fleet<'_>,
origin: &str,
timeout: Duration,
with_freshness: bool,
) -> Result<NodeInfo> {
let (session, base) = (fleet.session(), fleet.base());
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()
.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::bus::query::fleet_get(
fleet,
&introspect,
&crate::bus::query::GetOpts::new(timeout),
)
.await
.unwrap_or_default();
let served: Vec<zenkey::slice::RegistrySlice> = answers
.into_iter()
.filter(|a| a.origin == origin)
.filter_map(|a| {
let crate::bus::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.token().to_string()).collect())
.unwrap_or_default(),
media: slice
.map(|s| {
s.media
.iter()
.map(|m| MediaStreamInfo {
path: m.path.clone(),
encoding: m.encoding.as_encoding_str().to_string(),
})
.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::bus::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.is(&zenkey::Class::State) {
continue;
}
let age = samples
.iter()
.filter_map(|s| {
let parsed = zenkey::grammar::parse_full(base, &s.key)?;
let p = parsed.producer()?.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(
fleet: &crate::Fleet<'_>,
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 = fleet.wire(relative);
let answers =
crate::bus::query::fleet_get(fleet, &key, &crate::bus::query::GetOpts::new(timeout))
.await?;
let mut matches = Vec::new();
let seen = answers.len();
for a in &answers {
let crate::bus::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))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_token_names_its_origin_and_producer() {
assert_eq!(
token_identity("acme", "acme/v1/h-3fa9c2d41b7e/state/sysinfo/alive"),
Some(("h-3fa9c2d41b7e".into(), "sysinfo".into()))
);
assert_eq!(
token_identity("acme", "acme/v1/@catalog/state/alive"),
Some(("@catalog".into(), "catalog".into()))
);
assert_eq!(
token_identity("", "v1/h-3fa9c2d41b7e/state/sysinfo/alive"),
Some(("h-3fa9c2d41b7e".into(), "sysinfo".into()))
);
assert_eq!(token_identity("acme", "demo/not/a/token"), None);
}
#[test]
fn applying_a_token_reports_only_real_changes() {
let mut roster: BTreeMap<String, Vec<String>> = BTreeMap::new();
let key = "acme/v1/h-3fa9c2d41b7e/state/sysinfo/alive";
assert!(apply_token(&mut roster, "acme", key, true));
assert!(
!apply_token(&mut roster, "acme", key, true),
"a replayed history token is not a change"
);
assert_eq!(roster["h-3fa9c2d41b7e"], ["sysinfo"]);
let other = "acme/v1/h-3fa9c2d41b7e/state/alerts/alive";
assert!(apply_token(&mut roster, "acme", other, true));
assert_eq!(roster["h-3fa9c2d41b7e"], ["alerts", "sysinfo"]);
assert!(apply_token(&mut roster, "acme", key, false));
assert!(
!apply_token(&mut roster, "acme", key, false),
"retracting what is already gone is not a change"
);
assert_eq!(roster["h-3fa9c2d41b7e"], ["alerts"]);
assert!(apply_token(&mut roster, "acme", other, false));
assert!(roster.is_empty());
assert!(!apply_token(&mut roster, "acme", "demo/foreign", true));
}
#[test]
fn rows_say_whether_a_slice_was_even_asked_for() {
let mut roster: BTreeMap<String, Vec<String>> = BTreeMap::new();
roster.insert(
"h-3fa9c2d41b7e".into(),
vec!["sysinfo".into(), "sysinfo-2".into()],
);
let unasked = node_rows(&roster, None);
assert!(!unasked.slices_joined, "no join was attempted");
assert!(unasked.nodes.iter().all(|n| n.app.is_none()));
let slice = zenkey::parse_slice(
"[registry]\nversion = \"1.0\"\napp = \"demo\"\nconvention = 1\n\
[producer]\nname = \"sysinfo\"\n",
)
.expect("fixture slice parses");
let joined = node_rows(&roster, Some(&crate::SliceSet::from_slices(vec![slice])));
assert!(joined.slices_joined);
assert_eq!(joined.nodes.len(), 2);
for row in &joined.nodes {
assert_eq!(
row.app.as_deref(),
Some("demo"),
"an instance suffix shares the base producer's slice: {}",
row.producer
);
}
assert_eq!(
joined.nodes[1].producer, "sysinfo-2",
"the row keeps the suffix"
);
}
#[tokio::test(start_paused = true)]
async fn a_cancelled_poll_keeps_the_change_it_already_applied() {
let core = crate::MonitorCore::new(16);
let mut events = core.events();
let mut roster: BTreeMap<String, Vec<String>> = BTreeMap::new();
let mut pending = RosterChange::default();
core.node_event("v1/h-3fa9c2d41b7e/state/sysinfo/alive".into(), true);
let cancelled = tokio::time::timeout(
BURST_QUIET / 2,
next_change_in(&mut events, &mut roster, "", &mut pending),
)
.await;
assert!(cancelled.is_err(), "the poll is still draining the burst");
assert!(
roster.contains_key("h-3fa9c2d41b7e"),
"the token was applied before the drop"
);
let change = tokio::time::timeout(
BURST_QUIET / 2,
next_change_in(&mut events, &mut roster, "", &mut pending),
)
.await
.expect("the applied change is reported, not waited on")
.expect("a change, not a closed stream");
assert_eq!(
change,
RosterChange {
node_up: true,
node_down: false
}
);
}
}