use std::time::Duration;
use anyhow::Context as _;
use async_nats::jetstream::Context;
use async_nats::jetstream::kv::{Operation, Store};
use kanade_shared::kv::{
BUCKET_FLEET_CONFIG, BUCKET_SERVER_SETTINGS, KEY_SERVER_SETTINGS, KEY_SUPPORT_CODES,
};
use kanade_shared::wire::{ServerSettings, SupportCodesProjection};
use tokio::sync::Mutex;
use tracing::{info, warn};
const RESYNC_INTERVAL: Duration = Duration::from_secs(300);
const MAX_ATTEMPTS: usize = 8;
static SYNC_LOCK: Mutex<()> = Mutex::const_new(());
#[derive(Debug, PartialEq, Eq)]
enum Plan {
Unchanged,
Write(Vec<u8>),
}
fn plan(settings: Option<&[u8]>, current: Option<&[u8]>) -> anyhow::Result<Plan> {
let settings: ServerSettings = match settings {
Some(bytes) => serde_json::from_slice(bytes).context("decode server_settings")?,
None => ServerSettings::default(),
};
let desired = SupportCodesProjection::from_settings(&settings);
if let Some(cur) = current
&& serde_json::from_slice::<SupportCodesProjection>(cur).is_ok_and(|c| c == desired)
{
return Ok(Plan::Unchanged);
}
Ok(Plan::Write(
serde_json::to_vec(&desired).context("encode support codes projection")?,
))
}
async fn live_bytes(kv: &Store, key: &str) -> anyhow::Result<(Option<Vec<u8>>, Option<u64>)> {
match kv
.entry(key)
.await
.with_context(|| format!("kv entry '{key}'"))?
{
Some(e) if e.operation == Operation::Put => Ok((Some(e.value.to_vec()), Some(e.revision))),
Some(e) => Ok((None, Some(e.revision))),
None => Ok((None, None)),
}
}
pub async fn sync_support_codes_projection(js: &Context) -> anyhow::Result<bool> {
let _guard = SYNC_LOCK.lock().await;
let settings_kv = js
.get_key_value(BUCKET_SERVER_SETTINGS)
.await
.context("open server_settings bucket")?;
let fleet_kv = js
.get_key_value(BUCKET_FLEET_CONFIG)
.await
.context("open fleet_config bucket")?;
let mut wrote = false;
let mut last_err = None;
for _ in 0..MAX_ATTEMPTS {
let (current, revision) = live_bytes(&fleet_kv, KEY_SUPPORT_CODES).await?;
let (settings, _) = live_bytes(&settings_kv, KEY_SERVER_SETTINGS).await?;
let body = match plan(settings.as_deref(), current.as_deref())? {
Plan::Unchanged => return Ok(wrote),
Plan::Write(b) => b,
};
let res = match revision {
Some(rev) => fleet_kv
.update(KEY_SUPPORT_CODES, body.into(), rev)
.await
.map(|_| ())
.map_err(anyhow::Error::from),
None => fleet_kv
.create(KEY_SUPPORT_CODES, body.into())
.await
.map(|_| ())
.map_err(anyhow::Error::from),
};
match res {
Ok(()) => {
info!("support code projection published");
wrote = true;
}
Err(e) => last_err = Some(e),
}
}
Err(last_err
.unwrap_or_else(|| anyhow::anyhow!("projection kept changing"))
.context("publish support codes projection"))
}
pub async fn sync_or_warn(js: &Context) {
if let Err(e) = sync_support_codes_projection(js).await {
warn!(error = %format!("{e:#}"), "support codes projection sync failed; will retry");
}
}
pub fn spawn(js: Context) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut interval = tokio::time::interval(RESYNC_INTERVAL);
loop {
interval.tick().await;
sync_or_warn(&js).await;
}
})
}
#[cfg(test)]
mod tests {
use super::*;
use kanade_shared::config::{MailEncryption, MailSection};
use kanade_shared::wire::{AgentInstallSection, SupportCode};
use serde_json::{Value, json};
use std::collections::BTreeSet;
fn code(scope: &str, disabled: bool) -> SupportCode {
SupportCode {
scope: scope.into(),
hash: format!("$argon2id$hash-{scope}"),
label: Some("desk".into()),
ttl_minutes: Some(30),
disabled,
}
}
fn secret_settings(codes: Vec<SupportCode>) -> ServerSettings {
ServerSettings {
support_codes: codes,
agent_install: Some(AgentInstallSection {
nats_url: Some("nats://x:4222".into()),
nats_token: Some("super-secret-token".into()),
..Default::default()
}),
mail: Some(MailSection {
host: "smtp.example".into(),
port: 587,
encryption: MailEncryption::Starttls,
username: Some("mailer".into()),
from: "kanade@example".into(),
}),
agent_prune_days: Some(9),
..Default::default()
}
}
fn keys(v: &Value) -> BTreeSet<String> {
v.as_object().unwrap().keys().cloned().collect()
}
#[test]
fn projection_carries_only_support_code_fields() {
let settings = secret_settings(vec![code("support", false), code("admin", true)]);
let raw = serde_json::to_string(&SupportCodesProjection::from_settings(&settings)).unwrap();
assert!(!raw.contains("super-secret-token"));
assert!(!raw.contains("smtp.example"));
let v: Value = serde_json::from_str(&raw).unwrap();
assert_eq!(keys(&v), BTreeSet::from(["support_codes".to_string()]));
let allowed: BTreeSet<String> = ["scope", "hash", "label", "ttl_minutes", "disabled"]
.map(String::from)
.into();
for c in v["support_codes"].as_array().unwrap() {
assert!(
keys(c).is_subset(&allowed),
"unexpected keys: {:?}",
keys(c)
);
}
assert_eq!(v["support_codes"][1]["disabled"], json!(true));
}
#[test]
fn startup_reconcile_creates_the_projection_when_absent() {
let settings = serde_json::to_vec(&secret_settings(vec![code("support", false)])).unwrap();
let Plan::Write(body) = plan(Some(&settings), None).unwrap() else {
panic!("expected a write");
};
let p: SupportCodesProjection = serde_json::from_slice(&body).unwrap();
assert_eq!(p.support_codes, vec![code("support", false)]);
}
#[test]
fn unconfigured_deployment_gets_an_empty_projection_not_nothing() {
let Plan::Write(body) = plan(None, None).unwrap() else {
panic!("expected a write");
};
assert_eq!(body, br#"{"support_codes":[]}"#);
}
#[test]
fn a_change_updates_and_clearing_empties_it() {
let before = serde_json::to_vec(&secret_settings(vec![code("a", false)])).unwrap();
let Plan::Write(cur) = plan(Some(&before), None).unwrap() else {
panic!()
};
let after = serde_json::to_vec(&secret_settings(vec![code("a", true)])).unwrap();
let Plan::Write(updated) = plan(Some(&after), Some(&cur)).unwrap() else {
panic!("a changed code must rewrite the projection");
};
assert!(
serde_json::from_slice::<SupportCodesProjection>(&updated)
.unwrap()
.support_codes[0]
.disabled
);
let cleared = serde_json::to_vec(&secret_settings(vec![])).unwrap();
let Plan::Write(empty) = plan(Some(&cleared), Some(&updated)).unwrap() else {
panic!("clearing must rewrite, not skip");
};
assert_eq!(empty, br#"{"support_codes":[]}"#);
}
#[test]
fn unchanged_projection_is_not_rewritten() {
let settings = serde_json::to_vec(&secret_settings(vec![code("a", false)])).unwrap();
let Plan::Write(cur) = plan(Some(&settings), None).unwrap() else {
panic!()
};
assert_eq!(plan(Some(&settings), Some(&cur)).unwrap(), Plan::Unchanged);
}
#[test]
fn a_corrupt_settings_document_never_clobbers_the_projection() {
assert!(plan(Some(b"not json"), Some(br#"{"support_codes":[]}"#)).is_err());
}
}