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::{
MAX_AGENT_PRUNE_DAYS, MAX_CHECK_STATUS_STALE_DAYS, MAX_COLLECT_RETENTION_DAYS,
MAX_SESSION_TTL_HOURS, MAX_SUPPORT_UNLOCK_TTL_MINUTES, 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 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 kv = open_bucket(&s).await?;
let mut collect_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;
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
})
.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,
session_ttl_hours = ?merged.session_ttl_hours,
controller_group = ?merged.controller_group,
mail_configured = merged.mail.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",
),
}
}
audit::record(
&s.nats,
"operator",
"server_settings_set",
Some(KEY_SERVER_SETTINGS),
Some(&caller),
redact_support_hashes(doc),
)
.await;
Ok(Json(merged.redacted()))
}
fn redact_support_hashes(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");
}
}
}
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 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(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)?;
}
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::{
MIN_SUPPORT_CODE_LEN, ServerSettings, SupportCodeBody, hash_support_code, merge_field,
normalize, redact_support_hashes, 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_support_hashes(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_redaction_tolerates_a_missing_or_odd_field() {
assert_eq!(
redact_support_hashes(json!({"agent_prune_days": 7}))["agent_prune_days"],
7
);
assert_eq!(
redact_support_hashes(json!({"support_codes": "nonsense"}))["support_codes"],
"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_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_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);
}
}