use async_nats::HeaderMap;
use async_nats::client::PublishError;
use bytes::Bytes;
use kanade_shared::signing::{self, Signer};
use tracing::{error, info, warn};
const ENV_SIGNING_KEY: &str = "KANADE_COMMAND_SIGNING_KEY";
const ENV_SIGNING_KID: &str = "KANADE_COMMAND_SIGNING_KID";
pub struct CommandPublisher {
nats: async_nats::Client,
signer: Option<Signer>,
}
impl CommandPublisher {
pub fn new(nats: async_nats::Client, signer: Option<Signer>) -> Self {
Self { nats, signer }
}
pub fn from_host(nats: async_nats::Client) -> Self {
let signer = match resolve_signer() {
Ok(Some(s)) => {
info!(
kid = s.kid(),
"signing every published command — agents that lack this key will report \
command_signature_unknown_key"
);
Some(s)
}
Ok(None) => {
warn!(
"no command-signing key on this host — commands are published unsigned. Run \
`kanade-backend command-key-generate` and distribute the public key to the \
fleet first (#1165)."
);
None
}
Err(e) => {
error!(
error = %e,
"command-signing key is configured but unusable — commands are published \
UNSIGNED. This is not a warning about a missing key: something is set and \
wrong."
);
None
}
};
Self::new(nats, signer)
}
pub fn kid(&self) -> Option<&str> {
self.signer.as_ref().map(Signer::kid)
}
pub async fn publish(&self, subject: String, payload: Bytes) -> Result<(), PublishError> {
match headers_for(
self.signer.as_ref(),
&payload,
chrono::Utc::now().timestamp_millis(),
) {
Some(headers) => {
self.nats
.publish_with_headers(subject, headers, payload)
.await
}
None => self.nats.publish(subject, payload).await,
}
}
}
fn headers_for(signer: Option<&Signer>, body: &[u8], at_ms: i64) -> Option<HeaderMap> {
let h = signer?.headers(body, at_ms);
let mut map = HeaderMap::new();
for (name, value) in [
(signing::SIG, h.sig_b64.as_deref()),
(signing::SIG_KID, h.kid.as_deref()),
(signing::SIG_ALG, h.alg.as_deref()),
(signing::SIG_AT, h.at_ms.as_deref()),
] {
if let Some(v) = value {
map.insert(name, v);
}
}
Some(map)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Source {
Registry,
Env,
}
impl Source {
fn describe(self) -> String {
match self {
Source::Registry => format!(
"HKLM\\{}\\{{{}, {}}}",
signing::REG_BACKEND_SUBKEY,
signing::REG_SIGNING_KEY,
signing::REG_SIGNING_KID
),
Source::Env => format!("${ENV_SIGNING_KEY} / ${ENV_SIGNING_KID}"),
}
}
}
fn resolve_signer() -> Result<Option<Signer>, String> {
let registry = (
kanade_shared::secrets::read_hklm_value(
signing::REG_BACKEND_SUBKEY,
signing::REG_SIGNING_KEY,
),
kanade_shared::secrets::read_hklm_value(
signing::REG_BACKEND_SUBKEY,
signing::REG_SIGNING_KID,
),
);
let env = (read_env(ENV_SIGNING_KEY), read_env(ENV_SIGNING_KID));
match choose(registry, env)? {
Some((secret, kid, source)) => Signer::from_secret(&secret, &kid)
.map(Some)
.map_err(|e| format!("{} in {}", e, source.describe())),
None => Ok(None),
}
}
fn read_env(name: &str) -> Option<String> {
match std::env::var(name) {
Ok(v) if !v.is_empty() => Some(v),
_ => None,
}
}
fn choose(
registry: (Option<String>, Option<String>),
env: (Option<String>, Option<String>),
) -> Result<Option<(String, String, Source)>, String> {
for (source, (secret, kid)) in [(Source::Registry, registry), (Source::Env, env)] {
match (secret, kid) {
(Some(secret), Some(kid)) => return Ok(Some((secret, kid, source))),
(Some(_), None) => {
return Err(format!(
"a command-signing key is present in {} but its key id is missing. The two \
are only meaningful together — signing under the wrong id is reported by \
every agent as an invalid signature.",
source.describe()
));
}
(None, Some(kid)) => {
return Err(format!(
"a command-signing key id ({kid}) is present in {} but the key itself is \
missing.",
source.describe()
));
}
(None, None) => continue,
}
}
Ok(None)
}
#[cfg(test)]
mod tests {
use super::*;
use kanade_shared::signing::{KeyPolicy, KeyRing, SigHeaders, VerifyError, verify};
fn secret() -> String {
signing::encode_secret(&signing::generate_keypair().unwrap())
}
#[test]
fn a_complete_pair_is_taken_from_whichever_store_has_it() {
let s = secret();
let (got, kid, src) = choose((Some(s.clone()), Some("backend-1".into())), (None, None))
.unwrap()
.unwrap();
assert_eq!(
(got, kid, src),
(s.clone(), "backend-1".to_string(), Source::Registry)
);
let (_, kid, src) = choose((None, None), (Some(s), Some("backend-2".into())))
.unwrap()
.unwrap();
assert_eq!((kid, src), ("backend-2".to_string(), Source::Env));
}
#[test]
fn the_registry_wins_when_both_stores_are_complete() {
let (_, kid, src) = choose(
(Some(secret()), Some("from-registry".into())),
(Some(secret()), Some("from-env".into())),
)
.unwrap()
.unwrap();
assert_eq!((kid.as_str(), src), ("from-registry", Source::Registry));
}
#[test]
fn half_a_pair_is_an_error_rather_than_a_fallback() {
let err = choose(
(Some(secret()), None),
(Some(secret()), Some("from-env".into())),
)
.unwrap_err();
assert!(err.contains("CommandSigningKid"), "{err}");
let err = choose((None, Some("orphan".into())), (None, None)).unwrap_err();
assert!(err.contains("orphan"), "{err}");
}
#[test]
fn nothing_configured_is_not_an_error() {
assert!(choose((None, None), (None, None)).unwrap().is_none());
}
#[test]
fn the_headers_a_publisher_attaches_are_the_ones_an_agent_verifies() {
let key = signing::generate_keypair().unwrap();
let signer = Signer::from_secret(&signing::encode_secret(&key), "backend-1").unwrap();
let body = br#"{"id":"job","request_id":"r1"}"#;
let at = 1_700_000_000_000;
let map = headers_for(Some(&signer), body, at).expect("a signing backend attaches headers");
let headers = SigHeaders {
sig_b64: map.get(signing::SIG).map(|v| v.to_string()),
kid: map.get(signing::SIG_KID).map(|v| v.to_string()),
alg: map.get(signing::SIG_ALG).map(|v| v.to_string()),
at_ms: map.get(signing::SIG_AT).map(|v| v.to_string()),
};
let mut ring = KeyRing::new();
ring.insert(
"backend-1",
key.verifying_key(),
KeyPolicy::backend("backend"),
);
assert_eq!(verify(&ring, body, &headers, at).unwrap().kid, "backend-1");
for name in [
signing::SIG,
signing::SIG_KID,
signing::SIG_ALG,
signing::SIG_AT,
] {
let v = map.get(name).unwrap_or_else(|| panic!("{name} missing"));
assert!(!v.to_string().is_empty(), "{name} is empty");
}
}
#[test]
fn a_non_signing_backend_attaches_nothing_at_all() {
assert!(headers_for(None, b"body", 0).is_none());
assert_eq!(
verify(&KeyRing::new(), b"body", &SigHeaders::default(), 0),
Err(VerifyError::Unsigned)
);
}
}