use std::path::PathBuf;
use tokio::sync::mpsc;
use super::super::cloud::tamper::TamperEvent;
#[derive(Clone)]
pub struct TamperLogger {
tx: mpsc::Sender<TamperEvent>,
}
pub struct TamperLoggerHandle {
_join_handle: tokio::task::JoinHandle<()>,
}
impl TamperLoggerHandle {
pub fn from_task(join_handle: tokio::task::JoinHandle<()>) -> Self {
Self {
_join_handle: join_handle,
}
}
}
impl TamperLogger {
pub fn new(log_dir: PathBuf) -> (Self, TamperLoggerHandle) {
let (logger, mut rx) = Self::channel();
let handle = TamperLoggerHandle {
_join_handle: tokio::spawn(async move { run_tamper_writer(log_dir, &mut rx).await }),
};
(logger, handle)
}
pub fn channel() -> (Self, mpsc::Receiver<TamperEvent>) {
let (tx, rx) = mpsc::channel(256);
(Self { tx }, rx)
}
pub fn log(&self, event: TamperEvent) {
if self.tx.try_send(event).is_err() {
tracing::warn!("tamper log channel full or closed, event dropped");
}
}
}
pub async fn run_tamper_writer(log_dir: PathBuf, rx: &mut mpsc::Receiver<TamperEvent>) {
let path = log_dir.join("tamper.jsonl");
if let Some(parent) = path.parent() {
let _ = std::fs::create_dir_all(parent);
}
while let Some(event) = rx.recv().await {
match serde_json::to_string(&event) {
Ok(json_line) => {
use std::io::Write;
let mut file = match std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&path)
{
Ok(f) => f,
Err(e) => {
tracing::warn!(error = %e, "cannot open tamper.jsonl");
continue;
}
};
if let Err(e) = writeln!(file, "{json_line}") {
tracing::warn!(error = %e, "cannot write to tamper.jsonl");
}
}
Err(e) => {
tracing::warn!(error = %e, "cannot serialize tamper event");
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::cloud::tamper::{FieldDelta, HealOutcome, OcsfHeader, TamperPayload};
use chrono::{TimeZone, Utc};
fn event(id: &str, related: Option<&str>, deltas: Vec<FieldDelta>) -> TamperEvent {
TamperEvent {
ocsf: OcsfHeader {
class_uid: 2004,
activity_id: 1,
type_uid: 200402,
severity_id: 4,
time: Utc.with_ymd_and_hms(2026, 1, 2, 3, 4, 5).unwrap(),
},
tamper: TamperPayload {
event_id: id.to_string(),
entry_id: "entry-1".to_string(),
agent_type: "claude-code".to_string(),
settings_path_hash: "abc123".to_string(),
hook_event: "PreToolUse".to_string(),
detection_method: "hash_mismatch".to_string(),
field_deltas: deltas,
heal: HealOutcome {
outcome: "healed".to_string(),
attempt: 1,
circuit: "closed".to_string(),
},
related_event_id: related.map(str::to_string),
},
}
}
fn read_lines(dir: &std::path::Path) -> Vec<serde_json::Value> {
std::fs::read_to_string(dir.join("tamper.jsonl"))
.unwrap()
.lines()
.map(|l| serde_json::from_str(l).expect("each line is standalone JSON"))
.collect()
}
#[tokio::test]
async fn writer_appends_one_json_line_per_event_in_order() {
let dir = tempfile::TempDir::new().unwrap();
let (logger, mut rx) = TamperLogger::channel();
logger.log(event("evt-1", None, vec![]));
logger.log(event("evt-2", Some("evt-1"), vec![]));
drop(logger);
run_tamper_writer(dir.path().to_path_buf(), &mut rx).await;
let lines = read_lines(dir.path());
assert_eq!(lines.len(), 2);
assert_eq!(lines[0]["tamper"]["event_id"], "evt-1");
assert_eq!(lines[1]["tamper"]["event_id"], "evt-2");
assert_eq!(lines[1]["tamper"]["related_event_id"], "evt-1");
assert_eq!(lines[0]["ocsf"]["type_uid"], 200402);
}
#[tokio::test]
async fn writer_omits_empty_deltas_and_absent_related_id() {
let dir = tempfile::TempDir::new().unwrap();
let (logger, mut rx) = TamperLogger::channel();
logger.log(event("evt-1", None, vec![]));
drop(logger);
run_tamper_writer(dir.path().to_path_buf(), &mut rx).await;
let line = &read_lines(dir.path())[0];
assert!(line["tamper"].get("field_deltas").is_none());
assert!(line["tamper"].get("related_event_id").is_none());
}
#[tokio::test]
async fn writer_records_delta_paths_and_change_kinds_only() {
let dir = tempfile::TempDir::new().unwrap();
let (logger, mut rx) = TamperLogger::channel();
logger.log(event(
"evt-1",
None,
vec![FieldDelta {
field: "hooks[0].url".to_string(),
change: "modified".to_string(),
}],
));
drop(logger);
run_tamper_writer(dir.path().to_path_buf(), &mut rx).await;
let deltas = &read_lines(dir.path())[0]["tamper"]["field_deltas"];
assert_eq!(deltas.as_array().unwrap().len(), 1);
assert_eq!(deltas[0]["field"], "hooks[0].url");
assert_eq!(deltas[0]["change"], "modified");
assert_eq!(deltas[0].as_object().unwrap().len(), 2, "no value fields");
}
#[tokio::test]
async fn writer_appends_to_an_existing_file_without_truncating() {
let dir = tempfile::TempDir::new().unwrap();
let path = dir.path().join("tamper.jsonl");
std::fs::write(&path, "{\"prior\":true}\n").unwrap();
let (logger, mut rx) = TamperLogger::channel();
logger.log(event("evt-1", None, vec![]));
drop(logger);
run_tamper_writer(dir.path().to_path_buf(), &mut rx).await;
let content = std::fs::read_to_string(&path).unwrap();
let mut lines = content.lines();
assert_eq!(lines.next(), Some("{\"prior\":true}"));
assert!(lines.next().unwrap().contains("evt-1"));
assert_eq!(lines.next(), None);
}
#[tokio::test]
async fn writer_creates_a_missing_log_directory() {
let dir = tempfile::TempDir::new().unwrap();
let nested = dir.path().join("a").join("b");
let (logger, mut rx) = TamperLogger::channel();
logger.log(event("evt-1", None, vec![]));
drop(logger);
run_tamper_writer(nested.clone(), &mut rx).await;
assert_eq!(read_lines(&nested).len(), 1);
}
#[tokio::test]
async fn writer_survives_an_unopenable_target_and_keeps_draining() {
let dir = tempfile::TempDir::new().unwrap();
std::fs::create_dir(dir.path().join("tamper.jsonl")).unwrap();
let (logger, mut rx) = TamperLogger::channel();
logger.log(event("evt-1", None, vec![]));
logger.log(event("evt-2", None, vec![]));
drop(logger);
run_tamper_writer(dir.path().to_path_buf(), &mut rx).await;
assert!(rx.recv().await.is_none(), "channel fully drained");
}
#[tokio::test]
async fn restarted_writer_resumes_on_the_same_channel_without_losing_events() {
let dir = tempfile::TempDir::new().unwrap();
let (logger, mut rx) = TamperLogger::channel();
logger.log(event("evt-1", None, vec![]));
let first = tokio::time::timeout(
std::time::Duration::from_millis(200),
run_tamper_writer(dir.path().to_path_buf(), &mut rx),
)
.await;
assert!(first.is_err(), "writer blocks while the sender is alive");
logger.log(event("evt-2", None, vec![]));
drop(logger);
run_tamper_writer(dir.path().to_path_buf(), &mut rx).await;
let ids: Vec<_> = read_lines(dir.path())
.iter()
.map(|l| l["tamper"]["event_id"].as_str().unwrap().to_string())
.collect();
assert_eq!(ids, ["evt-1", "evt-2"]);
}
#[tokio::test]
async fn log_drops_events_without_blocking_when_the_channel_is_full() {
let (logger, mut rx) = TamperLogger::channel();
for i in 0..300 {
logger.log(event(&format!("evt-{i}"), None, vec![]));
}
drop(logger);
let mut received = 0;
while rx.recv().await.is_some() {
received += 1;
}
assert_eq!(received, 256, "bounded at capacity, overflow dropped");
}
#[tokio::test]
async fn log_after_the_receiver_closes_does_not_panic() {
let (logger, rx) = TamperLogger::channel();
drop(rx);
logger.log(event("evt-1", None, vec![]));
}
#[tokio::test]
async fn new_spawns_a_writer_that_persists_events() {
let dir = tempfile::TempDir::new().unwrap();
let (logger, _handle) = TamperLogger::new(dir.path().to_path_buf());
logger.log(event("evt-1", None, vec![]));
let path = dir.path().join("tamper.jsonl");
let mut waited = 0;
while !path.exists() || std::fs::read_to_string(&path).unwrap().is_empty() {
assert!(waited < 100, "writer never persisted the event");
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
waited += 1;
}
assert_eq!(read_lines(dir.path())[0]["tamper"]["event_id"], "evt-1");
}
}