use std::sync::{Arc, Mutex};
use std::time::SystemTime;
use serde_json::Value;
use sha2::{Digest, Sha256};
pub use nexo_tool_meta::admin::audit::{
AdminAuditResult, AdminAuditRow, AuditTailFilter, AuditTailPage,
};
#[async_trait::async_trait]
pub trait AdminAuditWriter: Send + Sync + std::fmt::Debug {
async fn append(&self, row: AdminAuditRow);
}
#[async_trait::async_trait]
pub trait AdminAuditReader: Send + Sync + std::fmt::Debug {
async fn tail(&self, filter: &AuditTailFilter) -> anyhow::Result<AuditTailPage>;
}
#[derive(Debug, Default, Clone)]
pub struct InMemoryAuditWriter {
rows: Arc<Mutex<Vec<AdminAuditRow>>>,
}
impl InMemoryAuditWriter {
pub fn new() -> Self {
Self::default()
}
pub fn rows(&self) -> Vec<AdminAuditRow> {
self.rows.lock().unwrap().clone()
}
pub fn last(&self) -> Option<AdminAuditRow> {
self.rows.lock().unwrap().last().cloned()
}
}
#[async_trait::async_trait]
impl AdminAuditWriter for InMemoryAuditWriter {
async fn append(&self, row: AdminAuditRow) {
self.rows.lock().unwrap().push(row);
}
}
pub fn hash_params(params: &Value) -> String {
let canonical = canonicalize(params);
let serialized = serde_json::to_string(&canonical).unwrap_or_default();
let digest = Sha256::digest(serialized.as_bytes());
hex_encode(&digest)
}
pub fn redact_for_audit(method: &str, params: &Value) -> Value {
if method == "nexo/admin/secrets/write" {
if let Some(obj) = params.as_object() {
let mut redacted = obj.clone();
if redacted.contains_key("value") {
redacted.insert("value".into(), Value::String("<redacted>".into()));
}
return Value::Object(redacted);
}
}
if method == "nexo/admin/credentials/register" {
if let Some(obj) = params.as_object() {
let mut redacted = obj.clone();
if let Some(payload) = redacted.get("payload").cloned() {
redacted.insert("payload".into(), redact_secret_keys(&payload));
}
if let Some(metadata) = redacted.get("metadata").cloned() {
redacted.insert("metadata".into(), redact_secret_keys(&metadata));
}
return Value::Object(redacted);
}
}
if method == "nexo/admin/llm_providers/upsert" {
if let Some(obj) = params.as_object() {
let mut redacted = obj.clone();
if redacted.contains_key("api_key_secret_value") {
redacted.insert(
"api_key_secret_value".into(),
Value::String("<redacted>".into()),
);
}
if let Some(fields) = redacted.get("fields").cloned() {
redacted.insert("fields".into(), redact_secret_keys(&fields));
}
return Value::Object(redacted);
}
}
params.clone()
}
fn redact_secret_keys(value: &Value) -> Value {
const SECRET_KEYS: &[&str] = &[
"token",
"password",
"xoauth2_token",
"api_key",
"secret",
"setup_token",
"access_token",
"refresh_token",
"oauth_bundle",
];
match value {
Value::Object(map) => {
let mut out = serde_json::Map::with_capacity(map.len());
for (k, v) in map {
if SECRET_KEYS.iter().any(|s| s.eq_ignore_ascii_case(k)) {
out.insert(k.clone(), Value::String("<redacted>".into()));
} else {
out.insert(k.clone(), redact_secret_keys(v));
}
}
Value::Object(out)
}
Value::Array(arr) => Value::Array(arr.iter().map(redact_secret_keys).collect()),
other => other.clone(),
}
}
fn canonicalize(value: &Value) -> Value {
match value {
Value::Object(map) => {
let mut sorted: Vec<(String, Value)> = map
.iter()
.map(|(k, v)| (k.clone(), canonicalize(v)))
.collect();
sorted.sort_by(|a, b| a.0.cmp(&b.0));
Value::Object(sorted.into_iter().collect())
}
Value::Array(arr) => Value::Array(arr.iter().map(canonicalize).collect()),
other => other.clone(),
}
}
fn hex_encode(bytes: &[u8]) -> String {
let mut s = String::with_capacity(bytes.len() * 2);
for byte in bytes {
s.push_str(&format!("{byte:02x}"));
}
s
}
pub fn extract_tenant_id(params: &Value) -> Option<String> {
params
.as_object()
.and_then(|m| m.get("tenant_id"))
.and_then(|v| v.as_str())
.map(|s| s.to_string())
}
pub fn now_epoch_ms() -> u64 {
SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn in_memory_writer_records_each_call() {
let writer = InMemoryAuditWriter::new();
let row = AdminAuditRow {
microapp_id: "agent-creator".into(),
method: "nexo/admin/agents/list".into(),
capability: "agents_crud".into(),
args_hash: "abc123".into(),
started_at_ms: 1_700_000_000_000,
result: AdminAuditResult::Ok,
duration_ms: 5,
tenant_id: None,
};
writer.append(row.clone()).await;
assert_eq!(writer.rows().len(), 1);
assert_eq!(writer.last().unwrap(), row);
}
#[tokio::test]
async fn in_memory_writer_records_denial() {
let writer = InMemoryAuditWriter::new();
writer
.append(AdminAuditRow {
microapp_id: "agent-creator".into(),
method: "nexo/admin/llm_providers/upsert".into(),
capability: "llm_keys_crud".into(),
args_hash: hash_params(&serde_json::json!({})),
started_at_ms: now_epoch_ms(),
result: AdminAuditResult::Denied,
duration_ms: 0,
tenant_id: None,
})
.await;
let last = writer.last().unwrap();
assert_eq!(last.result, AdminAuditResult::Denied);
assert_eq!(last.capability, "llm_keys_crud");
}
#[test]
fn hash_params_is_deterministic_with_key_order() {
let a = serde_json::json!({ "z": 1, "a": 2 });
let b = serde_json::json!({ "a": 2, "z": 1 });
assert_eq!(hash_params(&a), hash_params(&b));
}
#[test]
fn hash_params_differs_for_different_payloads() {
let a = serde_json::json!({ "x": 1 });
let b = serde_json::json!({ "x": 2 });
assert_ne!(hash_params(&a), hash_params(&b));
}
#[test]
fn audit_result_as_str_table() {
assert_eq!(AdminAuditResult::Ok.as_str(), "ok");
assert_eq!(AdminAuditResult::Error.as_str(), "error");
assert_eq!(AdminAuditResult::Denied.as_str(), "denied");
}
#[test]
fn redact_for_audit_redacts_secrets_write_value() {
let original = serde_json::json!({
"name": "MINIMAX_API_KEY",
"value": "sk-leak-this-not"
});
let redacted = redact_for_audit("nexo/admin/secrets/write", &original);
assert_eq!(redacted["name"], "MINIMAX_API_KEY");
assert_eq!(redacted["value"], "<redacted>");
assert_ne!(hash_params(&original), hash_params(&redacted));
}
#[test]
fn redact_for_audit_redacts_llm_upsert_secret_value() {
let original = serde_json::json!({
"id": "minimax-cliente-a",
"base_url": "https://api.minimax.chat/v1",
"factory_type": "minimax",
"api_key_secret_value": "sk-leak-this-not"
});
let redacted = redact_for_audit("nexo/admin/llm_providers/upsert", &original);
assert_eq!(redacted["id"], "minimax-cliente-a");
assert_eq!(redacted["factory_type"], "minimax");
assert_eq!(redacted["api_key_secret_value"], "<redacted>");
assert_ne!(hash_params(&original), hash_params(&redacted));
}
#[test]
fn redact_for_audit_redacts_llm_upsert_schema_fields() {
let original = serde_json::json!({
"id": "minimax-cliente-a",
"base_url": "https://api.minimax.chat/v1",
"factory_type": "minimax",
"auth_mode": "api_key",
"fields": {
"api_key": "sk-leak-this-not",
"group_id": "1234567890123",
"region": "global",
"setup_token": "sk-ant-oat01-leak-too",
}
});
let redacted = redact_for_audit("nexo/admin/llm_providers/upsert", &original);
let fields = &redacted["fields"];
assert_eq!(fields["api_key"], "<redacted>");
assert_eq!(fields["setup_token"], "<redacted>");
assert_eq!(fields["group_id"], "1234567890123");
assert_eq!(fields["region"], "global");
assert_eq!(redacted["auth_mode"], "api_key");
assert_eq!(redacted["factory_type"], "minimax");
}
#[test]
fn redact_for_audit_masks_oauth_bundle_keys_in_fields() {
let original = serde_json::json!({
"id": "anthropic-personal",
"base_url": "https://api.anthropic.com/v1",
"factory_type": "anthropic",
"auth_mode": "oauth_bundle_import",
"fields": {
"oauth_bundle": "{\"access_token\":\"sk-...\",\"refresh_token\":\"r-...\"}",
}
});
let redacted = redact_for_audit("nexo/admin/llm_providers/upsert", &original);
assert_eq!(redacted["fields"]["oauth_bundle"], "<redacted>");
}
#[test]
fn redact_for_audit_keeps_llm_upsert_secret_id_visible() {
let original = serde_json::json!({
"id": "minimax-cliente-a",
"base_url": "https://api.minimax.chat/v1",
"api_key_secret_id": "LLM_MINIMAX_CLIENTE_A"
});
let redacted = redact_for_audit("nexo/admin/llm_providers/upsert", &original);
assert_eq!(redacted["api_key_secret_id"], "LLM_MINIMAX_CLIENTE_A");
}
#[test]
fn redact_for_audit_passthrough_for_other_methods() {
let original = serde_json::json!({"agent_id": "ana"});
let result = redact_for_audit("nexo/admin/agents/get", &original);
assert_eq!(result, original);
}
#[test]
fn redact_for_audit_redacts_telegram_token() {
let original = serde_json::json!({
"channel": "telegram",
"instance": "kate",
"agent_ids": ["kate"],
"payload": { "token": "bot12345:ABCDEF" }
});
let redacted = redact_for_audit("nexo/admin/credentials/register", &original);
assert_eq!(redacted["channel"], "telegram");
assert_eq!(redacted["payload"]["token"], "<redacted>");
assert_ne!(hash_params(&original), hash_params(&redacted));
}
#[test]
fn redact_for_audit_redacts_email_password_and_xoauth2() {
let with_pw = serde_json::json!({
"channel": "email",
"instance": "ops",
"payload": { "address": "ops@example.com", "password": "p@ss" }
});
let r1 = redact_for_audit("nexo/admin/credentials/register", &with_pw);
assert_eq!(r1["payload"]["password"], "<redacted>");
assert_eq!(r1["payload"]["address"], "ops@example.com");
let with_xoauth = serde_json::json!({
"channel": "email",
"instance": "ops",
"payload": { "address": "ops@example.com", "xoauth2_token": "xo.AAA" }
});
let r2 = redact_for_audit("nexo/admin/credentials/register", &with_xoauth);
assert_eq!(r2["payload"]["xoauth2_token"], "<redacted>");
}
#[test]
fn redact_for_audit_redacts_nested_metadata_secrets() {
let original = serde_json::json!({
"channel": "email",
"payload": { "address": "ops@example.com" },
"metadata": {
"imap": { "host": "imap.example.com", "password": "leaked-here" },
"api_key": "should-not-leak",
"tls": "implicit_tls"
}
});
let redacted = redact_for_audit("nexo/admin/credentials/register", &original);
assert_eq!(redacted["metadata"]["imap"]["password"], "<redacted>");
assert_eq!(redacted["metadata"]["api_key"], "<redacted>");
assert_eq!(redacted["metadata"]["imap"]["host"], "imap.example.com");
assert_eq!(redacted["metadata"]["tls"], "implicit_tls");
}
#[test]
fn redact_for_audit_preserves_non_secret_fields() {
let original = serde_json::json!({
"channel": "email",
"instance": "ops",
"agent_ids": ["ana"],
"payload": { "address": "ops@example.com" },
"metadata": {
"imap": { "host": "imap.example.com", "port": 993, "tls": "implicit_tls" }
}
});
let redacted = redact_for_audit("nexo/admin/credentials/register", &original);
assert_eq!(redacted["channel"], "email");
assert_eq!(redacted["agent_ids"][0], "ana");
assert_eq!(redacted["metadata"]["imap"]["host"], "imap.example.com");
assert_eq!(redacted["metadata"]["imap"]["port"], 993);
}
#[test]
fn extract_tenant_id_reads_string_field() {
let p = serde_json::json!({ "tenant_id": "acme", "other": 1 });
assert_eq!(extract_tenant_id(&p), Some("acme".into()));
}
#[test]
fn extract_tenant_id_missing_field_yields_none() {
let p = serde_json::json!({ "agent_id": "a1" });
assert_eq!(extract_tenant_id(&p), None);
}
#[test]
fn extract_tenant_id_non_string_yields_none() {
let p = serde_json::json!({ "tenant_id": 42 });
assert_eq!(extract_tenant_id(&p), None);
}
#[test]
fn extract_tenant_id_non_object_yields_none() {
let p = serde_json::json!(["a", "b"]);
assert_eq!(extract_tenant_id(&p), None);
let p = serde_json::json!("scalar");
assert_eq!(extract_tenant_id(&p), None);
}
#[test]
fn audit_row_serde_skips_none_tenant_id() {
let row = AdminAuditRow {
microapp_id: "a".into(),
method: "nexo/admin/echo".into(),
capability: "echo".into(),
args_hash: "h".into(),
started_at_ms: 1,
result: AdminAuditResult::Ok,
duration_ms: 1,
tenant_id: None,
};
let json = serde_json::to_string(&row).unwrap();
assert!(!json.contains("tenant_id"), "None tenant skips field");
}
#[test]
fn audit_row_serde_emits_some_tenant_id() {
let row = AdminAuditRow {
microapp_id: "a".into(),
method: "nexo/admin/echo".into(),
capability: "echo".into(),
args_hash: "h".into(),
started_at_ms: 1,
result: AdminAuditResult::Ok,
duration_ms: 1,
tenant_id: Some("acme".into()),
};
let json = serde_json::to_string(&row).unwrap();
assert!(json.contains("\"tenant_id\":\"acme\""));
}
}