use std::sync::Mutex;
use std::time::{Duration, Instant};
use kanade_shared::signing::{KeyPolicy, KeyRing, SigHeaders, VerifyError, verify};
use kanade_shared::wire::ObsEvent;
use serde::Deserialize;
use tracing::{info, warn};
const REG_SUBKEY: &str = r"SOFTWARE\kanade\agent";
const REG_VALUE: &str = "CommandKeys";
const SOURCE: &str = "command_signature";
const RELOAD_MIN_INTERVAL: Duration = Duration::from_secs(30);
#[derive(Debug, Deserialize)]
struct KeyEntry {
kid: String,
public_key: String,
#[serde(default)]
label: Option<String>,
#[serde(default)]
max_age_secs: Option<u64>,
#[serde(default)]
audit_every_use: bool,
}
fn read_keyring() -> Result<KeyRing, String> {
let Some(raw) = kanade_shared::secrets::read_hklm_value(REG_SUBKEY, REG_VALUE) else {
return Ok(KeyRing::new());
};
parse_keyring(&raw)
}
fn parse_keyring(raw: &str) -> Result<KeyRing, String> {
use base64::Engine;
let entries: Vec<KeyEntry> = serde_json::from_str(raw).map_err(|e| e.to_string())?;
let mut ring = KeyRing::new();
let mut seen: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
for e in entries {
if !seen.insert(e.kid.clone()) {
return Err(format!(
"key {} appears twice — two different keys must never share an id, and a ring \
keyed by id cannot hold both",
e.kid
));
}
let bytes = base64::engine::general_purpose::STANDARD
.decode(&e.public_key)
.map_err(|err| format!("key {}: {err}", e.kid))?;
let arr: [u8; 32] = bytes
.as_slice()
.try_into()
.map_err(|_| format!("key {}: expected 32 bytes, got {}", e.kid, bytes.len()))?;
let vk = ed25519_dalek::VerifyingKey::from_bytes(&arr)
.map_err(|err| format!("key {}: {err}", e.kid))?;
let label = e.label.unwrap_or_else(|| e.kid.clone());
let policy = match e.max_age_secs {
Some(secs) => KeyPolicy::break_glass(label, std::time::Duration::from_secs(secs)),
None => {
let mut p = KeyPolicy::backend(label);
p.audit_every_use = e.audit_every_use;
p
}
};
ring.insert(e.kid, vk, policy);
}
Ok(ring)
}
pub fn headers_of(msg: &async_nats::Message) -> SigHeaders {
let get = |name: &str| {
msg.headers
.as_ref()
.and_then(|h| h.get(name))
.map(|v| v.to_string())
};
SigHeaders {
sig_b64: get(kanade_shared::signing::SIG),
kid: get(kanade_shared::signing::SIG_KID),
alg: get(kanade_shared::signing::SIG_ALG),
at_ms: get(kanade_shared::signing::SIG_AT),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Outcome {
Verified,
Unsigned,
Unprovisioned,
UnknownKid,
Invalid,
Stale,
}
impl Outcome {
#[cfg(test)]
const ALL: [Outcome; 6] = [
Outcome::Verified,
Outcome::Unsigned,
Outcome::Unprovisioned,
Outcome::UnknownKid,
Outcome::Invalid,
Outcome::Stale,
];
fn kind(self) -> &'static str {
match self {
Outcome::Verified => "command_signature_ok",
Outcome::Unsigned => "command_signature_absent",
Outcome::Unprovisioned => "command_signature_unprovisioned",
Outcome::UnknownKid => "command_signature_unknown_key",
Outcome::Invalid => "command_signature_invalid",
Outcome::Stale => "command_signature_stale",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Reload {
Done,
RateLimited,
Failed,
}
impl Reload {
fn as_str(self) -> &'static str {
match self {
Reload::Done => "reloaded",
Reload::RateLimited => "rate-limited",
Reload::Failed => "reload-failed",
}
}
}
fn lock<T>(m: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
m.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
}
type Loader = Box<dyn Fn() -> Result<KeyRing, String> + Send + Sync>;
pub struct Verifier {
ring: Mutex<KeyRing>,
loader: Loader,
last_reload: Mutex<Instant>,
pc_id: String,
obs_dir: std::path::PathBuf,
last: Mutex<Option<Outcome>>,
}
impl Verifier {
pub fn new(pc_id: String, obs_dir: std::path::PathBuf) -> Self {
Self::with_loader(pc_id, obs_dir, Box::new(read_keyring))
}
fn with_loader(pc_id: String, obs_dir: std::path::PathBuf, loader: Loader) -> Self {
if let Err(e) = crate::obs_outbox::ensure_outbox_dir(&obs_dir) {
warn!(error = %e, "command_verify: outbox dir — reports may be dropped until it exists");
}
let ring = loader().unwrap_or_else(|e| {
warn!(error = %e, "command keyring is unreadable — treating as empty");
KeyRing::new()
});
if ring.is_empty() {
info!(
"command keyring is empty — signed commands will be reported unprovisioned until \
one is distributed (no restart needed; the ring reloads on demand)"
);
} else {
info!(kids = ?ring.kids().collect::<Vec<_>>(), "command keyring loaded");
}
Self {
ring: Mutex::new(ring),
loader,
last_reload: Mutex::new(Instant::now()),
pc_id,
obs_dir,
last: Mutex::new(None),
}
}
pub fn trusted_kids(&self) -> Vec<String> {
lock(&self.ring).kids().map(str::to_owned).collect()
}
pub fn observe(&self, body: &[u8], headers: &SigHeaders, request_id: &str) -> Outcome {
self.observe_at(
body,
headers,
request_id,
chrono::Utc::now().timestamp_millis(),
)
}
fn observe_at(
&self,
body: &[u8],
headers: &SigHeaders,
request_id: &str,
now_ms: i64,
) -> Outcome {
let outcome = self.classify(body, headers, request_id, now_ms, Instant::now());
self.report_transition(outcome);
outcome
}
fn classify(
&self,
body: &[u8],
headers: &SigHeaders,
request_id: &str,
now_ms: i64,
now: Instant,
) -> Outcome {
match self.check(body, headers, request_id, now_ms) {
Ok(outcome) => outcome,
Err(kid) => {
let reload = self.reload_if_due(now);
if reload != Reload::Done {
return self.report_missing(&kid, request_id, reload.as_str());
}
match self.check(body, headers, request_id, now_ms) {
Ok(outcome) => {
info!(
kid,
request_id, "keyring reload resolved a previously unknown key"
);
outcome
}
Err(kid) => self.report_missing(&kid, request_id, "reloaded"),
}
}
}
}
fn check(
&self,
body: &[u8],
headers: &SigHeaders,
request_id: &str,
now_ms: i64,
) -> Result<Outcome, String> {
let ring = lock(&self.ring);
match verify(&ring, body, headers, now_ms) {
Ok(v) => {
if v.policy.audit_every_use {
warn!(kid = v.kid, request_id, "command signed by an audited key");
}
Ok(Outcome::Verified)
}
Err(VerifyError::Unsigned) => Ok(Outcome::Unsigned),
Err(VerifyError::UnknownKid { kid }) => Err(kid),
Err(e @ VerifyError::Stale { .. }) => {
warn!(error = %e, request_id, "command signature is past its freshness bound");
Ok(Outcome::Stale)
}
Err(e) => {
warn!(error = %e, request_id, "command signature did not verify");
Ok(Outcome::Invalid)
}
}
}
fn report_missing(&self, kid: &str, request_id: &str, reload: &str) -> Outcome {
let ring = lock(&self.ring);
if ring.is_empty() {
warn!(
kid,
request_id,
reload,
"command is signed but this agent holds no keys — provision \
HKLM\\SOFTWARE\\kanade\\agent\\CommandKeys"
);
Outcome::Unprovisioned
} else {
warn!(
kid,
request_id,
reload,
known = ?ring.kids().collect::<Vec<_>>(),
"command signed by a key this agent does not have"
);
Outcome::UnknownKid
}
}
fn reload_if_due(&self, now: Instant) -> Reload {
let mut last = lock(&self.last_reload);
if now.checked_duration_since(*last).unwrap_or_default() < RELOAD_MIN_INTERVAL {
return Reload::RateLimited;
}
*last = now;
match (self.loader)() {
Ok(fresh) => {
*lock(&self.ring) = fresh;
Reload::Done
}
Err(e) => {
warn!(
error = %e,
"keyring reload failed — keeping the keys already loaded"
);
Reload::Failed
}
}
}
fn report_transition(&self, outcome: Outcome) {
let mut last = lock(&self.last);
let Some(transition) = step(*last, outcome) else {
return;
};
let event = build_event(&self.pc_id, &transition, chrono::Utc::now());
match crate::obs_outbox::enqueue(&self.obs_dir, &event) {
Ok(_path) => *last = Some(outcome),
Err(e) => warn!(
error = %e,
kind = outcome.kind(),
"command_verify: enqueue failed — will retry on the next command"
),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct Transition {
from: Option<Outcome>,
to: Outcome,
}
fn step(last: Option<Outcome>, outcome: Outcome) -> Option<Transition> {
if last == Some(outcome) {
return None;
}
Some(Transition {
from: last,
to: outcome,
})
}
fn build_event(pc_id: &str, t: &Transition, at: chrono::DateTime<chrono::Utc>) -> ObsEvent {
ObsEvent {
pc_id: pc_id.to_string(),
at,
kind: t.to.kind().to_string(),
source: SOURCE.to_string(),
event_record_id: Some(format!("{}:{}", t.to.kind(), at.timestamp_millis())),
payload: serde_json::json!({
"from": t.from.map(|p| p.kind()),
"to": t.to.kind(),
}),
}
}
#[cfg(test)]
mod tests {
use super::*;
use base64::Engine;
use ed25519_dalek::SigningKey;
use kanade_shared::signing::sign;
fn b64(bytes: &[u8]) -> String {
base64::engine::general_purpose::STANDARD.encode(bytes)
}
fn test_dir() -> std::path::PathBuf {
std::env::temp_dir().join("kanade-command-verify-test")
}
fn verifier_with(ring: KeyRing) -> Verifier {
Verifier::with_loader("PC1".into(), test_dir(), Box::new(move || Ok(ring.clone())))
}
fn backend_ring(kid: &str, sk: &SigningKey) -> KeyRing {
let mut r = KeyRing::new();
r.insert(kid, sk.verifying_key(), KeyPolicy::backend("backend"));
r
}
#[derive(Clone)]
struct Store {
inner: std::sync::Arc<Mutex<(Result<KeyRing, String>, usize)>>,
}
impl Default for Store {
fn default() -> Self {
Self {
inner: std::sync::Arc::new(Mutex::new((Ok(KeyRing::new()), 0))),
}
}
}
impl Store {
fn provision(&self, ring: KeyRing) {
lock(&self.inner).0 = Ok(ring);
}
fn corrupt(&self) {
lock(&self.inner).0 = Err("expected value at line 1 column 3".into());
}
fn reads(&self) -> usize {
lock(&self.inner).1
}
fn loader(&self) -> Loader {
let inner = self.inner.clone();
Box::new(move || {
let mut g = lock(&inner);
g.1 += 1;
g.0.clone()
})
}
}
#[test]
fn keyring_parses_the_provisioned_shape() {
let sk = SigningKey::from_bytes(&[7u8; 32]);
let raw = format!(
r#"[{{"kid":"backend-1","public_key":"{}","label":"backend"}}]"#,
b64(sk.verifying_key().as_bytes())
);
let ring = parse_keyring(&raw).expect("parses");
let (kid, _, policy) = ring.get("backend-1").expect("present");
assert_eq!(kid, "backend-1");
assert_eq!(policy.label, "backend");
assert_eq!(policy.max_age, None);
assert!(!policy.audit_every_use);
}
#[test]
fn a_max_age_entry_becomes_a_break_glass_policy() {
let sk = SigningKey::from_bytes(&[8u8; 32]);
let raw = format!(
r#"[{{"kid":"bg","public_key":"{}","max_age_secs":300}}]"#,
b64(sk.verifying_key().as_bytes())
);
let ring = parse_keyring(&raw).unwrap();
let (_, _, policy) = ring.get("bg").unwrap();
assert_eq!(policy.max_age, Some(std::time::Duration::from_secs(300)));
assert!(policy.audit_every_use);
assert_eq!(policy.label, "bg");
}
#[test]
fn a_malformed_keyring_is_rejected_rather_than_half_loaded() {
let good = b64(SigningKey::from_bytes(&[1u8; 32])
.verifying_key()
.as_bytes());
let raw = format!(
r#"[{{"kid":"ok","public_key":"{good}"}},{{"kid":"bad","public_key":"not base64"}}]"#
);
assert!(parse_keyring(&raw).is_err());
let short = b64(&[0u8; 31]);
assert!(parse_keyring(&format!(r#"[{{"kid":"s","public_key":"{short}"}}]"#)).is_err());
assert!(parse_keyring("{{{").is_err());
}
#[test]
fn a_duplicate_kid_is_refused_rather_than_silently_collapsed() {
let a = b64(SigningKey::from_bytes(&[1u8; 32])
.verifying_key()
.as_bytes());
let b = b64(SigningKey::from_bytes(&[2u8; 32])
.verifying_key()
.as_bytes());
let raw =
format!(r#"[{{"kid":"bg","public_key":"{a}"}},{{"kid":"bg","public_key":"{b}"}}]"#);
let err = parse_keyring(&raw).unwrap_err();
assert!(err.contains("bg"), "the error must name the id: {err}");
assert!(err.contains("twice"), "{err}");
let raw = format!(
r#"[{{"kid":"backend-1","public_key":"{a}"}},{{"kid":"bg","public_key":"{b}","max_age_secs":900}}]"#
);
let ring = parse_keyring(&raw).expect("two distinct ids are fine");
assert!(ring.get("backend-1").is_some());
assert!(ring.get("bg").is_some());
}
#[test]
fn an_empty_ring_reports_unsigned_and_unknown_but_never_verifies() {
let ring = parse_keyring("[]").unwrap();
assert!(ring.is_empty());
let sk = SigningKey::from_bytes(&[3u8; 32]);
assert_eq!(
verify(&ring, b"body", &SigHeaders::default(), 0),
Err(VerifyError::Unsigned)
);
let headers = sign(&sk, "backend-1", b"body", 0);
assert!(matches!(
verify(&ring, b"body", &headers, 0),
Err(VerifyError::UnknownKid { .. })
));
}
#[test]
fn outcome_kinds_are_distinct_and_stable() {
let kinds: Vec<_> = Outcome::ALL.iter().map(|o| o.kind()).collect();
let unique: std::collections::BTreeSet<_> = kinds.iter().collect();
assert_eq!(unique.len(), kinds.len(), "kinds must not collide");
assert!(kinds.iter().all(|k| k.starts_with("command_signature")));
}
#[test]
fn step_reports_the_baseline_then_only_on_change() {
assert_eq!(
step(None, Outcome::Unsigned),
Some(Transition {
from: None,
to: Outcome::Unsigned
})
);
assert_eq!(step(Some(Outcome::Unsigned), Outcome::Unsigned), None);
assert_eq!(
step(Some(Outcome::Unsigned), Outcome::Verified),
Some(Transition {
from: Some(Outcome::Unsigned),
to: Outcome::Verified
})
);
assert_eq!(
step(Some(Outcome::Verified), Outcome::UnknownKid),
Some(Transition {
from: Some(Outcome::Verified),
to: Outcome::UnknownKid
})
);
}
#[test]
fn a_flapping_state_reports_each_way() {
let mut last = None;
let mut reported = Vec::new();
for o in [
Outcome::Verified,
Outcome::Verified,
Outcome::UnknownKid,
Outcome::Verified,
] {
if let Some(t) = step(last, o) {
reported.push(t.to);
last = Some(o);
}
}
assert_eq!(
reported,
vec![Outcome::Verified, Outcome::UnknownKid, Outcome::Verified]
);
}
#[test]
fn a_failed_enqueue_leaves_the_transition_unreported_so_it_retries() {
let mut last: Option<Outcome> = None;
let enqueue_ok = false;
if let Some(_t) = step(last, Outcome::UnknownKid)
&& enqueue_ok
{
last = Some(Outcome::UnknownKid);
}
assert_eq!(last, None, "a failed enqueue must not mark it reported");
assert!(
step(last, Outcome::UnknownKid).is_some(),
"the retry has to be possible"
);
}
#[test]
fn the_event_carries_both_ends_of_the_transition() {
let at = chrono::DateTime::from_timestamp(1_700_000_000, 0).unwrap();
let e = build_event(
"PC1",
&Transition {
from: Some(Outcome::Verified),
to: Outcome::UnknownKid,
},
at,
);
assert_eq!(e.pc_id, "PC1");
assert_eq!(e.source, SOURCE);
assert_eq!(e.kind, "command_signature_unknown_key");
assert_eq!(e.payload["from"], "command_signature_ok");
assert_eq!(e.payload["to"], "command_signature_unknown_key");
assert_eq!(
e.event_record_id.as_deref(),
Some("command_signature_unknown_key:1700000000000")
);
}
#[test]
fn the_first_event_reports_no_previous_state() {
let at = chrono::DateTime::from_timestamp(0, 0).unwrap();
let e = build_event(
"PC1",
&Transition {
from: None,
to: Outcome::Unsigned,
},
at,
);
assert!(e.payload["from"].is_null());
}
#[test]
fn a_replayed_break_glass_command_is_reported_stale_not_invalid() {
let sk = SigningKey::from_bytes(&[9u8; 32]);
let mut ring = KeyRing::new();
ring.insert(
"break-glass",
sk.verifying_key(),
KeyPolicy::break_glass("break-glass", std::time::Duration::from_secs(300)),
);
let v = verifier_with(ring);
let now = 1_700_000_000_000i64;
let body = b"emergency";
let week = 7 * 24 * 60 * 60 * 1000;
assert_eq!(
v.observe_at(body, &sign(&sk, "break-glass", body, now - week), "r1", now),
Outcome::Stale
);
assert_eq!(
v.observe_at(
body,
&sign(&sk, "break-glass", body, now - 1_000),
"r2",
now
),
Outcome::Verified
);
}
#[test]
fn a_key_provisioned_after_boot_takes_effect_without_a_restart() {
let sk = SigningKey::from_bytes(&[11u8; 32]);
let store = Store::default();
let v = Verifier::with_loader("PC1".into(), test_dir(), store.loader());
let now = 1_700_000_000_000i64;
let body = b"job";
let headers = sign(&sk, "backend-1", body, now);
let boot = Instant::now();
assert_eq!(
v.classify(body, &headers, "r1", now, boot),
Outcome::Unprovisioned
);
store.provision(backend_ring("backend-1", &sk));
assert_eq!(
v.classify(body, &headers, "r2", now, boot),
Outcome::Unprovisioned
);
let later = boot + RELOAD_MIN_INTERVAL;
assert_eq!(
v.classify(body, &headers, "r3", now, later),
Outcome::Verified
);
}
#[test]
fn a_reload_that_cannot_be_read_keeps_the_working_ring() {
let sk = SigningKey::from_bytes(&[21u8; 32]);
let store = Store::default();
store.provision(backend_ring("backend-1", &sk));
let v = Verifier::with_loader("PC1".into(), test_dir(), store.loader());
let now = 1_700_000_000_000i64;
let good = sign(&sk, "backend-1", b"job", now);
let t = Instant::now() + RELOAD_MIN_INTERVAL * 2;
assert_eq!(v.classify(b"job", &good, "r1", now, t), Outcome::Verified);
store.corrupt();
let unknown = sign(&sk, "backend-2", b"job", now);
assert_eq!(
v.classify(b"job", &unknown, "r2", now, t + RELOAD_MIN_INTERVAL),
Outcome::UnknownKid,
"a failed reload must not turn this into Unprovisioned"
);
assert_eq!(
v.classify(b"job", &good, "r3", now, t + RELOAD_MIN_INTERVAL * 2),
Outcome::Verified,
"the working ring must survive a failed reload"
);
}
#[test]
fn only_an_unknown_key_reaches_the_store() {
let sk = SigningKey::from_bytes(&[12u8; 32]);
let store = Store::default();
store.provision(backend_ring("backend-1", &sk));
let v = Verifier::with_loader("PC1".into(), test_dir(), store.loader());
let after_boot = Instant::now() + RELOAD_MIN_INTERVAL * 10;
let now = 1_700_000_000_000i64;
let reads = store.reads();
assert_eq!(
v.classify(b"x", &SigHeaders::default(), "r1", now, after_boot),
Outcome::Unsigned
);
let forged = sign(&SigningKey::from_bytes(&[99u8; 32]), "backend-1", b"x", now);
assert_eq!(
v.classify(b"x", &forged, "r2", now, after_boot),
Outcome::Invalid
);
let partial = SigHeaders {
sig_b64: Some("AAAA".into()),
kid: None,
alg: None,
at_ms: None,
};
assert_eq!(
v.classify(b"x", &partial, "r3", now, after_boot),
Outcome::Invalid
);
assert_eq!(store.reads(), reads, "none of these may touch the store");
let unknown = sign(&sk, "backend-2", b"x", now);
assert_eq!(
v.classify(b"x", &unknown, "r4", now, after_boot),
Outcome::UnknownKid
);
assert_eq!(store.reads(), reads + 1);
}
#[test]
fn repeated_unknown_keys_reload_at_most_once_per_interval() {
let sk = SigningKey::from_bytes(&[13u8; 32]);
let store = Store::default();
store.provision(backend_ring("backend-1", &sk));
let v = Verifier::with_loader("PC1".into(), test_dir(), store.loader());
let now = 1_700_000_000_000i64;
let headers = sign(&sk, "backend-2", b"x", now);
let base = Instant::now() + RELOAD_MIN_INTERVAL;
let reads = store.reads();
for i in 0..20 {
let t = base + RELOAD_MIN_INTERVAL / 40 * i;
assert_eq!(v.classify(b"x", &headers, "r", now, t), Outcome::UnknownKid);
}
assert_eq!(
store.reads(),
reads + 1,
"a flood of invented key ids must not become a flood of store reads"
);
assert_eq!(
v.classify(b"x", &headers, "r", now, base + RELOAD_MIN_INTERVAL * 2),
Outcome::UnknownKid
);
assert_eq!(store.reads(), reads + 2);
}
#[test]
fn an_empty_ring_and_a_missing_key_are_different_states() {
let sk = SigningKey::from_bytes(&[14u8; 32]);
let now = 1_700_000_000_000i64;
let headers = sign(&sk, "backend-2", b"x", now);
let t = Instant::now() + RELOAD_MIN_INTERVAL * 2;
let empty =
Verifier::with_loader("PC1".into(), test_dir(), Box::new(|| Ok(KeyRing::new())));
assert_eq!(
empty.classify(b"x", &headers, "r1", now, t),
Outcome::Unprovisioned
);
let other = verifier_with(backend_ring("backend-1", &sk));
assert_eq!(
other.classify(b"x", &headers, "r2", now, t),
Outcome::UnknownKid
);
}
#[test]
fn a_non_monotonic_clock_does_not_panic() {
let sk = SigningKey::from_bytes(&[15u8; 32]);
let v = verifier_with(backend_ring("backend-1", &sk));
let now = 1_700_000_000_000i64;
let headers = sign(&sk, "backend-2", b"x", now);
let Some(past) = Instant::now().checked_sub(RELOAD_MIN_INTERVAL * 3) else {
return;
};
assert_eq!(
v.classify(b"x", &headers, "r1", now, past),
Outcome::UnknownKid
);
}
#[test]
fn the_ordinary_signer_is_never_stale() {
let sk = SigningKey::from_bytes(&[4u8; 32]);
let mut ring = KeyRing::new();
ring.insert(
"backend-1",
sk.verifying_key(),
KeyPolicy::backend("backend"),
);
let v = verifier_with(ring);
let now = 1_700_000_000_000i64;
let week = 7 * 24 * 60 * 60 * 1000;
assert_eq!(
v.observe_at(
b"job",
&sign(&sk, "backend-1", b"job", now - week),
"r1",
now
),
Outcome::Verified
);
}
}