use std::sync::Mutex;
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";
#[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,
}
pub fn load_keyring() -> KeyRing {
let Some(raw) = kanade_shared::secrets::read_hklm_value(REG_SUBKEY, REG_VALUE) else {
return KeyRing::new();
};
parse_keyring(&raw).unwrap_or_else(|e| {
warn!(error = %e, "command keyring is unreadable — treating as empty");
KeyRing::new()
})
}
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();
for e in entries {
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,
UnknownKid,
Invalid,
Stale,
}
impl Outcome {
fn kind(self) -> &'static str {
match self {
Outcome::Verified => "command_signature_ok",
Outcome::Unsigned => "command_signature_absent",
Outcome::UnknownKid => "command_signature_unknown_key",
Outcome::Invalid => "command_signature_invalid",
Outcome::Stale => "command_signature_stale",
}
}
}
pub struct Verifier {
ring: KeyRing,
pc_id: String,
obs_dir: std::path::PathBuf,
last: Mutex<Option<Outcome>>,
}
impl Verifier {
pub fn new(ring: KeyRing, pc_id: String, obs_dir: std::path::PathBuf) -> 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");
}
if ring.is_empty() {
info!("command keyring is empty — signatures will be reported, not checked");
} else {
info!(kids = ?ring.kids().collect::<Vec<_>>(), "command keyring loaded");
}
Self {
ring,
pc_id,
obs_dir,
last: Mutex::new(None),
}
}
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 = match verify(&self.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");
}
Outcome::Verified
}
Err(VerifyError::Unsigned) => Outcome::Unsigned,
Err(VerifyError::UnknownKid { kid }) => {
warn!(
kid,
request_id,
known = ?self.ring.kids().collect::<Vec<_>>(),
"command signed by a key this agent does not have"
);
Outcome::UnknownKid
}
Err(e @ VerifyError::Stale { .. }) => {
warn!(error = %e, request_id, "command signature is past its freshness bound");
Outcome::Stale
}
Err(e) => {
warn!(error = %e, request_id, "command signature did not verify");
Outcome::Invalid
}
};
self.report_transition(outcome);
outcome
}
fn report_transition(&self, outcome: Outcome) {
let mut last = match self.last.lock() {
Ok(g) => g,
Err(poisoned) => poisoned.into_inner(),
};
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)
}
#[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 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::Verified,
Outcome::Unsigned,
Outcome::UnknownKid,
Outcome::Invalid,
Outcome::Stale,
]
.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 dir = std::env::temp_dir().join("kanade-command-verify-test");
let v = Verifier::new(ring, "PC1".into(), dir);
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 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 dir = std::env::temp_dir().join("kanade-command-verify-test");
let v = Verifier::new(ring, "PC1".into(), dir);
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
);
}
}