use async_nats::jetstream::kv::Config as KvConfig;
use axum::Json;
use axum::extract::{Path, State};
use axum::http::StatusCode;
use axum::http::header::HeaderMap;
use futures::TryStreamExt;
use kanade_shared::kv::{BUCKET_GROUP_DEFS, BUCKET_GROUP_DEFS_YAML};
use kanade_shared::manifest::GroupDef;
use serde::Serialize;
use tracing::{info, warn};
use crate::api::AppState;
use crate::api::yaml_body::{YamlOrJson, mirror_yaml, yaml_headers};
use crate::audit;
use crate::audit::Caller;
#[derive(Serialize)]
pub struct GroupDefSummary {
pub id: String,
pub description: Option<String>,
pub kind: &'static str,
pub member_count: Option<usize>,
pub tags: Vec<String>,
}
#[derive(Serialize)]
pub struct GroupMembers {
pub id: String,
pub kind: &'static str,
pub count: usize,
pub members: Vec<String>,
}
fn kind_of(g: &GroupDef) -> &'static str {
if g.dynamic_query().is_some() {
"dynamic"
} else {
"static"
}
}
fn guard_id(id: &str) -> Result<(), (StatusCode, String)> {
if kanade_shared::manifest::is_valid_resource_id(id) {
Ok(())
} else {
Err((
StatusCode::BAD_REQUEST,
format!("invalid group id '{id}' (allowed: [A-Za-z0-9._-])"),
))
}
}
pub async fn list(State(s): State<AppState>) -> Result<Json<Vec<GroupDef>>, (StatusCode, String)> {
let kv = match s.jetstream.get_key_value(BUCKET_GROUP_DEFS).await {
Ok(k) => k,
Err(_) => return Ok(Json(Vec::new())),
};
let keys: Vec<String> = kv
.keys()
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("kv keys: {e}")))?
.try_collect()
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("kv keys: {e}")))?;
let mut out = Vec::with_capacity(keys.len());
for k in keys {
match kv.get(&k).await {
Ok(Some(bytes)) => match serde_json::from_slice::<GroupDef>(&bytes) {
Ok(g) => out.push(g),
Err(e) => {
warn!(group_id = %k, error = %e, "group_defs: skipping undecodable entry")
}
},
Ok(None) => {}
Err(e) => warn!(group_id = %k, error = %e, "group_defs: kv get failed, skipping"),
}
}
out.sort_by(|a, b| a.id.cmp(&b.id));
Ok(Json(out))
}
pub async fn create(
State(s): State<AppState>,
caller: Caller,
body: YamlOrJson<GroupDef>,
) -> Result<Json<GroupDefSummary>, (StatusCode, String)> {
let YamlOrJson {
value: group,
raw_yaml,
} = body;
if let Err(e) = group.validate() {
return Err((StatusCode::BAD_REQUEST, format!("invalid group: {e}")));
}
if let Some(query) = group.dynamic_query()
&& let Err(e) = crate::api::query::validate_read_only(query)
{
return Err((
StatusCode::BAD_REQUEST,
format!("invalid group: query is not a read-only query: {e}"),
));
}
let kv = s
.jetstream
.create_key_value(KvConfig {
bucket: BUCKET_GROUP_DEFS.into(),
history: 5,
..Default::default()
})
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("ensure KV: {e}")))?;
let body_bytes = serde_json::to_vec(&group)
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("serialize: {e}")))?;
kv.put(&group.id, body_bytes.into())
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("KV put: {e}")))?;
s.group_cache
.lock()
.expect("group_cache mutex")
.remove(&group.id);
let yaml_source = raw_yaml.unwrap_or_else(|| {
serde_yaml::to_string(&group)
.unwrap_or_else(|_| String::from("# YAML mirror unavailable for this entry"))
});
if let Err(e) = mirror_yaml(&s, BUCKET_GROUP_DEFS_YAML, &group.id, &yaml_source).await {
warn!(error = %e, group_id = %group.id, "group_defs: YAML mirror put failed; JSON catalog is current");
}
let kind = kind_of(&group);
let member_count = (kind == "static").then_some(group.members.len());
info!(group_id = %group.id, kind, "group def upserted");
audit::record(
&s.nats,
"operator",
"group_def_upsert",
Some(&group.id),
Some(&caller),
serde_json::json!({ "kind": kind }),
)
.await;
Ok(Json(GroupDefSummary {
id: group.id,
description: group.description,
kind,
member_count,
tags: group.tags,
}))
}
pub async fn get_yaml(
State(s): State<AppState>,
Path(id): Path<String>,
) -> Result<(StatusCode, HeaderMap, String), (StatusCode, String)> {
guard_id(&id)?;
if let Ok(kv) = s.jetstream.get_key_value(BUCKET_GROUP_DEFS_YAML).await
&& let Ok(Some(bytes)) = kv.get(&id).await
&& let Ok(text) = String::from_utf8(bytes.to_vec())
{
return Ok((StatusCode::OK, yaml_headers(), text));
}
let kv = s
.jetstream
.get_key_value(BUCKET_GROUP_DEFS)
.await
.map_err(|_| (StatusCode::NOT_FOUND, format!("group '{id}' not found")))?;
let bytes = kv
.get(&id)
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("KV get: {e}")))?
.ok_or_else(|| (StatusCode::NOT_FOUND, format!("group '{id}' not found")))?;
let group: GroupDef = serde_json::from_slice(&bytes)
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("decode: {e}")))?;
let yaml = serde_yaml::to_string(&group).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("encode YAML: {e}"),
)
})?;
Ok((StatusCode::OK, yaml_headers(), yaml))
}
pub async fn members(
State(s): State<AppState>,
Path(id): Path<String>,
) -> Result<Json<GroupMembers>, (StatusCode, String)> {
guard_id(&id)?;
let kv = s
.jetstream
.get_key_value(BUCKET_GROUP_DEFS)
.await
.map_err(|_| (StatusCode::NOT_FOUND, format!("group '{id}' not found")))?;
let bytes = kv
.get(&id)
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("KV get: {e}")))?
.ok_or_else(|| (StatusCode::NOT_FOUND, format!("group '{id}' not found")))?;
let group: GroupDef = serde_json::from_slice(&bytes)
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("decode: {e}")))?;
let kind = kind_of(&group);
let mut members = crate::api::group_sql::resolve_group_members(&s, &group)
.await
.map_err(|e| {
(
StatusCode::BAD_REQUEST,
format!("resolve group '{id}': {e}"),
)
})?;
members.sort();
members.dedup();
Ok(Json(GroupMembers {
id,
kind,
count: members.len(),
members,
}))
}
pub async fn delete(
State(s): State<AppState>,
Path(id): Path<String>,
caller: Caller,
) -> Result<StatusCode, (StatusCode, String)> {
guard_id(&id)?;
let kv = match s.jetstream.get_key_value(BUCKET_GROUP_DEFS).await {
Ok(k) => k,
Err(e) => {
warn!(error = %e, "group_defs KV missing on delete");
return Err((StatusCode::NOT_FOUND, "group_defs bucket missing".into()));
}
};
if kv
.get(&id)
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("kv get: {e}")))?
.is_none()
{
return Err((StatusCode::NOT_FOUND, format!("group '{id}' not found")));
}
kv.delete(&id)
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("kv delete: {e}")))?;
if let Ok(ykv) = s.jetstream.get_key_value(BUCKET_GROUP_DEFS_YAML).await {
let _ = ykv.delete(&id).await;
}
s.group_cache.lock().expect("group_cache mutex").remove(&id);
info!(group_id = %id, "group def deleted");
audit::record(
&s.nats,
"operator",
"group_def_delete",
Some(&id),
Some(&caller),
serde_json::json!({}),
)
.await;
Ok(StatusCode::NO_CONTENT)
}