openlatch-client 0.6.0

OpenLatch runtime enforcement node — the capture-and-enforce adapter that evaluates every covered action against a coding agent's Autonomy Zone before it runs
//! A CLI process's last telemetry events, handed to the running daemon.
//!
//! A command captures `command_invoked` on its way out of `main`, after the last
//! batch tick it will ever see, so nothing sends it. Sending it from the CLI would
//! put a PostHog POST — a fresh TLS handshake — on the exit path of every command.
//! Instead the process gives its unsent events to the local daemon over the
//! loopback admin API it already talks to, and they ride the daemon's next batch
//! with the properties they were captured with.
//!
//! Best effort by construction. No daemon, no token, a refused connection, a
//! daemon whose telemetry is off: each drops the events (I4), and none of it is
//! reported anywhere (I10). The whole exchange is bounded by [`DRAIN_WAIT`] plus
//! [`TOTAL_TIMEOUT`], far below a delay a person notices.
//!
//! The body is the client's own contract with itself — a CLI and the daemon it
//! hands to — like every other `/admin/*` body, not a platform wire type.

use std::path::Path;
use std::time::Duration;

use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};

use super::client::{QueuedEvent, TelemetryHandle};

/// The daemon route that accepts a hand-off, behind the `/admin/*` bearer auth.
pub const PATH: &str = "/admin/telemetry/events";

/// Most events one hand-off may carry. A CLI process holds at most one batch's
/// worth at exit; the headroom is for a command that captured in a burst.
pub const MAX_EVENTS: usize = 50;

/// Largest hand-off body, in bytes. An event with its super-properties is well
/// under 1 KiB.
pub const MAX_BYTES: usize = 64 * 1024;

/// Longest event name accepted, matching every name `events.rs` builds.
const MAX_NAME_LEN: usize = 64;

/// How long the exiting process waits for its batch task to give its buffer back.
const DRAIN_WAIT: Duration = Duration::from_millis(25);

/// Loopback connect deadline. A daemon that is up accepts in microseconds.
const CONNECT_TIMEOUT: Duration = Duration::from_millis(50);

/// The whole request, connect included.
const TOTAL_TIMEOUT: Duration = Duration::from_millis(100);

/// The hand-off body: telemetry events exactly as their process captured them.
#[derive(Debug, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct HandoffBatch {
    pub events: Vec<HandoffEvent>,
}

/// One captured event, super-properties already merged.
#[derive(Debug, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct HandoffEvent {
    pub event: String,
    pub properties: Map<String, Value>,
}

impl HandoffBatch {
    /// Whether this is a batch of telemetry events within the caps — the only
    /// thing the daemon relays. The caller has already bounded the body at
    /// [`MAX_BYTES`].
    ///
    /// A telemetry event is a name `events.rs` could have built (`snake_case`,
    /// or `$`-prefixed like `$create_alias`) and the properties every capture
    /// merges in, `distinct_id` and `agent_id` among them.
    pub fn validate(&self) -> Result<(), &'static str> {
        if self.events.is_empty() || self.events.len() > MAX_EVENTS {
            return Err("a hand-off carries between 1 and MAX_EVENTS events");
        }
        for event in &self.events {
            if !is_event_name(&event.event) {
                return Err("an event name is not a telemetry event name");
            }
            for key in ["distinct_id", "agent_id"] {
                if !event
                    .properties
                    .get(key)
                    .and_then(Value::as_str)
                    .is_some_and(|v| !v.is_empty())
                {
                    return Err("an event is missing its distinct_id or agent_id");
                }
            }
        }
        Ok(())
    }

    /// The events, in the shape the batch loop POSTs.
    pub fn into_queued(self) -> impl Iterator<Item = QueuedEvent> {
        self.events.into_iter().map(|e| QueuedEvent {
            name: e.event,
            properties: e.properties,
        })
    }
}

fn is_event_name(name: &str) -> bool {
    let bare = name.strip_prefix('$').unwrap_or(name);
    name.len() <= MAX_NAME_LEN
        && bare.starts_with(|c: char| c.is_ascii_lowercase())
        && bare
            .chars()
            .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '_')
}

/// Give this process's unsent events to the daemon that owns `openlatch_dir`.
///
/// Returns once the daemon answered, or the bound ran out, or there was
/// nothing to hand off. Never fails and never reports: the events are a metric,
/// and the process is exiting.
pub fn hand_off(handle: &TelemetryHandle, openlatch_dir: &Path) {
    // The daemon writes its PID file when it starts and removes it when it
    // stops, so without one there is nobody to hand to. Checking is a stat;
    // not checking costs a connect to a closed port, which Windows retries
    // until `CONNECT_TIMEOUT` rather than refusing at once.
    if !openlatch_dir.join("daemon.pid").exists() {
        return;
    }
    let Ok(cfg) = crate::config::Config::load(None, None, false) else {
        return;
    };
    let events = handle.take_unsent(DRAIN_WAIT);
    if !events.is_empty() {
        send(events, openlatch_dir, cfg.port);
    }
}

/// POST `events` to the daemon. `true` when the daemon took them.
fn send(events: Vec<QueuedEvent>, openlatch_dir: &Path, port: u16) -> bool {
    let Ok(token) = std::fs::read_to_string(openlatch_dir.join("daemon.token")) else {
        return false;
    };
    let token = token.trim();
    if token.is_empty() {
        return false;
    }
    let batch = HandoffBatch {
        events: events
            .into_iter()
            .take(MAX_EVENTS)
            .map(|e| HandoffEvent {
                event: e.name,
                properties: e.properties,
            })
            .collect(),
    };
    let Ok(body) = serde_json::to_vec(&batch) else {
        return false;
    };
    if body.len() > MAX_BYTES {
        return false;
    }
    let agent = ureq::Agent::new_with_config(
        ureq::Agent::config_builder()
            .timeout_connect(Some(CONNECT_TIMEOUT))
            .timeout_global(Some(TOTAL_TIMEOUT))
            .http_status_as_error(false)
            // SECURITY: the destination is the literal loopback daemon. Honouring
            // ALL_PROXY / HTTP_PROXY would send the daemon token to a proxy, and a
            // redirect would take the request off 127.0.0.1 — as for the hook.
            .proxy(None)
            .max_redirects(0)
            .build(),
    );
    agent
        .post(format!("http://127.0.0.1:{port}{PATH}"))
        .header("Authorization", format!("Bearer {token}"))
        .header("Content-Type", "application/json")
        .send(&body[..])
        .is_ok_and(|resp| resp.status().is_success())
}

#[cfg(test)]
mod tests {
    use super::*;
    use serde_json::json;

    fn event(name: &str) -> HandoffEvent {
        let mut properties = Map::new();
        properties.insert("distinct_id".into(), json!("agt_a"));
        properties.insert("agent_id".into(), json!("agt_a"));
        properties.insert("command".into(), json!("status"));
        HandoffEvent {
            event: name.into(),
            properties,
        }
    }

    #[test]
    fn a_batch_of_captured_events_is_accepted() {
        let batch = HandoffBatch {
            events: vec![event("command_invoked"), event("$create_alias")],
        };
        assert_eq!(batch.validate(), Ok(()));
    }

    #[test]
    fn anything_but_a_telemetry_event_shape_is_refused() {
        let empty = HandoffBatch { events: vec![] };
        assert!(empty.validate().is_err(), "an empty hand-off");

        let too_many = HandoffBatch {
            events: (0..=MAX_EVENTS).map(|_| event("command_invoked")).collect(),
        };
        assert!(too_many.validate().is_err(), "more than MAX_EVENTS");

        for name in [
            "",
            "Command",
            "command-invoked",
            "$",
            "a b",
            &"x".repeat(65),
        ] {
            let batch = HandoffBatch {
                events: vec![event(name)],
            };
            assert!(batch.validate().is_err(), "name {name:?}");
        }

        for key in ["distinct_id", "agent_id"] {
            let mut e = event("command_invoked");
            e.properties.remove(key);
            let batch = HandoffBatch { events: vec![e] };
            assert!(batch.validate().is_err(), "without {key}");
        }
    }

    #[test]
    fn unknown_fields_are_refused_at_parse() {
        let extra_top = json!({ "events": [], "api_key": "phc_x" });
        assert!(serde_json::from_value::<HandoffBatch>(extra_top).is_err());
        let extra_event = json!({ "events": [{ "event": "x", "properties": {}, "timestamp": 1 }] });
        assert!(serde_json::from_value::<HandoffBatch>(extra_event).is_err());
    }

    /// The drop path: no daemon on the port. The CLI must move on inside the
    /// bound, having sent nothing and said nothing.
    #[test]
    fn no_daemon_drops_the_events_inside_the_bound() {
        let dir = tempfile::tempdir().expect("tempdir");
        std::fs::write(dir.path().join("daemon.token"), "t0ken").expect("token");
        // Bound, then released: nothing listens on it now.
        let port = std::net::TcpListener::bind("127.0.0.1:0")
            .expect("bind")
            .local_addr()
            .expect("addr")
            .port();
        let events = vec![QueuedEvent {
            name: "command_invoked".into(),
            properties: event("command_invoked").properties,
        }];
        let started = std::time::Instant::now();
        assert!(!send(events, dir.path(), port));
        assert!(
            started.elapsed() < TOTAL_TIMEOUT * 3,
            "a refused hand-off took {:?}",
            started.elapsed()
        );
    }

    /// No PID file is the common "no daemon" case: the hand-off returns at
    /// once, without draining the batch or touching the network. The batch
    /// still holding its event is the proof: every path past the PID check
    /// drains it first. No wall-clock bound — on a loaded runner a stat can
    /// outlast any tight one, and time is not what this test is about.
    #[test]
    fn no_pid_file_hands_nothing_off() {
        use super::super::client::{start, ClientConfig};
        use super::super::consent::{ConsentState, DecidedBy, Resolved};

        let dir = tempfile::tempdir().expect("tempdir");
        let handle = start(ClientConfig {
            resolved: Resolved {
                state: ConsentState::Enabled,
                decided_by: DecidedBy::ConfigFile,
            },
            super_props: super::super::super_props::SuperProps::new("agt_a".into(), false),
            debug_stderr: false,
            baked_key_present: true,
            egress: crate::egress::EgressConfig::direct(),
            host: "http://127.0.0.1:9".into(),
        });
        handle.capture(super::super::events::Event::command_invoked(
            "status", None, 0, 1,
        ));

        hand_off(&handle, dir.path());
        assert_eq!(
            handle.take_unsent(Duration::from_secs(2)).len(),
            1,
            "the batch was never drained"
        );
    }

    /// No token means no daemon this process could talk to: nothing is sent.
    #[test]
    fn no_token_sends_nothing() {
        let dir = tempfile::tempdir().expect("tempdir");
        let events = vec![QueuedEvent {
            name: "command_invoked".into(),
            properties: event("command_invoked").properties,
        }];
        assert!(!send(events, dir.path(), 9));
    }
}