use std::collections::VecDeque;
use std::sync::Mutex;
use std::time::{SystemTime, UNIX_EPOCH};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AuditAction {
Hit,
Miss,
Set,
Delete,
Evict,
Expired,
Clear,
}
impl AuditAction {
pub fn as_str(&self) -> &'static str {
match self {
AuditAction::Hit => "hit",
AuditAction::Miss => "miss",
AuditAction::Set => "set",
AuditAction::Delete => "delete",
AuditAction::Evict => "evict",
AuditAction::Expired => "expired",
AuditAction::Clear => "clear",
}
}
}
#[derive(Debug, Clone)]
pub struct AuditEvent {
pub action: AuditAction,
pub key: Option<String>,
pub namespace: Option<String>,
pub operator: Option<String>,
pub timestamp_ms: u64,
pub metadata: Vec<(String, String)>,
}
impl AuditEvent {
pub fn new(action: AuditAction) -> Self {
Self {
action,
key: None,
namespace: None,
operator: None,
timestamp_ms: now_ms(),
metadata: Vec::new(),
}
}
pub fn with_key(mut self, redacted_key: impl Into<String>) -> Self {
self.key = Some(redacted_key.into());
self
}
pub fn with_namespace(mut self, namespace: impl Into<String>) -> Self {
self.namespace = Some(namespace.into());
self
}
pub fn with_operator(mut self, operator: impl Into<String>) -> Self {
self.operator = Some(operator.into());
self
}
pub fn with_metadata(mut self, k: impl Into<String>, v: impl Into<String>) -> Self {
self.metadata.push((k.into(), v.into()));
self
}
}
fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or(std::time::Duration::ZERO)
.as_millis() as u64
}
pub fn redact_key_for_audit(key: &str) -> String {
const SENSITIVE: [&str; 4] = ["token", "password", "secret", "api_key"];
const MAX_LEN: usize = 32;
let lower = key.to_ascii_lowercase();
if SENSITIVE.iter().any(|p| lower.contains(p)) {
let suffix = if key.len() > 2 {
&key[key.len() - 2..]
} else {
key
};
return format!("<sensitive>…{suffix}");
}
if key.len() > MAX_LEN {
format!("{}…(len={})", &key[..MAX_LEN], key.len())
} else {
key.to_string()
}
}
pub trait AuditEventPublisher: Send + Sync + 'static {
fn publish(&self, event: AuditEvent);
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NoOpAuditPublisher;
impl AuditEventPublisher for NoOpAuditPublisher {
fn publish(&self, _event: AuditEvent) {}
}
pub struct InMemoryAuditPublisher {
events: Mutex<VecDeque<AuditEvent>>,
capacity: usize,
}
impl InMemoryAuditPublisher {
pub fn new(capacity: usize) -> Self {
Self {
events: Mutex::new(VecDeque::new()),
capacity: capacity.max(1),
}
}
pub fn snapshot(&self) -> Vec<AuditEvent> {
self.events
.lock()
.map(|q| q.iter().cloned().collect())
.unwrap_or_default()
}
pub fn len(&self) -> usize {
self.events.lock().map(|q| q.len()).unwrap_or(0)
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
impl AuditEventPublisher for InMemoryAuditPublisher {
fn publish(&self, event: AuditEvent) {
if let Ok(mut q) = self.events.lock() {
if q.len() >= self.capacity {
q.pop_front();
}
q.push_back(event);
}
}
}
#[cfg(feature = "telemetry")]
pub struct TracingAuditPublisher;
#[cfg(feature = "telemetry")]
impl AuditEventPublisher for TracingAuditPublisher {
fn publish(&self, event: AuditEvent) {
use crate::i18n::messages::{MSG_LOG_AUDIT_EVENT, t};
tracing::info!(
target: "oxcache::audit",
action = event.action.as_str(),
key = event.key.as_deref().unwrap_or(""),
namespace = event.namespace.as_deref().unwrap_or(""),
operator = event.operator.as_deref().unwrap_or(""),
timestamp_ms = event.timestamp_ms,
"{}",
t(MSG_LOG_AUDIT_EVENT, &[])
);
}
}
#[cfg(feature = "inklog")]
pub const BRIDGE_CHANNEL_CAPACITY: usize = 1024;
#[cfg(feature = "inklog")]
use std::sync::atomic::{AtomicU64, Ordering};
#[cfg(feature = "inklog")]
use std::sync::{Arc, Once};
#[cfg(feature = "inklog")]
struct InklogBridgeInner {
sink: Arc<dyn inklog::sink::LogSink>,
tx: tokio::sync::mpsc::Sender<inklog::LogRecord>,
rx: Mutex<Option<tokio::sync::mpsc::Receiver<inklog::LogRecord>>>,
writer_started: Once,
dropped: Arc<AtomicU64>,
write_failures: Arc<AtomicU64>,
}
#[cfg(feature = "inklog")]
struct WriterRxGuard {
rx: tokio::sync::mpsc::Receiver<inklog::LogRecord>,
dropped: Arc<AtomicU64>,
}
#[cfg(feature = "inklog")]
impl Drop for WriterRxGuard {
fn drop(&mut self) {
while self.rx.try_recv().is_ok() {
self.dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
#[cfg(feature = "inklog")]
#[derive(Clone)]
pub struct InklogAuditPublisher {
inner: Arc<InklogBridgeInner>,
}
#[cfg(feature = "inklog")]
impl InklogAuditPublisher {
pub fn new(sink: Arc<dyn inklog::sink::LogSink>) -> Self {
let (tx, rx) = tokio::sync::mpsc::channel(BRIDGE_CHANNEL_CAPACITY);
Self {
inner: Arc::new(InklogBridgeInner {
sink,
tx,
rx: Mutex::new(Some(rx)),
writer_started: Once::new(),
dropped: Arc::new(AtomicU64::new(0)),
write_failures: Arc::new(AtomicU64::new(0)),
}),
}
}
pub fn dropped_count(&self) -> u64 {
self.inner.dropped.load(Ordering::Relaxed)
}
pub fn write_failure_count(&self) -> u64 {
self.inner.write_failures.load(Ordering::Relaxed)
}
}
#[cfg(feature = "inklog")]
fn audit_event_to_log_record(event: &AuditEvent) -> inklog::LogRecord {
use crate::i18n::messages::{MSG_LOG_AUDIT_EVENT, t};
use serde_json::Value;
let level = match event.action {
AuditAction::Set | AuditAction::Delete | AuditAction::Clear => inklog::tracing::Level::INFO,
AuditAction::Hit | AuditAction::Miss | AuditAction::Evict | AuditAction::Expired => {
inklog::tracing::Level::DEBUG
}
};
let mut record = inklog::LogRecord::new(
level,
"oxcache::audit".to_string(),
t(MSG_LOG_AUDIT_EVENT, &[]),
);
record
.fields
.insert("action".to_string(), Value::from(event.action.as_str()));
if let Some(key) = &event.key {
record
.fields
.insert("key".to_string(), Value::from(key.as_str()));
}
if let Some(namespace) = &event.namespace {
record
.fields
.insert("namespace".to_string(), Value::from(namespace.as_str()));
}
if let Some(operator) = &event.operator {
record
.fields
.insert("operator".to_string(), Value::from(operator.as_str()));
}
record
.fields
.insert("timestamp_ms".to_string(), Value::from(event.timestamp_ms));
for (k, v) in &event.metadata {
record
.fields
.insert(format!("meta.{k}"), Value::from(v.as_str()));
}
record
}
#[cfg(feature = "inklog")]
impl AuditEventPublisher for InklogAuditPublisher {
fn publish(&self, event: AuditEvent) {
let record = audit_event_to_log_record(&event);
match tokio::runtime::Handle::try_current() {
Ok(handle) => {
self.inner.writer_started.call_once(|| {
use crate::i18n::messages::{
MSG_PANIC_AUDIT_WRITER_LOCK, MSG_PANIC_AUDIT_WRITER_RX, t,
};
let rx = self
.inner
.rx
.lock()
.unwrap_or_else(|e| {
panic!("{}: {e:?}", t(MSG_PANIC_AUDIT_WRITER_LOCK, &[]))
})
.take()
.unwrap_or_else(|| panic!("{}", t(MSG_PANIC_AUDIT_WRITER_RX, &[])));
let sink = self.inner.sink.clone();
let dropped = self.inner.dropped.clone();
let write_failures = self.inner.write_failures.clone();
handle.spawn(async move {
let mut guard = WriterRxGuard { rx, dropped };
while let Some(record) = guard.rx.recv().await {
if sink.write(&record).await.is_err() {
write_failures.fetch_add(1, Ordering::Relaxed);
}
}
});
});
match self.inner.tx.try_send(record) {
Ok(()) => {}
Err(tokio::sync::mpsc::error::TrySendError::Full(_))
| Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => {
self.inner.dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
Err(_) => {
self.inner.dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
#[test]
fn noop_publisher_accepts_events() {
let publisher = NoOpAuditPublisher;
publisher.publish(AuditEvent::new(AuditAction::Set).with_key("k"));
publisher.publish(AuditEvent::new(AuditAction::Evict));
}
#[test]
fn in_memory_publisher_ring_buffers() {
let publisher = InMemoryAuditPublisher::new(3);
for i in 0..5u32 {
publisher.publish(AuditEvent::new(AuditAction::Hit).with_key(format!("k{i}")));
}
assert_eq!(publisher.len(), 3, "环形缓冲应保持容量上限");
let events = publisher.snapshot();
assert_eq!(events[0].key.as_deref(), Some("k2"));
assert_eq!(events[2].key.as_deref(), Some("k4"));
}
#[test]
fn events_carry_action_key_timestamp_and_metadata() {
let event = AuditEvent::new(AuditAction::Expired)
.with_key("user:1")
.with_namespace("users")
.with_operator("svc-cache")
.with_metadata("reason", "ttl");
assert_eq!(event.action, AuditAction::Expired);
assert_eq!(event.key.as_deref(), Some("user:1"));
assert_eq!(event.namespace.as_deref(), Some("users"));
assert_eq!(event.operator.as_deref(), Some("svc-cache"));
assert_eq!(event.metadata[0], ("reason".to_string(), "ttl".to_string()));
assert!(event.timestamp_ms > 0);
}
#[test]
fn action_display_names() {
assert_eq!(AuditAction::Hit.as_str(), "hit");
assert_eq!(AuditAction::Miss.as_str(), "miss");
assert_eq!(AuditAction::Evict.as_str(), "evict");
assert_eq!(AuditAction::Expired.as_str(), "expired");
}
#[test]
fn sensitive_keys_are_masked() {
let redacted = redact_key_for_audit("user:api_key:abcdef");
assert!(redacted.starts_with("<sensitive>"), "got {redacted}");
assert!(!redacted.contains("abcdef"), "敏感值不得完整出现");
let long = "k".repeat(100);
let redacted = redact_key_for_audit(&long);
assert!(redacted.contains("len=100"), "长键应截断并带长度标注");
assert!(redacted.len() < 50);
let normal = redact_key_for_audit("user:1");
assert_eq!(normal, "user:1");
}
#[test]
fn publisher_is_object_safe() {
let concrete = Arc::new(InMemoryAuditPublisher::new(8));
let publisher: Arc<dyn AuditEventPublisher> = concrete.clone();
publisher.publish(AuditEvent::new(AuditAction::Clear));
assert_eq!(concrete.len(), 1);
}
#[cfg(feature = "inklog")]
mod inklog_bridge {
use super::*;
use crate::features::audit::{InklogAuditPublisher, audit_event_to_log_record};
use inklog::sink::LogSink;
use std::collections::HashMap;
use std::sync::mpsc::Sender;
#[derive(Debug, Clone)]
struct Captured {
target: String,
level: String,
fields: HashMap<String, serde_json::Value>,
}
struct CollectingSink {
captured: Arc<Mutex<Vec<Captured>>>,
done: Sender<()>,
fail: bool,
}
#[async_trait::async_trait]
impl LogSink for CollectingSink {
async fn write(&self, record: &inklog::LogRecord) -> Result<(), inklog::InklogError> {
if !self.fail {
self.captured.lock().unwrap().push(Captured {
target: record.target.clone(),
level: record.level.clone(),
fields: record.fields.clone(),
});
}
let _ = self.done.send(());
if self.fail {
Err(inklog::InklogError::RuntimeError("sink write fault".into()))
} else {
Ok(())
}
}
async fn flush(&self) -> Result<(), inklog::InklogError> {
Ok(())
}
async fn shutdown(&self) -> Result<(), inklog::InklogError> {
Ok(())
}
}
fn publisher_with_sink(
fail: bool,
) -> (
InklogAuditPublisher,
Arc<Mutex<Vec<Captured>>>,
std::sync::mpsc::Receiver<()>,
) {
let captured = Arc::new(Mutex::new(Vec::new()));
let (tx, rx) = std::sync::mpsc::channel();
let sink = Arc::new(CollectingSink {
captured: captured.clone(),
done: tx,
fail,
});
(InklogAuditPublisher::new(sink), captured, rx)
}
#[tokio::test(flavor = "multi_thread")]
async fn publish_forwards_audit_event_as_log_record() {
let (publisher, captured, rx) = publisher_with_sink(false);
let dyn_publisher: Arc<dyn AuditEventPublisher> = Arc::new(publisher.clone());
dyn_publisher.publish(
AuditEvent::new(AuditAction::Set)
.with_key("user:1")
.with_namespace("users")
.with_operator("svc-a")
.with_metadata("reason", "test"),
);
rx.recv_timeout(std::time::Duration::from_secs(5))
.expect("写入完成信号");
let snapshot = captured.lock().unwrap();
assert_eq!(snapshot.len(), 1);
let record = &snapshot[0];
assert_eq!(record.target, "oxcache::audit");
assert_eq!(record.level, "INFO", "写操作(Set)应映射 INFO");
assert_eq!(
record.fields.get("action").and_then(|v| v.as_str()),
Some("set")
);
assert_eq!(
record.fields.get("key").and_then(|v| v.as_str()),
Some("user:1")
);
assert_eq!(
record.fields.get("namespace").and_then(|v| v.as_str()),
Some("users")
);
assert_eq!(
record.fields.get("operator").and_then(|v| v.as_str()),
Some("svc-a")
);
assert_eq!(
record.fields.get("meta.reason").and_then(|v| v.as_str()),
Some("test")
);
assert!(record.fields.contains_key("timestamp_ms"));
assert_eq!(publisher.write_failure_count(), 0);
assert_eq!(publisher.dropped_count(), 0);
}
#[tokio::test(flavor = "multi_thread")]
async fn read_actions_map_to_debug_level() {
let (publisher, captured, rx) = publisher_with_sink(false);
publisher.publish(AuditEvent::new(AuditAction::Hit).with_key("k"));
publisher.publish(AuditEvent::new(AuditAction::Miss).with_key("k2"));
publisher.publish(AuditEvent::new(AuditAction::Evict).with_key("k3"));
for _ in 0..3 {
rx.recv_timeout(std::time::Duration::from_secs(5))
.expect("写入完成信号");
}
let snapshot = captured.lock().unwrap();
assert_eq!(snapshot.len(), 3);
for record in snapshot.iter() {
assert_eq!(record.level, "DEBUG", "读探测/容量清理应映射 DEBUG");
}
}
#[tokio::test(flavor = "multi_thread")]
async fn sink_write_failure_is_counted() {
let (publisher, captured, rx) = publisher_with_sink(true);
publisher.publish(AuditEvent::new(AuditAction::Set).with_key("k"));
rx.recv_timeout(std::time::Duration::from_secs(5))
.expect("写入尝试完成信号");
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(1);
while publisher.write_failure_count() == 0 && std::time::Instant::now() < deadline {
std::thread::sleep(std::time::Duration::from_millis(2));
}
assert!(
captured.lock().unwrap().is_empty(),
"失败写入不得留下半条记录"
);
assert_eq!(publisher.write_failure_count(), 1, "sink 写失败应显性计数");
}
#[test]
fn publish_without_runtime_is_dropped_and_counted() {
let (publisher, captured, _rx) = publisher_with_sink(false);
let moved = publisher.clone();
std::thread::spawn(move || {
moved.publish(AuditEvent::new(AuditAction::Set).with_key("no-runtime"));
})
.join()
.expect("无 runtime publish 不得 panic");
assert_eq!(
publisher.dropped_count(),
1,
"无 runtime 上下文的事件应计入 dropped"
);
assert!(captured.lock().unwrap().is_empty());
}
struct GatedSink {
started: std::sync::mpsc::Sender<()>,
gate: Arc<tokio::sync::Notify>,
}
#[async_trait::async_trait]
impl LogSink for GatedSink {
async fn write(&self, _record: &inklog::LogRecord) -> Result<(), inklog::InklogError> {
let _ = self.started.send(());
self.gate.notified().await;
Ok(())
}
async fn flush(&self) -> Result<(), inklog::InklogError> {
Ok(())
}
async fn shutdown(&self) -> Result<(), inklog::InklogError> {
Ok(())
}
}
fn gated_publisher() -> (
InklogAuditPublisher,
std::sync::mpsc::Receiver<()>,
Arc<tokio::sync::Notify>,
) {
let (started_tx, started_rx) = std::sync::mpsc::channel();
let gate = Arc::new(tokio::sync::Notify::new());
(
InklogAuditPublisher::new(Arc::new(GatedSink {
started: started_tx,
gate: gate.clone(),
})),
started_rx,
gate,
)
}
#[tokio::test(flavor = "multi_thread")]
async fn channel_overload_drops_are_counted() {
let (publisher, started_rx, gate) = gated_publisher();
publisher.publish(AuditEvent::new(AuditAction::Set).with_key("first"));
started_rx
.recv_timeout(std::time::Duration::from_secs(5))
.expect("writer 应已消费首条事件并被闸门挂起");
for _ in 0..BRIDGE_CHANNEL_CAPACITY {
publisher.publish(AuditEvent::new(AuditAction::Miss).with_key("fill"));
}
let overflow = 7u64;
for _ in 0..overflow {
publisher.publish(AuditEvent::new(AuditAction::Miss).with_key("overflow"));
}
assert_eq!(
publisher.dropped_count(),
overflow,
"通道满后的 publish 应显性丢弃计数而非阻塞"
);
assert_eq!(publisher.write_failure_count(), 0);
gate.notify_waiters();
}
#[test]
fn runtime_shutdown_losses_are_counted() {
let (publisher, started_rx, _gate) = gated_publisher();
let rt1 = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.unwrap();
rt1.block_on(async {
publisher.publish(AuditEvent::new(AuditAction::Set).with_key("in-write"));
});
started_rx
.recv_timeout(std::time::Duration::from_secs(5))
.expect("writer 应已被钉在闸门写入内");
rt1.block_on(async {
publisher.publish(AuditEvent::new(AuditAction::Set).with_key("buffered-1"));
publisher.publish(AuditEvent::new(AuditAction::Set).with_key("buffered-2"));
});
drop(rt1);
let rt2 = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.unwrap();
rt2.block_on(async {
publisher.publish(AuditEvent::new(AuditAction::Set).with_key("after-shutdown"));
});
assert_eq!(
publisher.dropped_count(),
3,
"滞留事件 2 条 + 通道关闭后 1 条均应显性计数"
);
assert_eq!(publisher.write_failure_count(), 0);
}
#[test]
fn level_mapping_is_explicit_and_total() {
let cases = [
(AuditAction::Set, "INFO"),
(AuditAction::Delete, "INFO"),
(AuditAction::Clear, "INFO"),
(AuditAction::Hit, "DEBUG"),
(AuditAction::Miss, "DEBUG"),
(AuditAction::Evict, "DEBUG"),
(AuditAction::Expired, "DEBUG"),
];
for (action, expected) in cases {
let record = audit_event_to_log_record(&AuditEvent::new(action).with_key("k"));
assert_eq!(record.level, expected, "{action:?} 级别映射");
assert_eq!(record.target, "oxcache::audit");
}
}
#[test]
fn optional_fields_omit_empty_values() {
let record = audit_event_to_log_record(&AuditEvent::new(AuditAction::Miss));
assert!(!record.fields.contains_key("key"));
assert!(!record.fields.contains_key("namespace"));
assert!(!record.fields.contains_key("operator"));
assert_eq!(
record.fields.get("action").and_then(|v| v.as_str()),
Some("miss")
);
}
}
}