use std::collections::BTreeMap;
use axum::Json;
use axum::extract::{Path, State};
use axum::http::StatusCode;
use futures::StreamExt;
use serde::Serialize;
use tracing::{info, warn};
use kanade_shared::config_field::{FieldUpdate, FieldValue, apply_field};
use kanade_shared::kv::{
BUCKET_AGENT_CONFIG, BUCKET_AGENT_GROUPS, KEY_AGENT_CONFIG_GLOBAL, agent_config_group_key,
agent_config_pc_key, parse_agent_config_group_key,
};
use kanade_shared::wire::{AgentGroups, ConfigScope, EffectiveConfig, ResolutionWarning, resolve};
use super::AppState;
use crate::audit::{self, Caller};
pub async fn get_global(
State(state): State<AppState>,
) -> Result<Json<ConfigScope>, (StatusCode, String)> {
let kv = open_cfg(&state).await?;
Ok(Json(
read_scope_or_default(&kv, KEY_AGENT_CONFIG_GLOBAL).await?,
))
}
pub async fn put_global(
State(state): State<AppState>,
caller: Caller,
Json(scope): Json<ConfigScope>,
) -> Result<Json<ConfigScope>, (StatusCode, String)> {
let kv = open_cfg(&state).await?;
write_scope(&kv, KEY_AGENT_CONFIG_GLOBAL, &scope).await?;
info!(scope = ?scope, "agent_config.global replaced");
audit_scope_put(&state, &caller, KEY_AGENT_CONFIG_GLOBAL, &scope).await;
Ok(Json(scope))
}
pub async fn delete_global(
State(state): State<AppState>,
caller: Caller,
) -> Result<StatusCode, (StatusCode, String)> {
let kv = open_cfg(&state).await?;
delete_key(&kv, KEY_AGENT_CONFIG_GLOBAL).await?;
info!("agent_config.global deleted");
audit_scope_clear(&state, &caller, KEY_AGENT_CONFIG_GLOBAL).await;
Ok(StatusCode::NO_CONTENT)
}
pub async fn set_field_global(
State(state): State<AppState>,
caller: Caller,
Path(field): Path<String>,
Json(body): Json<FieldValue>,
) -> Result<Json<FieldUpdate>, (StatusCode, String)> {
set_field(&state, &caller, KEY_AGENT_CONFIG_GLOBAL, &field, body.value).await
}
pub async fn unset_field_global(
State(state): State<AppState>,
caller: Caller,
Path(field): Path<String>,
) -> Result<Json<FieldUpdate>, (StatusCode, String)> {
unset_field(&state, &caller, KEY_AGENT_CONFIG_GLOBAL, &field).await
}
pub async fn defaults() -> Json<EffectiveConfig> {
Json(EffectiveConfig::builtin_defaults())
}
pub async fn get_group(
State(state): State<AppState>,
Path(name): Path<String>,
) -> Result<Json<ConfigScope>, (StatusCode, String)> {
let kv = open_cfg(&state).await?;
Ok(Json(
read_scope_or_default(&kv, &agent_config_group_key(&name)).await?,
))
}
pub async fn put_group(
State(state): State<AppState>,
caller: Caller,
Path(name): Path<String>,
Json(scope): Json<ConfigScope>,
) -> Result<Json<ConfigScope>, (StatusCode, String)> {
let kv = open_cfg(&state).await?;
let key = agent_config_group_key(&name);
write_scope(&kv, &key, &scope).await?;
info!(group = %name, scope = ?scope, "agent_config.groups.<name> replaced");
audit_scope_put(&state, &caller, &key, &scope).await;
Ok(Json(scope))
}
pub async fn delete_group(
State(state): State<AppState>,
caller: Caller,
Path(name): Path<String>,
) -> Result<StatusCode, (StatusCode, String)> {
let kv = open_cfg(&state).await?;
let key = agent_config_group_key(&name);
delete_key(&kv, &key).await?;
info!(group = %name, "agent_config.groups.<name> deleted");
audit_scope_clear(&state, &caller, &key).await;
Ok(StatusCode::NO_CONTENT)
}
pub async fn set_field_group(
State(state): State<AppState>,
caller: Caller,
Path((name, field)): Path<(String, String)>,
Json(body): Json<FieldValue>,
) -> Result<Json<FieldUpdate>, (StatusCode, String)> {
set_field(
&state,
&caller,
&agent_config_group_key(&name),
&field,
body.value,
)
.await
}
pub async fn unset_field_group(
State(state): State<AppState>,
caller: Caller,
Path((name, field)): Path<(String, String)>,
) -> Result<Json<FieldUpdate>, (StatusCode, String)> {
unset_field(&state, &caller, &agent_config_group_key(&name), &field).await
}
pub async fn get_pc(
State(state): State<AppState>,
Path(pc_id): Path<String>,
) -> Result<Json<ConfigScope>, (StatusCode, String)> {
let kv = open_cfg(&state).await?;
Ok(Json(
read_scope_or_default(&kv, &agent_config_pc_key(&pc_id)).await?,
))
}
pub async fn put_pc(
State(state): State<AppState>,
caller: Caller,
Path(pc_id): Path<String>,
Json(scope): Json<ConfigScope>,
) -> Result<Json<ConfigScope>, (StatusCode, String)> {
let kv = open_cfg(&state).await?;
let key = agent_config_pc_key(&pc_id);
write_scope(&kv, &key, &scope).await?;
info!(pc_id = %pc_id, scope = ?scope, "agent_config.pcs.<pc_id> replaced");
audit_scope_put(&state, &caller, &key, &scope).await;
Ok(Json(scope))
}
pub async fn delete_pc(
State(state): State<AppState>,
caller: Caller,
Path(pc_id): Path<String>,
) -> Result<StatusCode, (StatusCode, String)> {
let kv = open_cfg(&state).await?;
let key = agent_config_pc_key(&pc_id);
delete_key(&kv, &key).await?;
info!(pc_id = %pc_id, "agent_config.pcs.<pc_id> deleted");
audit_scope_clear(&state, &caller, &key).await;
Ok(StatusCode::NO_CONTENT)
}
pub async fn set_field_pc(
State(state): State<AppState>,
caller: Caller,
Path((pc_id, field)): Path<(String, String)>,
Json(body): Json<FieldValue>,
) -> Result<Json<FieldUpdate>, (StatusCode, String)> {
set_field(
&state,
&caller,
&agent_config_pc_key(&pc_id),
&field,
body.value,
)
.await
}
pub async fn unset_field_pc(
State(state): State<AppState>,
caller: Caller,
Path((pc_id, field)): Path<(String, String)>,
) -> Result<Json<FieldUpdate>, (StatusCode, String)> {
unset_field(&state, &caller, &agent_config_pc_key(&pc_id), &field).await
}
async fn set_field(
state: &AppState,
caller: &Caller,
key: &str,
field: &str,
value: String,
) -> Result<Json<FieldUpdate>, (StatusCode, String)> {
let kv = open_cfg(state).await?;
let update = update_field(&kv, key, field, Some(&value)).await?;
info!(
key,
field,
changed = update.changed,
"agent_config field set"
);
audit::record(
&state.nats,
"operator",
"config_field_set",
Some(key),
Some(caller),
serde_json::json!({ "field": field, "value": value, "changed": update.changed }),
)
.await;
Ok(Json(update))
}
async fn unset_field(
state: &AppState,
caller: &Caller,
key: &str,
field: &str,
) -> Result<Json<FieldUpdate>, (StatusCode, String)> {
let kv = open_cfg(state).await?;
let update = update_field(&kv, key, field, None).await?;
info!(
key,
field,
changed = update.changed,
"agent_config field unset"
);
audit::record(
&state.nats,
"operator",
"config_field_unset",
Some(key),
Some(caller),
serde_json::json!({ "field": field, "changed": update.changed }),
)
.await;
Ok(Json(update))
}
async fn update_field(
kv: &async_nats::jetstream::kv::Store,
key: &str,
field: &str,
value: Option<&str>,
) -> Result<FieldUpdate, (StatusCode, String)> {
validate_field(field, value)?;
let mut changed = false;
let scope = kanade_shared::kv_cas::read_modify_write(kv, key, |scope: &mut ConfigScope| {
let before = scope.clone();
let _ = apply_field(scope, field, value);
changed = *scope != before;
changed
})
.await
.map_err(|e| {
warn!(error = %e, key, field, "agent_config field update");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("update {key}: {e:#}"),
)
})?;
Ok(FieldUpdate { scope, changed })
}
fn validate_field(field: &str, value: Option<&str>) -> Result<(), (StatusCode, String)> {
apply_field(&mut ConfigScope::default(), field, value)
.map_err(|e| (StatusCode::BAD_REQUEST, format!("{e:#}")))
}
async fn audit_scope_put(state: &AppState, caller: &Caller, key: &str, scope: &ConfigScope) {
audit::record(
&state.nats,
"operator",
"config_scope_put",
Some(key),
Some(caller),
serde_json::json!({ "scope": scope }),
)
.await;
}
async fn audit_scope_clear(state: &AppState, caller: &Caller, key: &str) {
audit::record(
&state.nats,
"operator",
"config_scope_clear",
Some(key),
Some(caller),
serde_json::json!({}),
)
.await;
}
#[derive(Serialize)]
pub struct EffectiveConfigResponse {
pub pc_id: String,
pub effective: EffectiveConfig,
pub warnings: Vec<String>,
pub my_groups: Vec<String>,
}
pub async fn effective(
State(state): State<AppState>,
Path(pc_id): Path<String>,
) -> Result<Json<EffectiveConfigResponse>, (StatusCode, String)> {
let cfg_kv = open_cfg(&state).await?;
let groups_kv = open_groups(&state).await?;
let (global_scope, group_scopes) = collect_global_and_groups(&cfg_kv).await?;
let pc_scope = read_optional_scope(&cfg_kv, &agent_config_pc_key(&pc_id)).await?;
let my_groups = pc_group_memberships(&groups_kv, &pc_id).await;
let (effective, warns) = resolve(
global_scope.as_ref(),
&group_scopes,
pc_scope.as_ref(),
&my_groups,
);
Ok(Json(EffectiveConfigResponse {
pc_id,
effective,
warnings: warns.into_iter().map(render_warning).collect(),
my_groups,
}))
}
pub async fn pc_inherited(
State(state): State<AppState>,
Path(pc_id): Path<String>,
) -> Result<Json<EffectiveConfig>, (StatusCode, String)> {
let cfg_kv = open_cfg(&state).await?;
let groups_kv = open_groups(&state).await?;
let (global_scope, group_scopes) = collect_global_and_groups(&cfg_kv).await?;
let my_groups = pc_group_memberships(&groups_kv, &pc_id).await;
let (inherited, _warns) = resolve(global_scope.as_ref(), &group_scopes, None, &my_groups);
Ok(Json(inherited))
}
pub async fn group_inherited(
State(state): State<AppState>,
Path(_name): Path<String>,
) -> Result<Json<EffectiveConfig>, (StatusCode, String)> {
let cfg_kv = open_cfg(&state).await?;
let global_scope = read_optional_scope(&cfg_kv, KEY_AGENT_CONFIG_GLOBAL).await?;
let (inherited, _warns) = resolve(global_scope.as_ref(), &BTreeMap::new(), None, &[]);
Ok(Json(inherited))
}
fn render_warning(w: ResolutionWarning) -> String {
match w {
ResolutionWarning::MultiGroupConflict { field, groups } => format!(
"multi-group conflict on `{field}` — set by [{}]; alphabetical last wins (=> {})",
groups.join(", "),
groups.last().map(String::as_str).unwrap_or("<none>"),
),
}
}
async fn collect_global_and_groups(
cfg_kv: &async_nats::jetstream::kv::Store,
) -> Result<(Option<ConfigScope>, BTreeMap<String, ConfigScope>), (StatusCode, String)> {
let global_scope = read_optional_scope(cfg_kv, KEY_AGENT_CONFIG_GLOBAL).await?;
let mut group_scopes: BTreeMap<String, ConfigScope> = BTreeMap::new();
match cfg_kv.keys().await {
Ok(mut keys) => {
while let Some(k) = keys.next().await {
let key = match k {
Ok(k) => k,
Err(e) => {
warn!(error = %e, "agent_config keys()");
continue;
}
};
if let Some(group) = parse_agent_config_group_key(&key)
&& let Ok(Some(bytes)) = cfg_kv.get(&key).await
&& let Ok(scope) = serde_json::from_slice::<ConfigScope>(&bytes)
{
group_scopes.insert(group.to_string(), scope);
}
}
}
Err(e) => {
warn!(error = %e, "agent_config keys() for effective");
}
}
Ok((global_scope, group_scopes))
}
async fn pc_group_memberships(
groups_kv: &async_nats::jetstream::kv::Store,
pc_id: &str,
) -> Vec<String> {
match groups_kv.get(pc_id).await {
Ok(Some(bytes)) => match serde_json::from_slice::<AgentGroups>(&bytes) {
Ok(g) => g.groups,
Err(e) => {
warn!(error = %e, pc_id, "decode AgentGroups — treating as no groups");
Vec::new()
}
},
Ok(None) => Vec::new(),
Err(e) => {
warn!(error = %e, pc_id, "read agent_groups — treating as no groups");
Vec::new()
}
}
}
async fn open_cfg(
state: &AppState,
) -> Result<async_nats::jetstream::kv::Store, (StatusCode, String)> {
state
.jetstream
.get_key_value(BUCKET_AGENT_CONFIG)
.await
.map_err(|e| {
warn!(error = %e, bucket = BUCKET_AGENT_CONFIG, "open agent_config KV bucket");
(
StatusCode::SERVICE_UNAVAILABLE,
format!("agent_config KV bucket unavailable: {e}"),
)
})
}
async fn open_groups(
state: &AppState,
) -> Result<async_nats::jetstream::kv::Store, (StatusCode, String)> {
state
.jetstream
.get_key_value(BUCKET_AGENT_GROUPS)
.await
.map_err(|e| {
warn!(error = %e, bucket = BUCKET_AGENT_GROUPS, "open agent_groups KV bucket");
(
StatusCode::SERVICE_UNAVAILABLE,
format!("agent_groups KV bucket unavailable: {e}"),
)
})
}
async fn read_scope_or_default(
kv: &async_nats::jetstream::kv::Store,
key: &str,
) -> Result<ConfigScope, (StatusCode, String)> {
match read_optional_scope(kv, key).await? {
Some(s) => Ok(s),
None => Ok(ConfigScope::default()),
}
}
async fn read_optional_scope(
kv: &async_nats::jetstream::kv::Store,
key: &str,
) -> Result<Option<ConfigScope>, (StatusCode, String)> {
match kv.get(key).await {
Ok(Some(bytes)) => serde_json::from_slice(&bytes).map(Some).map_err(|e| {
warn!(error = %e, key, "decode ConfigScope");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("decode ConfigScope at {key}: {e}"),
)
}),
Ok(None) => Ok(None),
Err(e) => {
warn!(error = %e, key, "read ConfigScope");
Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("read ConfigScope at {key}: {e}"),
))
}
}
}
async fn write_scope(
kv: &async_nats::jetstream::kv::Store,
key: &str,
scope: &ConfigScope,
) -> Result<(), (StatusCode, String)> {
let bytes = serde_json::to_vec(scope).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("encode ConfigScope: {e}"),
)
})?;
kv.put(key, bytes.into()).await.map_err(|e| {
warn!(error = %e, key, "write ConfigScope");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("write ConfigScope at {key}: {e}"),
)
})?;
Ok(())
}
async fn delete_key(
kv: &async_nats::jetstream::kv::Store,
key: &str,
) -> Result<(), (StatusCode, String)> {
kv.delete(key).await.map_err(|e| {
warn!(error = %e, key, "delete KV key");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("delete KV key {key}: {e}"),
)
})?;
Ok(())
}
#[cfg(test)]
mod tests {
use std::process::Stdio;
use std::time::Duration;
use super::*;
struct Harness {
kv: async_nats::jetstream::kv::Store,
_server: tokio::process::Child,
_storage: tempfile::TempDir,
}
async fn spawn_broker() -> (async_nats::Client, tokio::process::Child, tempfile::TempDir) {
let port = portpicker::pick_unused_port().expect("pick port");
let storage = tempfile::TempDir::new().expect("storage tempdir");
let server = tokio::process::Command::new("nats-server")
.arg("-js")
.arg("-p")
.arg(port.to_string())
.arg("-sd")
.arg(storage.path())
.stdout(Stdio::null())
.stderr(Stdio::null())
.kill_on_drop(true)
.spawn()
.expect("spawn nats-server (is it in PATH?)");
let url = format!("nats://127.0.0.1:{port}");
let mut client = None;
for _ in 0..50 {
match async_nats::connect(&url).await {
Ok(c) => {
client = Some(c);
break;
}
Err(_) => tokio::time::sleep(Duration::from_millis(100)).await,
}
}
(
client.expect("nats-server did not come up in 5s"),
server,
storage,
)
}
async fn harness() -> Harness {
let (client, server, storage) = spawn_broker().await;
let js = async_nats::jetstream::new(client);
let kv = js
.create_key_value(async_nats::jetstream::kv::Config {
bucket: BUCKET_AGENT_CONFIG.to_string(),
history: 5,
..Default::default()
})
.await
.expect("create agent_config bucket");
Harness {
kv,
_server: server,
_storage: storage,
}
}
async fn revision(kv: &async_nats::jetstream::kv::Store, key: &str) -> Option<u64> {
kv.entry(key).await.unwrap().map(|e| e.revision)
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn concurrent_sets_of_different_fields_both_survive() {
let h = harness().await;
let key = agent_config_group_key("canary");
let mut tasks = Vec::new();
for (field, value) in [
("target_version", "1.2.3"),
("heartbeat_interval", "15s"),
("host_perf_interval", "2m"),
("target_version_jitter", "30m"),
("process_perf_top_n", "20"),
("client_display_name", "tool"),
] {
let kv = h.kv.clone();
let key = key.clone();
tasks.push(tokio::spawn(async move {
update_field(&kv, &key, field, Some(value)).await.unwrap()
}));
}
for t in tasks {
assert!(t.await.unwrap().changed);
}
let scope = read_scope_or_default(&h.kv, &key).await.unwrap();
assert_eq!(scope.target_version.as_deref(), Some("1.2.3"));
assert_eq!(scope.heartbeat_interval.as_deref(), Some("15s"));
assert_eq!(scope.host_perf_interval.as_deref(), Some("2m"));
assert_eq!(scope.target_version_jitter.as_deref(), Some("30m"));
assert_eq!(scope.process_perf_top_n, Some(20));
assert_eq!(scope.client_display_name.as_deref(), Some("tool"));
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn unchanged_set_and_unset_do_not_bump_the_revision() {
let h = harness().await;
let key = agent_config_pc_key("PC-01");
let r = update_field(&h.kv, &key, "target_version", None)
.await
.unwrap();
assert!(!r.changed);
assert_eq!(revision(&h.kv, &key).await, None);
let r = update_field(&h.kv, &key, "heartbeat_interval", Some("15s"))
.await
.unwrap();
assert!(r.changed);
let rev = revision(&h.kv, &key).await.unwrap();
let r = update_field(&h.kv, &key, "heartbeat_interval", Some("15s"))
.await
.unwrap();
assert!(!r.changed);
assert_eq!(r.scope.heartbeat_interval.as_deref(), Some("15s"));
let r = update_field(&h.kv, &key, "target_version", None)
.await
.unwrap();
assert!(!r.changed);
assert_eq!(revision(&h.kv, &key).await, Some(rev));
let r = update_field(&h.kv, &key, "heartbeat_interval", None)
.await
.unwrap();
assert!(r.changed);
assert!(r.scope.heartbeat_interval.is_none());
assert!(revision(&h.kv, &key).await.unwrap() > rev);
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn field_update_recreates_a_deleted_row_and_global_row_can_be_cleared() {
let h = harness().await;
update_field(
&h.kv,
KEY_AGENT_CONFIG_GLOBAL,
"target_version",
Some("1.0.0"),
)
.await
.unwrap();
delete_key(&h.kv, KEY_AGENT_CONFIG_GLOBAL).await.unwrap();
assert!(
read_optional_scope(&h.kv, KEY_AGENT_CONFIG_GLOBAL)
.await
.unwrap()
.is_none()
);
let r = update_field(&h.kv, KEY_AGENT_CONFIG_GLOBAL, "target_version", None)
.await
.unwrap();
assert!(!r.changed);
let r = update_field(
&h.kv,
KEY_AGENT_CONFIG_GLOBAL,
"target_version",
Some("2.0.0"),
)
.await
.unwrap();
assert!(r.changed);
assert_eq!(r.scope.target_version.as_deref(), Some("2.0.0"));
}
#[test]
fn invalid_field_is_a_400_with_the_cli_message() {
for (field, value, needle) in [
("nope", Some("x"), "unknown field 'nope'"),
("heartbeat_interval", Some("soon"), "humantime duration"),
("process_perf_enabled", Some("maybe"), "expected true|false"),
(
"max_local_concurrent",
Some("0"),
"expected an integer >= 1",
),
] {
let (status, msg) = validate_field(field, value).unwrap_err();
assert_eq!(status, StatusCode::BAD_REQUEST);
assert!(msg.contains(needle), "{field}: {msg}");
}
assert!(validate_field("heartbeat_interval", Some("15s")).is_ok());
assert!(validate_field("heartbeat_interval", None).is_ok());
}
use axum::Router;
use axum::body::Body;
use axum::http::Request;
use tower::ServiceExt;
struct App {
router: Router,
nats: async_nats::Client,
kv: async_nats::jetstream::kv::Store,
groups_kv: async_nats::jetstream::kv::Store,
_server: tokio::process::Child,
_storage: tempfile::TempDir,
}
fn mint(sub: &str) -> String {
use jsonwebtoken::{Algorithm, EncodingKey, Header, encode};
let claims = crate::auth::Claims {
sub: sub.into(),
exp: 4_102_444_800,
aud: Some(crate::auth::EXPECTED_AUDIENCE.to_string()),
roles: Vec::new(),
allowed_features: None,
};
encode(
&Header::new(Algorithm::HS256),
&claims,
&EncodingKey::from_secret(crate::auth::signing_secret().as_bytes()),
)
.expect("mint")
}
async fn app() -> App {
let (nats, server, storage) = spawn_broker().await;
let js = async_nats::jetstream::new(nats.clone());
let mk = |bucket: &'static str| {
let js = js.clone();
async move {
js.create_key_value(async_nats::jetstream::kv::Config {
bucket: bucket.to_string(),
history: 5,
..Default::default()
})
.await
.expect("create bucket")
}
};
let kv = mk(BUCKET_AGENT_CONFIG).await;
let groups_kv = mk(BUCKET_AGENT_GROUPS).await;
let pool = sqlx::SqlitePool::connect("sqlite::memory:").await.unwrap();
sqlx::migrate!("./migrations").run(&pool).await.unwrap();
for (name, role, features) in [
("op", "operator", None),
("view", "viewer", None),
("locked", "operator", Some("[]")),
] {
sqlx::query(
"INSERT INTO users (username, password_hash, role, allowed_features) VALUES (?, 'x', ?, ?)",
)
.bind(name)
.bind(role)
.bind(features)
.execute(&pool)
.await
.unwrap();
}
let state = AppState {
pool: pool.clone(),
query_pool: pool.clone(),
commands: std::sync::Arc::new(crate::command_publisher::CommandPublisher::new(
nats.clone(),
None,
)),
nats: nats.clone(),
jetstream: js,
explode_spec_cache: Default::default(),
sql_view_cache: crate::api::view_sql::new_cache(),
group_cache: crate::api::group_sql::new_cache(),
mailer: None,
public_url: None,
nats_url: String::new(),
login_throttle: Default::default(),
};
let router = crate::api::router(state).layer(axum::middleware::from_fn_with_state(
pool,
crate::auth::verify,
));
App {
router,
nats,
kv,
groups_kv,
_server: server,
_storage: storage,
}
}
async fn call(
app: &App,
method: &str,
uri: &str,
who: Option<&str>,
body: Option<serde_json::Value>,
) -> (StatusCode, String) {
let mut req = Request::builder().method(method).uri(uri);
if let Some(who) = who {
req = req.header("authorization", format!("Bearer {}", mint(who)));
}
req = req.header("x-kanade-source", "cli");
let body = match body {
Some(v) => {
req = req.header("content-type", "application/json");
Body::from(v.to_string())
}
None => Body::empty(),
};
let resp = app
.router
.clone()
.oneshot(req.body(body).unwrap())
.await
.unwrap();
let status = resp.status();
let bytes = axum::body::to_bytes(resp.into_body(), usize::MAX)
.await
.unwrap();
(status, String::from_utf8_lossy(&bytes).into_owned())
}
fn field_url(scope: &str, field: &str) -> String {
match scope {
"global" => format!("/api/config/fields/{field}"),
s => format!("{}/fields/{field}", scope_url(s)),
}
}
fn scope_url(scope: &str) -> String {
match scope.split_once(':') {
None => "/api/config".to_string(),
Some(("group", n)) => format!("/api/groups/{n}/config"),
Some(("pc", n)) => format!("/api/pcs/{n}/config"),
_ => unreachable!(),
}
}
fn key_of(scope: &str) -> String {
match scope.split_once(':') {
None => KEY_AGENT_CONFIG_GLOBAL.to_string(),
Some(("group", n)) => agent_config_group_key(n),
Some(("pc", n)) => agent_config_pc_key(n),
_ => unreachable!(),
}
}
fn update(body: &str) -> FieldUpdate {
serde_json::from_str(body).unwrap_or_else(|e| panic!("{e}: {body}"))
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn field_routes_set_unset_and_skip_no_ops_on_every_scope() {
let app = app().await;
for scope in ["global", "group:canary", "pc:PC-01"] {
let key = key_of(scope);
let url = field_url(scope, "heartbeat_interval");
let set = || {
call(
&app,
"PUT",
&url,
Some("op"),
Some(serde_json::json!({"value": "15s"})),
)
};
let (st, body) = set().await;
assert_eq!(st, StatusCode::OK, "{scope}: {body}");
let u = update(&body);
assert!(u.changed, "{scope}");
assert_eq!(u.scope.heartbeat_interval.as_deref(), Some("15s"));
let rev = revision(&app.kv, &key).await.unwrap();
let (st, body) = set().await;
assert_eq!(st, StatusCode::OK);
assert!(!update(&body).changed, "{scope}");
assert_eq!(revision(&app.kv, &key).await, Some(rev), "{scope}");
let (st, body) = call(&app, "GET", &scope_url(scope), Some("view"), None).await;
assert_eq!(st, StatusCode::OK);
let got: ConfigScope = serde_json::from_str(&body).unwrap();
assert_eq!(got.heartbeat_interval.as_deref(), Some("15s"));
let (st, body) = call(&app, "DELETE", &url, Some("op"), None).await;
assert_eq!(st, StatusCode::OK);
assert!(update(&body).changed);
let rev = revision(&app.kv, &key).await.unwrap();
let (_, body) = call(&app, "DELETE", &url, Some("op"), None).await;
assert!(!update(&body).changed);
let other = field_url(scope, "target_version");
let (_, body) = call(&app, "DELETE", &other, Some("op"), None).await;
assert!(!update(&body).changed);
assert_eq!(revision(&app.kv, &key).await, Some(rev), "{scope}");
}
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn concurrent_field_puts_through_the_router_keep_every_write() {
let app = std::sync::Arc::new(app().await);
let mut tasks = Vec::new();
for (field, value) in [
("target_version", "1.2.3"),
("heartbeat_interval", "15s"),
("host_perf_interval", "2m"),
("target_version_jitter", "30m"),
("process_perf_top_n", "20"),
("client_display_name", "tool"),
] {
let app = app.clone();
tasks.push(tokio::spawn(async move {
call(
&app,
"PUT",
&field_url("group:canary", field),
Some("op"),
Some(serde_json::json!({ "value": value })),
)
.await
}));
}
for t in tasks {
let (st, body) = t.await.unwrap();
assert_eq!(st, StatusCode::OK, "{body}");
}
let (_, body) = call(&app, "GET", "/api/groups/canary/config", Some("view"), None).await;
let scope: ConfigScope = serde_json::from_str(&body).unwrap();
assert_eq!(scope.target_version.as_deref(), Some("1.2.3"));
assert_eq!(scope.heartbeat_interval.as_deref(), Some("15s"));
assert_eq!(scope.host_perf_interval.as_deref(), Some("2m"));
assert_eq!(scope.target_version_jitter.as_deref(), Some("30m"));
assert_eq!(scope.process_perf_top_n, Some(20));
assert_eq!(scope.client_display_name.as_deref(), Some("tool"));
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn invalid_field_is_a_400_over_http_and_writes_nothing() {
let app = app().await;
for (field, value, needle) in [
("nope", "x", "unknown field 'nope'"),
("heartbeat_interval", "soon", "humantime duration"),
("process_perf_enabled", "maybe", "expected true|false"),
] {
let (st, body) = call(
&app,
"PUT",
&field_url("pc:PC-01", field),
Some("op"),
Some(serde_json::json!({ "value": value })),
)
.await;
assert_eq!(st, StatusCode::BAD_REQUEST);
assert!(body.contains(needle), "{field}: {body}");
}
let (st, body) = call(
&app,
"DELETE",
&field_url("pc:PC-01", "nope"),
Some("op"),
None,
)
.await;
assert_eq!(st, StatusCode::BAD_REQUEST);
assert!(body.contains("unknown field 'nope'"), "{body}");
assert_eq!(revision(&app.kv, &key_of("pc:PC-01")).await, None);
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn mutations_need_an_operator_with_the_config_page() {
let app = app().await;
let body = || Some(serde_json::json!({"value": "15s"}));
for scope in ["global", "group:canary", "pc:PC-01"] {
let url = field_url(scope, "heartbeat_interval");
let (st, _) = call(&app, "PUT", &url, None, body()).await;
assert_eq!(st, StatusCode::UNAUTHORIZED, "{scope}");
for who in ["view", "locked"] {
let (st, _) = call(&app, "PUT", &url, Some(who), body()).await;
assert_eq!(st, StatusCode::FORBIDDEN, "{scope} PUT as {who}");
let (st, _) = call(&app, "DELETE", &url, Some(who), None).await;
assert_eq!(st, StatusCode::FORBIDDEN, "{scope} DELETE as {who}");
let (st, _) = call(&app, "DELETE", &scope_url(scope), Some(who), None).await;
assert_eq!(st, StatusCode::FORBIDDEN, "{scope} clear as {who}");
}
assert_eq!(revision(&app.kv, &key_of(scope)).await, None, "{scope}");
}
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn global_scope_can_be_cleared() {
let app = app().await;
call(
&app,
"PUT",
"/api/config/fields/target_version",
Some("op"),
Some(serde_json::json!({"value": "1.0.0"})),
)
.await;
let (st, _) = call(&app, "DELETE", "/api/config", Some("op"), None).await;
assert_eq!(st, StatusCode::NO_CONTENT);
let (_, body) = call(&app, "GET", "/api/config", Some("view"), None).await;
let scope: ConfigScope = serde_json::from_str(&body).unwrap();
assert_eq!(scope, ConfigScope::default());
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn every_mutation_is_audited_under_the_callers_account() {
use futures::StreamExt;
let app = app().await;
let mut sub = app.nats.subscribe("audit.>").await.unwrap();
app.nats.flush().await.unwrap();
let value = |v: &str| Some(serde_json::json!({ "value": v }));
call(
&app,
"PUT",
"/api/groups/canary/config/fields/heartbeat_interval",
Some("op"),
value("15s"),
)
.await;
call(
&app,
"PUT",
"/api/groups/canary/config/fields/heartbeat_interval",
Some("op"),
value("15s"),
)
.await;
call(
&app,
"DELETE",
"/api/groups/canary/config/fields/heartbeat_interval",
Some("op"),
None,
)
.await;
call(
&app,
"PUT",
"/api/pcs/PC-01/config",
Some("op"),
Some(serde_json::json!({"target_version": "2.0.0"})),
)
.await;
call(&app, "DELETE", "/api/pcs/PC-01/config", Some("op"), None).await;
call(&app, "DELETE", "/api/config", Some("op"), None).await;
let expected = [
("audit.operator.config_field_set.groups.canary", Some(true)),
("audit.operator.config_field_set.groups.canary", Some(false)),
(
"audit.operator.config_field_unset.groups.canary",
Some(true),
),
("audit.operator.config_scope_put.pcs.PC-01", None),
("audit.operator.config_scope_clear.pcs.PC-01", None),
("audit.operator.config_scope_clear.global", None),
];
for (subject, changed) in expected {
let msg = tokio::time::timeout(Duration::from_secs(5), sub.next())
.await
.expect("audit event")
.expect("subscription open");
assert_eq!(msg.subject.as_str(), subject);
let v: serde_json::Value = serde_json::from_slice(&msg.payload).unwrap();
assert_eq!(v["payload"]["sub"], "op", "{v}");
assert_eq!(v["payload"]["source"], "cli", "{v}");
if let Some(c) = changed {
assert_eq!(v["payload"]["changed"], c, "{v}");
}
}
}
#[tokio::test]
#[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
async fn effective_config_reports_the_pcs_groups() {
let app = app().await;
app.groups_kv
.put("PC-01", br#"{"groups":["canary"]}"#.as_slice().into())
.await
.unwrap();
call(
&app,
"PUT",
"/api/groups/canary/config/fields/heartbeat_interval",
Some("op"),
Some(serde_json::json!({"value": "15s"})),
)
.await;
let (st, body) = call(
&app,
"GET",
"/api/agents/PC-01/effective_config",
Some("view"),
None,
)
.await;
assert_eq!(st, StatusCode::OK, "{body}");
let v: serde_json::Value = serde_json::from_str(&body).unwrap();
assert_eq!(v["my_groups"], serde_json::json!(["canary"]));
}
}