openlatch-client 0.6.1

OpenLatch runtime enforcement node — the capture-and-enforce adapter that evaluates every covered action against a coding agent's Autonomy Zone before it runs
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 {
    /// Wrap an already-spawned writer task, so the daemon can run the writer
    /// under its in-process task supervisor instead of a bare `tokio::spawn`.
    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)
    }

    /// Sender + receiver without spawning the writer — see
    /// [`crate::core::logging::EventLogger::channel`] for why the spawn lives
    /// with the daemon rather than here.
    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");
        }
    }
}

/// Drain tamper events onto `tamper.jsonl` until the channel closes.
///
/// `&mut` receiver so a supervisor can restart it on the same channel after a
/// panic without losing queued detections.
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();
        // A directory sitting where the file should be makes every open fail.
        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![]));

        // A supervisor restarts the writer on the same receiver after a panic;
        // emulate that by abandoning one run mid-stream and starting another.
        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");
    }
}