use std::time::Duration;
use zenkey::schema::SchemaSet;
use crate::Result;
use crate::bus::query::{Answer, GetOpts, fleet_get};
use crate::model::decode::DescribedSchema;
use crate::model::registry::{SliceSet, rpc_key};
#[derive(Debug, Clone, Default)]
pub struct DescribeSweep {
pub answers: Vec<DescribedSchema>,
pub undescribed: Vec<String>,
}
impl DescribeSweep {
pub fn first_per_producer(&self) -> Vec<(String, SchemaSet)> {
let mut seen = std::collections::BTreeSet::new();
self.answers
.iter()
.filter(|d| seen.insert(d.producer.as_str()))
.map(|d| (d.producer.clone(), d.set.clone()))
.collect()
}
}
pub async fn describe_sweep(
fleet: &crate::Fleet<'_>,
slices: &SliceSet,
timeout: Duration,
) -> Result<DescribeSweep> {
let base = fleet.base();
let mut answers: Vec<DescribedSchema> = Vec::new();
let mut undescribed: Vec<String> = Vec::new();
let opts = GetOpts::new(timeout);
for slice in slices.slices() {
let key = rpc_key(base, slice, "describe")?;
let replies = fleet_get(fleet, &key, &opts).await?;
let before = answers.len();
for reply in replies {
let origin = reply.origin;
let Answer::Value(bytes) = reply.answer else {
continue;
};
let cow = bytes.to_bytes();
let Some(set) = std::str::from_utf8(&cow)
.ok()
.and_then(|t| SchemaSet::parse(t).ok())
else {
continue;
};
answers.push(DescribedSchema::new(origin, slice.name.clone(), set));
}
if answers.len() == before {
undescribed.push(slice.name.clone());
}
}
Ok(DescribeSweep {
answers,
undescribed,
})
}
#[cfg(test)]
mod tests {
use super::*;
use zenkey::schema::TypeSchema;
fn set(name: &str) -> SchemaSet {
SchemaSet::builder("t")
.entry(
name,
TypeSchema::json_schema(serde_json::json!({"type": "object"})),
)
.build()
}
#[test]
fn first_per_producer_keeps_the_first_host_and_every_answer_stays() {
let sweep = DescribeSweep {
answers: vec![
DescribedSchema::new("h-aaaaaaaaaaaa", "sysinfo", set("A")),
DescribedSchema::new("h-bbbbbbbbbbbb", "sysinfo", set("B")),
DescribedSchema::new("h-aaaaaaaaaaaa", "gnmi", set("C")),
],
undescribed: vec![],
};
let first = sweep.first_per_producer();
assert_eq!(first.len(), 2);
assert_eq!(first[0].0, "sysinfo");
assert!(first[0].1.get("A").is_some(), "the first host's set wins");
assert_eq!(first[1].0, "gnmi");
assert_eq!(sweep.answers.len(), 3, "the attributed list is untouched");
}
}