use std::collections::{BTreeMap, BTreeSet};
use std::time::Duration;
use anyhow::{Context, Result};
use async_nats::jetstream::kv::{Config as KvConfig, Store};
use futures::{StreamExt, TryStreamExt};
use kanade_shared::kv::{BUCKET_AGENT_GROUPS_DERIVED, BUCKET_GROUP_DEFS};
use kanade_shared::manifest::GroupDef;
use kanade_shared::wire::AgentGroups;
use tracing::{debug, info, warn};
use crate::api::AppState;
use crate::api::group_sql::resolve_group_members;
const RECONCILE_INTERVAL: Duration = Duration::from_secs(60);
pub async fn run(state: AppState) {
info!("group materializer started");
loop {
if let Err(e) = reconcile_once(&state).await {
warn!(error = %e, "group materializer reconcile failed; will retry");
}
wait_for_trigger(&state).await;
}
}
async fn wait_for_trigger(state: &AppState) {
let tick = tokio::time::sleep(RECONCILE_INTERVAL);
let Ok(kv) = state.jetstream.get_key_value(BUCKET_GROUP_DEFS).await else {
tick.await;
return;
};
let Ok(mut watch) = kv.watch_all().await else {
tick.await;
return;
};
tokio::select! {
_ = tick => {}
ev = watch.next() => {
if ev.is_some() {
debug!("group_defs changed — reconciling derived membership");
}
}
}
}
async fn reconcile_once(state: &AppState) -> Result<()> {
let defs = load_group_defs(state).await?;
let mut resolved: Vec<(String, Vec<String>)> = Vec::with_capacity(defs.len());
let mut failed: BTreeSet<String> = BTreeSet::new();
for g in &defs {
match resolve_group_members(state, g).await {
Ok(pcs) => resolved.push((g.id.clone(), pcs)),
Err(e) => {
warn!(group = %g.id, error = %e, "materializer: group resolve failed; preserving existing members");
failed.insert(g.id.clone());
}
}
}
let desired = build_desired(&resolved);
write_derived(state, &desired, &failed).await
}
fn build_desired(resolved: &[(String, Vec<String>)]) -> BTreeMap<String, BTreeSet<String>> {
let mut desired: BTreeMap<String, BTreeSet<String>> = BTreeMap::new();
for (group_id, pcs) in resolved {
for pc in pcs {
desired
.entry(pc.clone())
.or_default()
.insert(group_id.clone());
}
}
desired
}
async fn write_derived(
state: &AppState,
desired: &BTreeMap<String, BTreeSet<String>>,
failed: &BTreeSet<String>,
) -> Result<()> {
let kv = state
.jetstream
.create_key_value(KvConfig {
bucket: BUCKET_AGENT_GROUPS_DERIVED.into(),
history: 1,
..Default::default()
})
.await
.context("ensure agent_groups_derived KV")?;
let existing: Vec<String> = match kv.keys().await {
Ok(k) => k.try_collect().await.unwrap_or_default(),
Err(_) => Vec::new(),
};
let mut candidates: BTreeSet<String> = desired.keys().cloned().collect();
candidates.extend(existing);
const WRITE_CONCURRENCY: usize = 16;
let _: Vec<()> = futures::stream::iter(candidates)
.map(|pc| {
let kv = kv.clone();
async move { reconcile_pc(&kv, &pc, desired.get(&pc), failed).await }
})
.buffer_unordered(WRITE_CONCURRENCY)
.try_collect()
.await?;
Ok(())
}
async fn reconcile_pc(
kv: &Store,
pc: &str,
desired: Option<&BTreeSet<String>>,
failed: &BTreeSet<String>,
) -> Result<()> {
let current = read_derived(kv, pc).await;
let current_groups = current.as_ref().map(|g| g.groups.as_slice()).unwrap_or(&[]);
let target = target_for_pc(desired, current_groups, failed);
if target.is_empty() {
if current.is_some() {
kv.delete(pc)
.await
.with_context(|| format!("delete stale derived membership for {pc}"))?;
debug!(pc = %pc, "cleared derived membership (left all declared groups)");
}
return Ok(());
}
let want = AgentGroups::new(target);
if current.as_ref() == Some(&want) {
return Ok(()); }
let bytes = serde_json::to_vec(&want).context("encode derived AgentGroups")?;
kv.put(pc, bytes.into())
.await
.with_context(|| format!("put derived membership for {pc}"))?;
debug!(pc = %pc, groups = ?want.groups, "materialized derived membership");
Ok(())
}
fn target_for_pc(
desired: Option<&BTreeSet<String>>,
current: &[String],
failed: &BTreeSet<String>,
) -> BTreeSet<String> {
let mut target: BTreeSet<String> = desired.cloned().unwrap_or_default();
for g in current {
if failed.contains(g) {
target.insert(g.clone());
}
}
target
}
async fn read_derived(kv: &Store, pc: &str) -> Option<AgentGroups> {
match kv.get(pc).await {
Ok(Some(bytes)) => serde_json::from_slice::<AgentGroups>(&bytes).ok(),
_ => None,
}
}
async fn load_group_defs(state: &AppState) -> Result<Vec<GroupDef>> {
let Ok(kv) = state.jetstream.get_key_value(BUCKET_GROUP_DEFS).await else {
return Ok(Vec::new());
};
let keys: Vec<String> = match kv.keys().await {
Ok(k) => k.try_collect().await.unwrap_or_default(),
Err(_) => return Ok(Vec::new()),
};
let mut out = Vec::with_capacity(keys.len());
for k in keys {
if let Ok(Some(bytes)) = kv.get(&k).await
&& let Ok(g) = serde_json::from_slice::<GroupDef>(&bytes)
{
out.push(g);
}
}
Ok(out)
}
#[cfg(test)]
mod tests {
use super::*;
fn set(names: &[&str]) -> BTreeSet<String> {
names.iter().map(|s| s.to_string()).collect()
}
#[test]
fn build_desired_inverts_group_to_pc() {
let resolved = vec![
("clients".to_string(), vec!["PC-1".into(), "PC-2".into()]),
("servers".to_string(), vec!["SRV-1".into()]),
];
let desired = build_desired(&resolved);
assert_eq!(desired.get("PC-1"), Some(&set(&["clients"])));
assert_eq!(desired.get("PC-2"), Some(&set(&["clients"])));
assert_eq!(desired.get("SRV-1"), Some(&set(&["servers"])));
assert_eq!(desired.len(), 3);
}
#[test]
fn build_desired_unions_multiple_groups_per_pc() {
let resolved = vec![
("clients".to_string(), vec!["PC-1".into()]),
("win-24h2".to_string(), vec!["PC-1".into()]),
];
let desired = build_desired(&resolved);
assert_eq!(desired.get("PC-1"), Some(&set(&["clients", "win-24h2"])));
assert_eq!(desired.len(), 1);
}
fn slice(names: &[&str]) -> Vec<String> {
names.iter().map(|s| s.to_string()).collect()
}
#[test]
fn target_uses_desired_when_nothing_failed() {
let desired = set(&["clients"]);
let t = target_for_pc(
Some(&desired),
&slice(&["clients", "old"]),
&BTreeSet::new(),
);
assert_eq!(t, set(&["clients"]));
}
#[test]
fn target_preserves_a_failed_group_the_pc_currently_has() {
let desired = set(&["clients"]);
let failed = set(&["servers"]);
let t = target_for_pc(Some(&desired), &slice(&["clients", "servers"]), &failed);
assert_eq!(t, set(&["clients", "servers"]));
}
#[test]
fn target_does_not_add_a_failed_group_the_pc_lacks() {
let failed = set(&["servers"]);
let t = target_for_pc(None, &slice(&["clients"]), &failed);
assert!(t.is_empty());
}
#[test]
fn target_empty_means_delete() {
let t = target_for_pc(None, &slice(&["gone"]), &BTreeSet::new());
assert!(t.is_empty());
}
#[test]
fn build_desired_empty_when_no_groups() {
assert!(build_desired(&[]).is_empty());
let resolved = vec![("empty".to_string(), vec![])];
assert!(build_desired(&resolved).is_empty());
}
}