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,
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)?;
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 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());
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 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) {
let had_token = ai.remove("nats_token").is_some_and(|v| v.is_string());
ai.insert("nats_token_set".to_string(), Value::Bool(had_token));
}
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(),
}
}
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 k == "nats_token_set" {
continue;
}
if v.is_null() {
merged.remove(k);
} else {
merged.insert(k.clone(), v.clone());
}
}
merged.remove("nats_token_set");
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(),
));
}
Ok(())
}
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}"),
)
})?;
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}"),
)
})?;
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 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, normalize,
redact_secrets, validate, validate_support_code,
};
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,
}
}
#[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_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),
nats_token_set: false,
require_signed_commands: None,
}),
..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"
);
}
}
}