use axum::Json;
use axum::extract::{Path, State};
use axum::http::StatusCode;
use kanade_shared::config::MailSection;
use kanade_shared::kv::{BUCKET_SERVER_SETTINGS, KEY_SERVER_SETTINGS};
use kanade_shared::kv_cas;
use kanade_shared::wire::{
AgentInstallSection, MAX_AGENT_PRUNE_DAYS, MAX_CHECK_STATUS_STALE_DAYS,
MAX_COLLECT_RETENTION_DAYS, MAX_OBJECT_STORE_CAP_MIB, MAX_OBJECT_STORE_TOTAL_MIB,
MAX_RESULT_OUTPUT_RETENTION_DAYS, MAX_SESSION_TTL_HOURS, MAX_SUPPORT_UNLOCK_TTL_MINUTES,
NatsAuthMode, ObjectStoreCaps, ServerSettings, SupportCode,
};
use lettre::message::Mailbox;
use serde_json::{Map, Value};
use tracing::{info, warn};
use crate::api::AppState;
use crate::audit;
use crate::audit::Caller;
pub async fn get(State(s): State<AppState>) -> Result<Json<ServerSettings>, (StatusCode, String)> {
let kv = open_bucket(&s).await?;
match kv.get(KEY_SERVER_SETTINGS).await {
Ok(Some(bytes)) => serde_json::from_slice::<ServerSettings>(&bytes)
.map(|s| Json(s.redacted()))
.map_err(|e| {
warn!(error = %e, "decode server_settings");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("stored server_settings is corrupt: {e}"),
)
}),
Ok(None) => Ok(Json(ServerSettings::default())),
Err(e) => {
warn!(error = %e, "read server_settings");
Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("read server_settings: {e}"),
))
}
}
}
pub async fn defaults() -> Json<ServerSettings> {
Json(ServerSettings::defaults())
}
pub async fn put(
State(s): State<AppState>,
caller: Caller,
Json(incoming): Json<Value>,
) -> Result<Json<ServerSettings>, (StatusCode, String)> {
let incoming = match incoming {
Value::Object(m) => m,
_ => {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"server_settings body must be a JSON object".to_string(),
));
}
};
let typed: ServerSettings =
serde_json::from_value(Value::Object(incoming.clone())).map_err(|e| {
(
StatusCode::UNPROCESSABLE_ENTITY,
format!("invalid server_settings body: {e}"),
)
})?;
validate(&typed)?;
validate_user_pair_update(&incoming)?;
let typed = normalize(typed);
let prune_value = typed.agent_prune_days.map(Value::from);
let collect_value = typed.collect_retention_days.map(Value::from);
let result_output_value = typed.result_output_retention_days.map(Value::from);
let session_ttl_value = typed.session_ttl_hours.map(Value::from);
let check_stale_value = typed.check_status_stale_days.map(Value::from);
let controller_value = typed.controller_group.clone().map(Value::String);
let nats_mode_value = typed.nats_auth_mode.map(|m| Value::from(m.as_str()));
let now = chrono::Utc::now();
let mail_value = match typed.mail.as_ref() {
Some(m) => Some(serde_json::to_value(m).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("encode mail settings: {e}"),
)
})?),
None => None,
};
let caps_value = match typed.object_store_caps.as_ref() {
Some(c) => Some(serde_json::to_value(c).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("encode object_store_caps: {e}"),
)
})?),
None => None,
};
let kv = open_bucket(&s).await?;
{
let stored_caps: Option<ObjectStoreCaps> = match kv.get(KEY_SERVER_SETTINGS).await {
Ok(Some(bytes)) => serde_json::from_slice::<Value>(&bytes)
.ok()
.and_then(|v| v.get("object_store_caps").cloned())
.and_then(|v| serde_json::from_value(v).ok()),
_ => None,
};
let merged_caps = if incoming.contains_key("object_store_caps") {
typed.object_store_caps.clone()
} else {
stored_caps
};
if let Some(c) = merged_caps.as_ref() {
validate_object_store_caps(c)?;
}
}
let mut collect_changed = false;
let mut result_output_changed = false;
let mut caps_changed = false;
let merged_map =
kv_cas::read_modify_write::<Map<String, Value>, _>(&kv, KEY_SERVER_SETTINGS, |obj| {
let mut changed = false;
changed |= merge_field(obj, &incoming, "agent_prune_days", prune_value.clone());
collect_changed = merge_field(
obj,
&incoming,
"collect_retention_days",
collect_value.clone(),
);
changed |= collect_changed;
result_output_changed = merge_field(
obj,
&incoming,
"result_output_retention_days",
result_output_value.clone(),
);
changed |= result_output_changed;
changed |= merge_field(
obj,
&incoming,
"session_ttl_hours",
session_ttl_value.clone(),
);
changed |= merge_field(
obj,
&incoming,
"check_status_stale_days",
check_stale_value.clone(),
);
changed |= merge_field(obj, &incoming, "controller_group", controller_value.clone());
changed |= merge_field(obj, &incoming, "mail", mail_value.clone());
changed |= merge_nats_auth_mode(obj, &incoming, nats_mode_value.clone(), now);
caps_changed = merge_field(obj, &incoming, "object_store_caps", caps_value.clone());
changed |= caps_changed;
changed |= merge_agent_install(obj, &incoming);
changed
})
.await
.map_err(|e| {
warn!(error = %format!("{e:#}"), "write server_settings");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("write server_settings: {e}"),
)
})?;
let doc = Value::Object(merged_map);
let merged: ServerSettings = serde_json::from_value(doc.clone()).map_err(|e| {
warn!(error = %e, "decode merged server_settings");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("merged server_settings is corrupt: {e}"),
)
})?;
info!(
agent_prune_days = ?merged.agent_prune_days,
collect_retention_days = ?merged.collect_retention_days,
result_output_retention_days = ?merged.result_output_retention_days,
session_ttl_hours = ?merged.session_ttl_hours,
controller_group = ?merged.controller_group,
mail_configured = merged.mail.is_some(),
agent_install_configured = merged.agent_install.is_some(),
"server_settings merged",
);
if collect_changed {
let days = merged.effective_collect_retention_days();
match kanade_shared::bootstrap::reconcile_collect_retention(&s.jetstream, days).await {
Ok(true) => info!(
collect_retention_days = days,
"collect retention applied to Object Store"
),
Ok(false) => {}
Err(e) => warn!(
error = %format!("{e:#}"), collect_retention_days = days,
"collect retention: applied to KV but reconcile of the Object Store max_age failed; \
will be applied on the next backend restart",
),
}
}
if result_output_changed {
let days = merged.effective_result_output_retention_days();
match kanade_shared::bootstrap::reconcile_object_store_max_age(
&s.jetstream,
kanade_shared::kv::OBJECT_RESULT_OUTPUT,
days,
)
.await
{
Ok(true) => info!(
result_output_retention_days = days,
"result_output retention reconciled after save"
),
Ok(false) => {}
Err(e) => warn!(
error = %format!("{e:#}"), result_output_retention_days = days,
"result_output retention reconcile after save failed"
),
}
}
if caps_changed {
for (bucket, cap_mib) in merged.effective_object_store_caps().effective_all() {
match kanade_shared::bootstrap::reconcile_object_store_max_bytes(
&s.jetstream,
bucket,
cap_mib,
)
.await
{
Ok(true) => info!(bucket, cap_mib, "object store cap applied"),
Ok(false) => {}
Err(e) => warn!(
error = %format!("{e:#}"), bucket, cap_mib,
"object store cap: applied to KV but reconcile failed; \
will be applied on the next backend restart",
),
}
}
}
audit::record(
&s.nats,
"operator",
"server_settings_set",
Some(KEY_SERVER_SETTINGS),
Some(&caller),
redact_secrets(doc),
)
.await;
Ok(Json(merged.redacted()))
}
fn merge_nats_auth_mode(
obj: &mut Map<String, Value>,
incoming: &Map<String, Value>,
value: Option<Value>,
now: chrono::DateTime<chrono::Utc>,
) -> bool {
let effective = |obj: &Map<String, Value>| {
obj.get("nats_auth_mode")
.and_then(|v| serde_json::from_value::<NatsAuthMode>(v.clone()).ok())
.unwrap_or_default()
};
let before = effective(obj);
let mut changed = merge_field(obj, incoming, "nats_auth_mode", value);
if effective(obj) != before {
obj.insert(
"nats_auth_mode_changed_at".to_string(),
Value::String(now.to_rfc3339_opts(chrono::SecondsFormat::Millis, true)),
);
changed = true;
}
changed
}
fn redact_secrets(mut doc: Value) -> Value {
if let Some(codes) = doc.get_mut("support_codes").and_then(Value::as_array_mut) {
for c in codes {
if let Some(obj) = c.as_object_mut() {
obj.remove("hash");
}
}
}
if let Some(ai) = doc.get_mut("agent_install").and_then(Value::as_object_mut) {
for (secret, flag) in [
("nats_token", "nats_token_set"),
("nats_user", "nats_user_set"),
("nats_password", "nats_password_set"),
] {
let had = ai.remove(secret).is_some_and(|v| v.is_string());
ai.insert(flag.to_string(), Value::Bool(had));
}
}
doc
}
fn merge_field(
obj: &mut Map<String, Value>,
incoming: &Map<String, Value>,
key: &str,
value: Option<Value>,
) -> bool {
if !incoming.contains_key(key) {
return false;
}
match value {
Some(v) => {
if obj.get(key) == Some(&v) {
return false;
}
obj.insert(key.to_string(), v);
true
}
None => obj.remove(key).is_some(),
}
}
const INSTALL_SECRET_FLAGS: [&str; 3] = ["nats_token_set", "nats_user_set", "nats_password_set"];
fn merge_agent_install(obj: &mut Map<String, Value>, incoming: &Map<String, Value>) -> bool {
let Some(inc) = incoming.get("agent_install") else {
return false;
};
if inc.is_null() {
return obj.remove("agent_install").is_some();
}
let Some(inc_obj) = inc.as_object() else {
return false;
};
let mut merged = match obj.get("agent_install").and_then(Value::as_object) {
Some(o) => o.clone(),
None => Map::new(),
};
for (k, v) in inc_obj {
if INSTALL_SECRET_FLAGS.contains(&k.as_str()) {
continue;
}
if v.is_null() {
merged.remove(k);
} else {
merged.insert(k.clone(), v.clone());
}
}
for flag in INSTALL_SECRET_FLAGS {
merged.remove(flag);
}
if merged.is_empty() {
return obj.remove("agent_install").is_some();
}
let new_val = Value::Object(merged);
if obj.get("agent_install") == Some(&new_val) {
return false;
}
obj.insert("agent_install".to_string(), new_val);
true
}
fn validate(s: &ServerSettings) -> Result<(), (StatusCode, String)> {
if let Some(days) = s.agent_prune_days {
if days == 0 {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"agent_prune_days must be >= 1; omit it or send null to disable pruning"
.to_string(),
));
}
if days > MAX_AGENT_PRUNE_DAYS {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!("agent_prune_days must be <= {MAX_AGENT_PRUNE_DAYS} (100 years)"),
));
}
}
if let Some(days) = s.collect_retention_days {
if days == 0 {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"collect_retention_days must be >= 1; omit it or send null to use the default"
.to_string(),
));
}
if days > MAX_COLLECT_RETENTION_DAYS {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!(
"collect_retention_days must be <= {MAX_COLLECT_RETENTION_DAYS} (10 years)"
),
));
}
}
if let Some(days) = s.result_output_retention_days {
if days == 0 {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"result_output_retention_days must be >= 1; omit it or send null to use the default"
.to_string(),
));
}
if days > MAX_RESULT_OUTPUT_RETENTION_DAYS {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!(
"result_output_retention_days must be <= {MAX_RESULT_OUTPUT_RETENTION_DAYS} (STREAM_RESULTS' own retention — beyond it there is nothing to replay)"
),
));
}
}
if let Some(hours) = s.session_ttl_hours {
if hours == 0 {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"session_ttl_hours must be >= 1; omit it or send null to use the default"
.to_string(),
));
}
if hours > MAX_SESSION_TTL_HOURS {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!("session_ttl_hours must be <= {MAX_SESSION_TTL_HOURS} (365 days)"),
));
}
}
if let Some(days) = s.check_status_stale_days {
if days > MAX_CHECK_STATUS_STALE_DAYS {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!(
"check_status_stale_days must be <= {MAX_CHECK_STATUS_STALE_DAYS} (10 years)"
),
));
}
}
if let Some(g) = s.controller_group.as_deref()
&& g.trim().is_empty()
{
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"controller_group must be a non-empty group name; omit it or send null to unset"
.to_string(),
));
}
if let Some(m) = s.mail.as_ref() {
validate_mail(m)?;
}
if let Some(c) = s.object_store_caps.as_ref() {
validate_object_store_caps(c)?;
}
if let Some(ai) = s.agent_install.as_ref() {
validate_agent_install(ai)?;
}
Ok(())
}
fn validate_agent_install(ai: &AgentInstallSection) -> Result<(), (StatusCode, String)> {
if let Some(url) = ai.nats_url.as_deref()
&& (url.is_empty() || url.contains('\'') || url.contains('\n') || url.contains('\r'))
{
return Err((
StatusCode::BAD_REQUEST,
"agent_install.nats_url must be non-empty and contain no single quote or newline \
(TOML literal-string safety)"
.into(),
));
}
if let Some(token) = ai.nats_token.as_deref()
&& (token.contains('\n') || token.contains('\r'))
{
return Err((
StatusCode::BAD_REQUEST,
"agent_install.nats_token must not contain a newline".into(),
));
}
for (name, v) in [
("nats_user", ai.nats_user.as_deref()),
("nats_password", ai.nats_password.as_deref()),
] {
if let Some(v) = v
&& (v.is_empty() || v.contains('\n') || v.contains('\r'))
{
return Err((
StatusCode::BAD_REQUEST,
format!("agent_install.{name} must be non-empty and must not contain a newline"),
));
}
}
Ok(())
}
fn validate_user_pair_update(incoming: &Map<String, Value>) -> Result<(), (StatusCode, String)> {
let Some(ai) = incoming.get("agent_install").and_then(Value::as_object) else {
return Ok(());
};
let state = |k: &str| ai.get(k).map(|v| !v.is_null());
match (state("nats_user"), state("nats_password")) {
(None, None) | (Some(false), Some(false)) | (Some(true), Some(true)) => Ok(()),
_ => Err((
StatusCode::BAD_REQUEST,
"agent_install.nats_user and agent_install.nats_password are a pair: send both \
(to set or rotate) or both as null (to clear), or omit both to keep the stored pair"
.into(),
)),
}
}
fn validate_object_store_caps(c: &ObjectStoreCaps) -> Result<(), (StatusCode, String)> {
for (field, v) in [
("result_output_mib", c.result_output_mib),
("agent_releases_mib", c.agent_releases_mib),
("app_packages_mib", c.app_packages_mib),
("scripts_mib", c.scripts_mib),
("collections_mib", c.collections_mib),
] {
let Some(v) = v else { continue };
if v == 0 {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!(
"object_store_caps.{field} must be >= 1; omit it or send null to use the default"
),
));
}
if v > MAX_OBJECT_STORE_CAP_MIB {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!("object_store_caps.{field} must be <= {MAX_OBJECT_STORE_CAP_MIB} (50 GiB)"),
));
}
}
let total: u64 = c.effective_all().iter().map(|(_, v)| *v as u64).sum();
if total > MAX_OBJECT_STORE_TOTAL_MIB as u64 {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!(
"object_store_caps total {total} MiB exceeds the broker-wide object-store budget \
{MAX_OBJECT_STORE_TOTAL_MIB} MiB (max_file_store minus stream reservations); \
lower one or more buckets"
),
));
}
Ok(())
}
fn validate_mail(m: &MailSection) -> Result<(), (StatusCode, String)> {
if m.host.trim().is_empty() {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"mail.host must not be empty".to_string(),
));
}
if m.port == 0 {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"mail.port must be between 1 and 65535".to_string(),
));
}
if m.from.trim().parse::<Mailbox>().is_err() {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!("mail.from is not a valid email address: {:?}", m.from),
));
}
Ok(())
}
fn normalize(mut s: ServerSettings) -> ServerSettings {
if let Some(g) = s.controller_group.as_mut() {
*g = g.trim().to_string();
}
if let Some(m) = s.mail.as_mut() {
m.host = m.host.trim().to_string();
m.from = m.from.trim().to_string();
m.username = m
.username
.as_deref()
.map(str::trim)
.filter(|u| !u.is_empty())
.map(String::from);
}
s
}
pub(crate) async fn load(s: &AppState) -> anyhow::Result<ServerSettings> {
load_from_js(&s.jetstream).await
}
pub(crate) async fn load_from_js(
js: &async_nats::jetstream::Context,
) -> anyhow::Result<ServerSettings> {
use anyhow::Context;
let kv = js
.get_key_value(BUCKET_SERVER_SETTINGS)
.await
.context("open server_settings KV")?;
match kv
.get(KEY_SERVER_SETTINGS)
.await
.context("get server_settings")?
{
Some(bytes) => serde_json::from_slice(&bytes).context("decode server_settings"),
None => Ok(ServerSettings::default()),
}
}
#[derive(serde::Deserialize)]
pub struct SupportCodeBody {
pub code: String,
#[serde(default)]
pub label: Option<String>,
#[serde(default)]
pub ttl_minutes: Option<u32>,
#[serde(default)]
pub disabled: bool,
}
const MIN_SUPPORT_CODE_LEN: usize = 8;
pub async fn put_support_code(
State(s): State<AppState>,
caller: Caller,
Path(scope): Path<String>,
Json(body): Json<SupportCodeBody>,
) -> Result<Json<ServerSettings>, (StatusCode, String)> {
let scope = scope.trim().to_string();
validate_support_code(&scope, &body)?;
let hash = hash_support_code(&body.code)?;
let entry = SupportCode {
scope: scope.clone(),
hash,
label: body.label.map(|l| l.trim().to_string()),
ttl_minutes: body.ttl_minutes,
disabled: body.disabled,
};
let entry_value = serde_json::to_value(&entry).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("encode support code: {e}"),
)
})?;
let kv = open_bucket(&s).await?;
let merged_map =
kv_cas::read_modify_write::<Map<String, Value>, _>(&kv, KEY_SERVER_SETTINGS, |obj| {
let codes = obj
.entry("support_codes".to_string())
.or_insert_with(|| Value::Array(Vec::new()));
if !codes.is_array() {
*codes = Value::Array(Vec::new());
}
let arr = codes.as_array_mut().expect("just ensured array");
arr.retain(|c| c.get("scope").and_then(Value::as_str) != Some(scope.as_str()));
arr.push(entry_value.clone());
true
})
.await
.map_err(|e| {
warn!(error = %format!("{e:#}"), "write support code");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("write support code: {e}"),
)
})?;
crate::support_unlock_config::sync_or_warn(&s.jetstream).await;
let merged = decode_merged(Value::Object(merged_map))?;
info!(
scope = %entry.scope,
disabled = entry.disabled,
"support code set",
);
audit::record(
&s.nats,
"operator",
"support_code_set",
Some(&entry.scope),
Some(&caller),
serde_json::json!({
"scope": entry.scope,
"label": entry.label,
"ttl_minutes": entry.ttl_minutes,
"disabled": entry.disabled,
}),
)
.await;
Ok(Json(merged.redacted()))
}
pub async fn delete_support_code(
State(s): State<AppState>,
caller: Caller,
Path(scope): Path<String>,
) -> Result<Json<ServerSettings>, (StatusCode, String)> {
let scope = scope.trim().to_string();
let kv = open_bucket(&s).await?;
let merged_map =
kv_cas::read_modify_write::<Map<String, Value>, _>(&kv, KEY_SERVER_SETTINGS, |obj| {
let Some(arr) = obj.get_mut("support_codes").and_then(Value::as_array_mut) else {
return false;
};
let before = arr.len();
arr.retain(|c| c.get("scope").and_then(Value::as_str) != Some(scope.as_str()));
arr.len() != before
})
.await
.map_err(|e| {
warn!(error = %format!("{e:#}"), "delete support code");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("delete support code: {e}"),
)
})?;
crate::support_unlock_config::sync_or_warn(&s.jetstream).await;
let merged = decode_merged(Value::Object(merged_map))?;
info!(scope = %scope, "support code deleted");
audit::record(
&s.nats,
"operator",
"support_code_deleted",
Some(&scope),
Some(&caller),
serde_json::json!({ "scope": scope }),
)
.await;
Ok(Json(merged.redacted()))
}
fn validate_support_code(scope: &str, body: &SupportCodeBody) -> Result<(), (StatusCode, String)> {
if !kanade_shared::manifest::is_valid_resource_id(scope) {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"scope must be a slug ([A-Za-z0-9._-]) matching a job's client.unlock".to_string(),
));
}
if body.code != body.code.trim() {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"code must not start or end with whitespace".to_string(),
));
}
if body.code.chars().count() < MIN_SUPPORT_CODE_LEN {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!("code must be at least {MIN_SUPPORT_CODE_LEN} characters"),
));
}
if let Some(ttl) = body.ttl_minutes {
if ttl == 0 || ttl > MAX_SUPPORT_UNLOCK_TTL_MINUTES {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!("ttl_minutes must be 1..={MAX_SUPPORT_UNLOCK_TTL_MINUTES}"),
));
}
}
if let Some(label) = body.label.as_ref() {
if label.trim().is_empty() {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"label must not be blank when set; omit it instead".to_string(),
));
}
}
Ok(())
}
fn hash_support_code(code: &str) -> Result<String, (StatusCode, String)> {
use argon2::password_hash::{PasswordHasher, SaltString, rand_core::OsRng};
let salt = SaltString::generate(&mut OsRng);
argon2::Argon2::default()
.hash_password(code.as_bytes(), &salt)
.map(|h| h.to_string())
.map_err(|e| {
warn!(error = %e, "hash support code");
(
StatusCode::INTERNAL_SERVER_ERROR,
"failed to hash the support code".to_string(),
)
})
}
fn decode_merged(doc: Value) -> Result<ServerSettings, (StatusCode, String)> {
serde_json::from_value(doc).map_err(|e| {
warn!(error = %e, "decode merged server_settings");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("merged server_settings is corrupt: {e}"),
)
})
}
async fn open_bucket(
s: &AppState,
) -> Result<async_nats::jetstream::kv::Store, (StatusCode, String)> {
s.jetstream
.get_key_value(BUCKET_SERVER_SETTINGS)
.await
.map_err(|e| {
warn!(error = %e, bucket = BUCKET_SERVER_SETTINGS, "open server_settings KV bucket");
(
StatusCode::SERVICE_UNAVAILABLE,
format!("server_settings KV bucket unavailable: {e}"),
)
})
}
#[cfg(test)]
mod tests {
use axum::http::StatusCode;
use kanade_shared::config::{MailEncryption, MailSection};
use serde_json::{Map, Value, json};
use super::{
MAX_OBJECT_STORE_CAP_MIB, MIN_SUPPORT_CODE_LEN, ObjectStoreCaps, ServerSettings,
SupportCodeBody, hash_support_code, merge_agent_install, merge_field, merge_nats_auth_mode,
normalize, redact_secrets, validate, validate_support_code, validate_user_pair_update,
};
fn obj(v: Value) -> Map<String, Value> {
v.as_object().expect("object literal").clone()
}
fn code_body(code: &str) -> SupportCodeBody {
SupportCodeBody {
code: code.to_string(),
label: None,
ttl_minutes: None,
disabled: false,
}
}
#[test]
fn support_code_validation_rejects_the_unusable() {
assert!(validate_support_code("support", &code_body("hunter2!!")).is_ok());
assert!(validate_support_code("has space", &code_body("hunter2!!")).is_err());
assert!(validate_support_code("", &code_body("hunter2!!")).is_err());
let short = "a".repeat(MIN_SUPPORT_CODE_LEN - 1);
assert!(validate_support_code("support", &code_body(&short)).is_err());
assert!(validate_support_code("support", &code_body(" hunter2!! ")).is_err());
let mut ttl = code_body("hunter2!!");
ttl.ttl_minutes = Some(0);
assert!(validate_support_code("support", &ttl).is_err());
ttl.ttl_minutes = Some(u32::MAX);
assert!(validate_support_code("support", &ttl).is_err());
let mut label = code_body("hunter2!!");
label.label = Some(" ".into());
assert!(validate_support_code("support", &label).is_err());
}
#[test]
fn hashed_code_verifies_and_is_salted() {
use argon2::{Argon2, PasswordHash, PasswordVerifier};
let a = hash_support_code("hunter2!!").unwrap();
let b = hash_support_code("hunter2!!").unwrap();
assert_ne!(a, b, "each hash must carry its own salt");
for h in [&a, &b] {
let parsed = PasswordHash::new(h).unwrap();
assert!(
Argon2::default()
.verify_password(b"hunter2!!", &parsed)
.is_ok()
);
assert!(
Argon2::default()
.verify_password(b"hunter3!!", &parsed)
.is_err()
);
}
}
#[test]
fn audit_copy_carries_no_hashes() {
let doc = json!({
"agent_prune_days": 7,
"support_codes": [
{"scope":"support","hash":"$argon2id$secret","label":"desk"},
{"scope":"admin","hash":"$argon2id$other"},
],
});
let redacted = redact_secrets(doc);
let text = redacted.to_string();
assert!(!text.contains("argon2"), "audit leaked a hash: {text}");
assert_eq!(redacted["agent_prune_days"], 7);
assert_eq!(redacted["support_codes"][0]["scope"], "support");
assert_eq!(redacted["support_codes"][0]["label"], "desk");
}
#[test]
fn audit_copy_carries_no_install_token() {
let doc = json!({
"agent_install": {"nats_url":"nats://b:4222","nats_token":"s3cret"},
});
let redacted = redact_secrets(doc);
let text = redacted.to_string();
assert!(!text.contains("s3cret"), "audit leaked the token: {text}");
assert_eq!(redacted["agent_install"]["nats_url"], "nats://b:4222");
assert_eq!(redacted["agent_install"]["nats_token_set"], true);
let doc = json!({ "agent_install": {"nats_url":"nats://b:4222"} });
let redacted = redact_secrets(doc);
assert_eq!(redacted["agent_install"]["nats_token_set"], false);
}
#[test]
fn audit_redaction_tolerates_a_missing_or_odd_field() {
assert_eq!(
redact_secrets(json!({"agent_prune_days": 7}))["agent_prune_days"],
7
);
assert_eq!(
redact_secrets(json!({"support_codes": "nonsense"}))["support_codes"],
"nonsense"
);
assert_eq!(
redact_secrets(json!({"agent_install": "nonsense"}))["agent_install"],
"nonsense"
);
}
fn sample_mail() -> MailSection {
MailSection {
host: "smtp.example.com".into(),
port: 587,
encryption: MailEncryption::Starttls,
from: "kanade-noreply@example.com".into(),
username: None,
}
}
fn nats_mode_doc(mode: Option<&str>, at: Option<&str>) -> Map<String, Value> {
let mut m = Map::new();
if let Some(mode) = mode {
m.insert("nats_auth_mode".into(), json!(mode));
}
if let Some(at) = at {
m.insert("nats_auth_mode_changed_at".into(), json!(at));
}
m
}
#[test]
fn a_mode_change_stamps_the_switch_time_and_a_resave_does_not() {
let now = chrono::Utc::now();
let mut doc = nats_mode_doc(None, None);
let inc = obj(json!({"nats_auth_mode": "users"}));
assert!(merge_nats_auth_mode(
&mut doc,
&inc,
Some(json!("users")),
now
));
assert_eq!(doc["nats_auth_mode"], "users");
assert!(doc.contains_key("nats_auth_mode_changed_at"));
let stamp = doc["nats_auth_mode_changed_at"].clone();
let later = now + chrono::Duration::minutes(30);
let inc = obj(json!({
"nats_auth_mode": "users",
"nats_auth_mode_changed_at": "2099-01-01T00:00:00Z"
}));
assert!(!merge_nats_auth_mode(
&mut doc,
&inc,
Some(json!("users")),
later
));
assert_eq!(doc["nats_auth_mode_changed_at"], stamp);
let inc = obj(json!({"nats_auth_mode": "token"}));
assert!(merge_nats_auth_mode(
&mut doc,
&inc,
Some(json!("token")),
later
));
assert_ne!(doc["nats_auth_mode_changed_at"], stamp);
}
#[test]
fn saving_token_over_unset_is_not_a_switch() {
let mut doc = nats_mode_doc(None, None);
let inc = obj(json!({"nats_auth_mode": "token"}));
merge_nats_auth_mode(&mut doc, &inc, Some(json!("token")), chrono::Utc::now());
assert!(!doc.contains_key("nats_auth_mode_changed_at"));
let mut doc = nats_mode_doc(Some("token"), None);
let inc = obj(json!({"nats_auth_mode": null}));
merge_nats_auth_mode(&mut doc, &inc, None, chrono::Utc::now());
assert!(!doc.contains_key("nats_auth_mode_changed_at"));
}
#[test]
fn a_request_that_omits_the_mode_never_touches_it_or_the_stamp() {
let mut doc = nats_mode_doc(Some("users"), Some("2026-01-01T00:00:00.000Z"));
let before = doc.clone();
let inc = obj(json!({"nats_auth_mode_changed_at": "2099-01-01T00:00:00Z"}));
assert!(!merge_nats_auth_mode(
&mut doc,
&inc,
None,
chrono::Utc::now()
));
assert_eq!(doc, before);
}
#[test]
fn merge_key_absent_is_left_untouched() {
let mut stored = obj(json!({ "a": 1, "unknown_future": true }));
assert!(!merge_field(
&mut stored,
&obj(json!({})),
"a",
Some(json!(2))
));
assert_eq!(stored.get("a"), Some(&json!(1)));
assert_eq!(stored.get("unknown_future"), Some(&json!(true)));
}
#[test]
fn merge_present_value_overwrites() {
let mut stored = obj(json!({ "a": 1 }));
assert!(merge_field(
&mut stored,
&obj(json!({ "a": 2 })),
"a",
Some(json!(2))
));
assert_eq!(stored.get("a"), Some(&json!(2)));
}
#[test]
fn merge_present_same_value_is_noop() {
let mut stored = obj(json!({ "a": 1 }));
assert!(!merge_field(
&mut stored,
&obj(json!({ "a": 1 })),
"a",
Some(json!(1))
));
assert_eq!(stored.get("a"), Some(&json!(1)));
}
#[test]
fn merge_present_null_unsets() {
let mut stored = obj(json!({ "a": 1 }));
assert!(merge_field(
&mut stored,
&obj(json!({ "a": Value::Null })),
"a",
None
));
assert!(!stored.contains_key("a"));
assert!(!merge_field(
&mut stored,
&obj(json!({ "a": Value::Null })),
"a",
None
));
}
#[test]
fn validate_rejects_zero_prune_days() {
let s = ServerSettings {
agent_prune_days: Some(0),
..Default::default()
};
assert!(validate(&s).is_err());
}
#[test]
fn validate_rejects_zero_or_oversize_result_output_retention() {
use kanade_shared::wire::MAX_RESULT_OUTPUT_RETENTION_DAYS;
assert!(
validate(&ServerSettings {
result_output_retention_days: Some(0),
..Default::default()
})
.is_err(),
"0 would wedge the SPA's min=1 field; omit / null is how you ask for the default"
);
assert!(
validate(&ServerSettings {
result_output_retention_days: Some(MAX_RESULT_OUTPUT_RETENTION_DAYS + 1),
..Default::default()
})
.is_err(),
"past STREAM_RESULTS' window there is nothing left to replay"
);
assert!(
validate(&ServerSettings {
result_output_retention_days: Some(MAX_RESULT_OUTPUT_RETENTION_DAYS),
..Default::default()
})
.is_ok(),
"the ceiling itself must be accepted, not rejected off-by-one"
);
}
#[test]
fn validate_rejects_zero_or_oversize_collect_retention() {
use kanade_shared::wire::MAX_COLLECT_RETENTION_DAYS;
assert!(
validate(&ServerSettings {
collect_retention_days: Some(0),
..Default::default()
})
.is_err()
);
assert!(
validate(&ServerSettings {
collect_retention_days: Some(MAX_COLLECT_RETENTION_DAYS + 1),
..Default::default()
})
.is_err()
);
assert!(
validate(&ServerSettings {
collect_retention_days: Some(90),
..Default::default()
})
.is_ok()
);
}
#[test]
fn validate_rejects_blank_controller_group() {
let s = ServerSettings {
controller_group: Some(" ".into()),
..Default::default()
};
assert!(validate(&s).is_err());
}
#[test]
fn validate_rejects_zero_or_oversize_object_store_caps() {
for caps in [
ObjectStoreCaps {
result_output_mib: Some(0),
..Default::default()
},
ObjectStoreCaps {
app_packages_mib: Some(MAX_OBJECT_STORE_CAP_MIB + 1),
..Default::default()
},
ObjectStoreCaps {
scripts_mib: Some(0),
..Default::default()
},
] {
assert!(
validate(&ServerSettings {
object_store_caps: Some(caps),
..Default::default()
})
.is_err()
);
}
assert!(
validate(&ServerSettings {
object_store_caps: Some(ObjectStoreCaps {
app_packages_mib: Some(8192),
..Default::default()
}),
..Default::default()
})
.is_ok()
);
}
#[test]
fn validate_rejects_aggregate_over_broker_budget() {
assert!(
validate(&ServerSettings {
object_store_caps: Some(ObjectStoreCaps {
result_output_mib: Some(12_000),
agent_releases_mib: Some(12_000),
app_packages_mib: Some(12_000),
collections_mib: Some(12_000),
scripts_mib: Some(12_000),
}),
..Default::default()
})
.is_err()
);
assert!(
validate(&ServerSettings {
object_store_caps: Some(ObjectStoreCaps {
app_packages_mib: Some(40_000),
agent_releases_mib: Some(12_000),
..Default::default()
}),
..Default::default()
})
.is_err()
);
}
#[test]
fn validate_rejects_bad_mail() {
let mut m = sample_mail();
m.host = "".into();
assert!(
validate(&ServerSettings {
mail: Some(m),
..Default::default()
})
.is_err()
);
let mut m = sample_mail();
m.port = 0;
assert!(
validate(&ServerSettings {
mail: Some(m),
..Default::default()
})
.is_err()
);
let mut m = sample_mail();
m.from = "not an address".into();
assert!(
validate(&ServerSettings {
mail: Some(m),
..Default::default()
})
.is_err()
);
}
#[test]
fn validate_accepts_good_mail() {
assert!(
validate(&ServerSettings {
mail: Some(sample_mail()),
..Default::default()
})
.is_ok()
);
}
#[test]
fn normalize_trims_and_blanks_username() {
let mut m = sample_mail();
m.host = " smtp.example.com ".into();
m.from = " kanade-noreply@example.com ".into();
m.username = Some(" ".into());
let s = normalize(ServerSettings {
controller_group: Some(" infra ".into()),
mail: Some(m),
..Default::default()
});
assert_eq!(s.controller_group.as_deref(), Some("infra"));
let m = s.mail.unwrap();
assert_eq!(m.host, "smtp.example.com");
assert_eq!(m.from, "kanade-noreply@example.com");
assert_eq!(m.username, None);
}
#[test]
fn agent_install_merge_preserves_the_stored_token_across_a_url_edit() {
let mut stored =
obj(json!({ "agent_install": {"nats_url":"nats://old:4222","nats_token":"s3cret"} }));
let incoming = obj(json!({ "agent_install": {"nats_url":"nats://new:4222"} }));
assert!(merge_agent_install(&mut stored, &incoming));
assert_eq!(
stored["agent_install"],
json!({"nats_url":"nats://new:4222","nats_token":"s3cret"}),
);
assert!(!merge_agent_install(&mut stored, &incoming));
}
#[test]
fn agent_install_merge_sets_and_clears_the_token_only_explicitly() {
let mut stored =
obj(json!({ "agent_install": {"nats_url":"nats://b:4222","nats_token":"old"} }));
let incoming = obj(json!({ "agent_install": {"nats_token":"new"} }));
assert!(merge_agent_install(&mut stored, &incoming));
assert_eq!(stored["agent_install"]["nats_token"], "new");
assert_eq!(stored["agent_install"]["nats_url"], "nats://b:4222");
let incoming = obj(json!({ "agent_install": {"nats_token": null} }));
assert!(merge_agent_install(&mut stored, &incoming));
assert!(stored["agent_install"].get("nats_token").is_none());
assert_eq!(stored["agent_install"]["nats_url"], "nats://b:4222");
assert!(!merge_agent_install(&mut stored, &obj(json!({}))));
assert_eq!(stored["agent_install"]["nats_url"], "nats://b:4222");
let incoming = obj(json!({ "agent_install": null }));
assert!(merge_agent_install(&mut stored, &incoming));
assert!(stored.get("agent_install").is_none());
}
#[test]
fn agent_install_merge_never_stores_the_indicator() {
let mut stored =
obj(json!({ "agent_install": {"nats_url":"nats://b:4222","nats_token_set":true} }));
let incoming =
obj(json!({ "agent_install": {"nats_url":"nats://b:4222","nats_token_set":true} }));
assert!(merge_agent_install(&mut stored, &incoming));
assert!(stored["agent_install"].get("nats_token_set").is_none());
}
#[test]
fn agent_install_merge_preserves_the_user_pair_when_omitted() {
let mut stored = obj(json!({ "agent_install": {
"nats_url":"nats://old:4222","nats_user":"u","nats_password":"p"} }));
let incoming = obj(json!({ "agent_install": {
"nats_url":"nats://new:4222","nats_user_set":true,"nats_password_set":true} }));
assert!(merge_agent_install(&mut stored, &incoming));
assert_eq!(
stored["agent_install"],
json!({"nats_url":"nats://new:4222","nats_user":"u","nats_password":"p"}),
);
let incoming = obj(json!({ "agent_install": {"nats_user":"u2","nats_password":"p2"} }));
assert!(merge_agent_install(&mut stored, &incoming));
assert_eq!(stored["agent_install"]["nats_user"], "u2");
assert_eq!(stored["agent_install"]["nats_password"], "p2");
let incoming = obj(json!({ "agent_install": {"nats_user":null,"nats_password":null} }));
assert!(merge_agent_install(&mut stored, &incoming));
assert!(stored["agent_install"].get("nats_user").is_none());
assert!(stored["agent_install"].get("nats_password").is_none());
}
#[test]
fn user_pair_update_must_be_all_or_nothing() {
let check = |v: Value| validate_user_pair_update(&obj(json!({ "agent_install": v })));
assert!(check(json!({"nats_url":"nats://b:4222"})).is_ok());
assert!(check(json!({"nats_user":"u","nats_password":"p"})).is_ok());
assert!(check(json!({"nats_user":null,"nats_password":null})).is_ok());
assert!(validate_user_pair_update(&obj(json!({}))).is_ok());
for bad in [
json!({"nats_user":"u-secret"}),
json!({"nats_password":"pw-secret"}),
json!({"nats_user":"u-secret","nats_password":null}),
json!({"nats_user":null,"nats_password":"pw-secret"}),
] {
let (code, msg) = check(bad).unwrap_err();
assert_eq!(code, StatusCode::BAD_REQUEST);
assert!(msg.contains("pair"), "{msg}");
assert!(!msg.contains("secret"), "message leaked a value: {msg}");
}
}
#[test]
fn validate_rejects_bad_user_pair_values_without_echoing_them() {
use kanade_shared::wire::AgentInstallSection;
let with = |u: &str, p: &str| ServerSettings {
agent_install: Some(AgentInstallSection {
nats_user: Some(u.into()),
nats_password: Some(p.into()),
..Default::default()
}),
..Default::default()
};
assert!(validate(&with("agent", "it's a $pw \"x\"")).is_ok());
for (u, p) in [("", "p"), ("u", ""), ("u\nx", "p"), ("u", "pw-secret\r")] {
let (_, msg) = validate(&with(u, p)).unwrap_err();
assert!(!msg.contains("secret"), "{msg}");
}
}
#[test]
fn redact_secrets_hides_the_user_pair_in_the_audit_copy() {
let redacted = redact_secrets(json!({"agent_install": {
"nats_url":"nats://b:4222","nats_user":"u-secret","nats_password":"pw-secret"}}));
let text = redacted.to_string();
assert!(!text.contains("secret"), "{text}");
assert_eq!(redacted["agent_install"]["nats_user_set"], true);
assert_eq!(redacted["agent_install"]["nats_password_set"], true);
assert_eq!(redacted["agent_install"]["nats_token_set"], false);
}
#[test]
fn agent_install_merge_drops_a_fully_cleared_section() {
let mut stored = obj(json!({ "agent_install": {"nats_url":"nats://b:4222"} }));
let incoming = obj(json!({ "agent_install": {"nats_url": null} }));
assert!(merge_agent_install(&mut stored, &incoming));
assert!(stored.get("agent_install").is_none());
}
#[test]
fn validate_rejects_injectable_agent_install_values() {
use kanade_shared::wire::AgentInstallSection;
let ok = |nats_url: Option<&str>, nats_token: Option<&str>| ServerSettings {
agent_install: Some(AgentInstallSection {
nats_url: nats_url.map(str::to_string),
nats_token: nats_token.map(str::to_string),
..Default::default()
}),
..Default::default()
};
assert!(validate(&ok(Some("nats://broker.corp:4222"), Some("tok"))).is_ok());
assert!(validate(&ok(None, None)).is_ok());
for bad in ["", "nats://evil'\nx='y'", "nats://a\nb", "nats://a\rb"] {
assert!(
validate(&ok(Some(bad), None)).is_err(),
"nats_url {bad:?} must be rejected"
);
}
for bad in ["a\nb", "a\rb"] {
assert!(
validate(&ok(None, Some(bad))).is_err(),
"nats_token {bad:?} must be rejected"
);
}
}
}