use kanade_shared::ipc::error::{ErrorKind, RpcError};
use kanade_shared::ipc::support::{
SupportLockParams, SupportLockResult, SupportStatusParams, SupportStatusResult,
SupportUnlockParams, SupportUnlockResult,
};
use kanade_shared::kv::{
BUCKET_FLEET_CONFIG, BUCKET_SERVER_SETTINGS, KEY_SERVER_SETTINGS, KEY_SUPPORT_CODES,
};
use kanade_shared::wire::{ObsEvent, ServerSettings, SupportCode, SupportCodesProjection};
use tracing::{info, warn};
use super::super::connection::ConnectionState;
use super::super::unlock;
use super::system::HandlerResult;
use crate::obs_outbox;
pub async fn handle_support_unlock(
conn: &ConnectionState,
params: SupportUnlockParams,
) -> HandlerResult<SupportUnlockResult> {
let sid = conn.peer.user_sid.clone();
if let Some(remaining) = unlock::lockout_remaining(&sid) {
warn!(
user = %conn.peer.user,
secs = remaining.as_secs(),
"support.unlock: refused, caller is rate-limited",
);
return Err(RpcError::new(
ErrorKind::RateLimit,
format!(
"too many failed attempts; try again in {} seconds",
remaining.as_secs().max(1)
),
));
}
let code = params.code.trim().to_string();
if code.is_empty() {
return Err(RpcError::new(
ErrorKind::InvalidParams,
"support.unlock: code must not be empty",
));
}
let client = conn.nats.as_ref().ok_or_else(|| {
RpcError::new(
ErrorKind::InternalError,
"support.unlock: NATS client not wired into the connection",
)
})?;
let codes = read_support_codes(client).await?;
let usable: Vec<SupportCode> = codes.iter().filter(|c| c.is_usable()).cloned().collect();
let matched = tokio::task::spawn_blocking(move || match_code(&code, &usable))
.await
.map_err(|e| {
warn!(error = %e, "support.unlock: verify task failed");
RpcError::new(ErrorKind::InternalError, "support.unlock: verify failed")
})?;
let Some(code) = matched else {
let lockout = unlock::record_failure(&sid);
warn!(
user = %conn.peer.user,
locked_out = lockout.is_some(),
"support.unlock: rejected an unrecognised code",
);
audit(conn, "support_unlock_failed", serde_json::json!({}));
return Err(RpcError::new(
ErrorKind::Unauthorized,
"support.unlock: code not recognised",
));
};
unlock::clear_failures(&sid);
let ttl = code.effective_ttl_minutes();
let grants = unlock::grant(&sid, &code.scope, code.label.clone(), ttl);
info!(
user = %conn.peer.user,
scope = %code.scope,
ttl_minutes = ttl,
"support.unlock: granted",
);
audit(
conn,
"support_unlock",
serde_json::json!({ "scope": code.scope, "ttl_minutes": ttl, "label": code.label }),
);
Ok(SupportUnlockResult { grants })
}
pub fn handle_support_lock(
conn: &ConnectionState,
_params: SupportLockParams,
) -> HandlerResult<SupportLockResult> {
let released = unlock::lock(&conn.peer.user_sid);
if released > 0 {
info!(user = %conn.peer.user, released, "support.lock: grants released");
audit(
conn,
"support_lock",
serde_json::json!({ "released": released }),
);
}
Ok(SupportLockResult { released })
}
pub fn handle_support_status(
conn: &ConnectionState,
_params: SupportStatusParams,
) -> HandlerResult<SupportStatusResult> {
Ok(SupportStatusResult {
grants: unlock::grants(&conn.peer.user_sid),
})
}
fn match_code(code: &str, codes: &[SupportCode]) -> Option<SupportCode> {
use argon2::{Argon2, PasswordHash, PasswordVerifier};
let mut hit: Option<SupportCode> = None;
for candidate in codes {
let parsed = match PasswordHash::new(&candidate.hash) {
Ok(p) => p,
Err(e) => {
warn!(scope = %candidate.scope, error = %e, "support code hash is unparseable");
continue;
}
};
if Argon2::default()
.verify_password(code.as_bytes(), &parsed)
.is_ok()
&& hit.is_none()
{
hit = Some(candidate.clone());
}
}
hit
}
async fn read_support_codes(client: &async_nats::Client) -> HandlerResult<Vec<SupportCode>> {
let js = async_nats::jetstream::new(client.clone());
let projection = read_projection_bytes(&js).await;
select_codes(projection, || read_legacy_codes(&js)).await
}
async fn select_codes<F, Fut>(
projection: HandlerResult<Option<Vec<u8>>>,
legacy: F,
) -> HandlerResult<Vec<SupportCode>>
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = HandlerResult<Vec<SupportCode>>>,
{
match projection? {
Some(bytes) => serde_json::from_slice::<SupportCodesProjection>(&bytes)
.map(|p| p.support_codes)
.map_err(|e| {
warn!(error = %e, "support.unlock: decode support_codes projection");
RpcError::new(
ErrorKind::InternalError,
"support.unlock: server settings are corrupt",
)
}),
None => legacy().await,
}
}
async fn read_projection_bytes(
js: &async_nats::jetstream::Context,
) -> HandlerResult<Option<Vec<u8>>> {
let kv = js.get_key_value(BUCKET_FLEET_CONFIG).await.map_err(|e| {
warn!(error = %e, "support.unlock: open fleet_config bucket");
RpcError::new(
ErrorKind::InternalError,
"support.unlock: server settings unavailable",
)
})?;
match kv.get(KEY_SUPPORT_CODES).await {
Ok(v) => Ok(v.map(|b| b.to_vec())),
Err(e) => {
warn!(error = %e, "support.unlock: read support_codes projection");
Err(RpcError::new(
ErrorKind::InternalError,
"support.unlock: server settings unavailable",
))
}
}
}
async fn read_legacy_codes(js: &async_nats::jetstream::Context) -> HandlerResult<Vec<SupportCode>> {
let kv = js
.get_key_value(BUCKET_SERVER_SETTINGS)
.await
.map_err(|e| {
warn!(error = %e, "support.unlock: open server_settings bucket");
RpcError::new(
ErrorKind::InternalError,
"support.unlock: server settings unavailable",
)
})?;
match kv.get(KEY_SERVER_SETTINGS).await {
Ok(Some(bytes)) => serde_json::from_slice::<ServerSettings>(&bytes)
.map(|s| s.support_codes)
.map_err(|e| {
warn!(error = %e, "support.unlock: decode server_settings");
RpcError::new(
ErrorKind::InternalError,
"support.unlock: server settings are corrupt",
)
}),
Ok(None) => Ok(Vec::new()),
Err(e) => {
warn!(error = %e, "support.unlock: read server_settings");
Err(RpcError::new(
ErrorKind::InternalError,
"support.unlock: server settings unavailable",
))
}
}
}
fn audit(conn: &ConnectionState, kind: &str, mut payload: serde_json::Value) {
if let Some(obj) = payload.as_object_mut() {
obj.insert("user".into(), conn.peer.user.clone().into());
obj.insert("user_sid".into(), conn.peer.user_sid.clone().into());
}
let event = ObsEvent {
pc_id: conn.pc_id.clone(),
at: chrono::Utc::now(),
kind: kind.to_string(),
source: "agent:support".to_string(),
event_record_id: Some(format!("support_{}", uuid::Uuid::new_v4().simple())),
payload,
};
let dir = obs_outbox::default_dir();
let res = obs_outbox::ensure_outbox_dir(&dir)
.and_then(|()| obs_outbox::enqueue(&dir, &event).map(|_| ()));
if let Err(e) = res {
warn!(error = %e, kind, "failed to queue support audit event");
}
}
#[cfg(test)]
mod tests {
use super::*;
fn hash_of(code: &str) -> 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)
.unwrap()
.to_string()
}
fn code(scope: &str, plain: &str) -> SupportCode {
SupportCode {
scope: scope.into(),
hash: hash_of(plain),
label: None,
ttl_minutes: None,
disabled: false,
}
}
#[test]
fn matches_the_right_scope() {
let codes = vec![code("support", "hunter2"), code("admin", "correct-horse")];
assert_eq!(match_code("hunter2", &codes).unwrap().scope, "support");
assert_eq!(match_code("correct-horse", &codes).unwrap().scope, "admin");
}
#[test]
fn rejects_a_wrong_code() {
let codes = vec![code("support", "hunter2")];
assert!(match_code("hunter3", &codes).is_none());
assert!(match_code("", &codes).is_none());
}
#[test]
fn no_configured_codes_matches_nothing() {
assert!(match_code("anything", &[]).is_none());
}
#[test]
fn skips_an_unparseable_hash_without_matching() {
let codes = vec![SupportCode {
scope: "broken".into(),
hash: "not-a-phc-string".into(),
..Default::default()
}];
assert!(match_code("not-a-phc-string", &codes).is_none());
assert!(match_code("anything", &codes).is_none());
}
#[test]
fn a_later_scope_still_matches_after_an_earlier_miss() {
let codes = vec![code("a", "aaa"), code("b", "bbb"), code("c", "ccc")];
assert_eq!(match_code("ccc", &codes).unwrap().scope, "c");
}
use std::sync::atomic::{AtomicUsize, Ordering};
fn projection_bytes(codes: Vec<SupportCode>) -> Vec<u8> {
serde_json::to_vec(&SupportCodesProjection {
support_codes: codes,
})
.unwrap()
}
async fn select(
projection: HandlerResult<Option<Vec<u8>>>,
legacy: HandlerResult<Vec<SupportCode>>,
) -> (HandlerResult<Vec<SupportCode>>, usize) {
let hits = AtomicUsize::new(0);
let out = select_codes(projection, || async {
hits.fetch_add(1, Ordering::SeqCst);
legacy
})
.await;
(out, hits.load(Ordering::SeqCst))
}
fn unavailable() -> RpcError {
RpcError::new(ErrorKind::InternalError, "unavailable")
}
async fn unlock_with(typed: &str, projection: Vec<SupportCode>) -> Option<String> {
let (codes, _) = select(Ok(Some(projection_bytes(projection))), Err(unavailable())).await;
let usable: Vec<SupportCode> = codes
.unwrap()
.into_iter()
.filter(|c| c.is_usable())
.collect();
match_code(typed, &usable).map(|c| c.scope)
}
#[tokio::test]
async fn projection_unlocks_with_the_right_code_and_refuses_a_wrong_one() {
let codes = vec![code("support", "hunter2")];
assert_eq!(
unlock_with("hunter2", codes.clone()).await.as_deref(),
Some("support")
);
assert_eq!(unlock_with("hunter3", codes).await, None);
}
#[tokio::test]
async fn a_disabled_scope_does_not_unlock() {
let mut c = code("support", "hunter2");
c.disabled = true;
assert_eq!(unlock_with("hunter2", vec![c]).await, None);
}
#[tokio::test]
async fn present_projection_never_touches_the_legacy_document() {
let (out, hits) = select(
Ok(Some(projection_bytes(vec![code("a", "aaa")]))),
Ok(vec![code("old", "old")]),
)
.await;
assert_eq!(out.unwrap()[0].scope, "a");
assert_eq!(hits, 0);
}
#[tokio::test]
async fn empty_projection_is_authoritative() {
let (out, hits) = select(
Ok(Some(projection_bytes(vec![]))),
Ok(vec![code("old", "o")]),
)
.await;
assert!(out.unwrap().is_empty());
assert_eq!(hits, 0);
}
#[tokio::test]
async fn unreadable_projection_is_unavailable_without_fallback() {
let (out, hits) = select(Err(unavailable()), Ok(vec![code("old", "o")])).await;
assert!(out.is_err());
assert_eq!(hits, 0);
}
#[tokio::test]
async fn corrupt_projection_is_surfaced_without_fallback() {
let (out, hits) = select(Ok(Some(b"{not json".to_vec())), Ok(vec![code("old", "o")])).await;
assert!(out.is_err());
assert_eq!(hits, 0);
}
#[tokio::test]
async fn absent_projection_falls_back_to_the_legacy_document() {
let (out, hits) = select(Ok(None), Ok(vec![code("old", "o")])).await;
assert_eq!(out.unwrap()[0].scope, "old");
assert_eq!(hits, 1);
let (out, hits) = select(Ok(None), Err(unavailable())).await;
assert!(out.is_err());
assert_eq!(hits, 1);
}
}