openlatch-client 0.6.3

OpenLatch runtime enforcement node — the capture-and-enforce adapter that evaluates every covered action against a coding agent's Autonomy Zone before it runs
//! Relay wiring drift and relay CA trust, reported as tamper and — for the wiring — healed by the
//! one writer, the wiring supervisor.
//!
//! **The supervisor re-wires; nothing here writes an agent file.** This is the bookkeeping around
//! its pass: which agents drifted, one `tamper_detected` per distinct loss, the pacing of the
//! re-write, and the linked `tamper_healed` once the pass has put the wiring back. The hook
//! reconciler never touches wiring, so there is exactly one writer of it, as there always was.
//!
//! The CA is report-only. The daemon never installs or re-trusts it (an install can raise a
//! password dialog on a schedule nobody chose); a loss is reported once, today's release path
//! (`react_to_trust_loss_in`) still fails the agent toward direct connections, and the heal is
//! reported once `init` or `doctor --fix` has put the trust back and a proof has passed.

use std::collections::BTreeMap;
use std::path::{Path, PathBuf};

use crate::config::Config;
use crate::core::cloud::tamper::{
    FieldDelta, TamperEvent, COMPONENT_MODEL_RELAY_WIRING, COMPONENT_RELAY_CA,
    DETECTION_CA_UNTRUSTED, DETECTION_RELAY_MODIFIED, DETECTION_RELAY_UNWIRED,
    RELAY_WIRING_HOOK_EVENT,
};
use crate::hooks::binding::{AgentBinding, EndpointConvention, ProxyDelivery};
use crate::hooks::WiringDrift;
use crate::model_relay::preflight::WiringState;

use super::endpoint_wiring::ReapplyGate;
use super::reconciler::TamperSinks;

/// A detection still waiting on its heal.
struct Pending {
    event: TamperEvent,
    attempts: u32,
    /// A failed heal is reported once; the agent's later re-wire then reports the success.
    failure_reported: bool,
}

/// The wiring supervisor's tamper state, held for its whole life (a supervised restart builds a
/// fresh one, as it does the file watch).
pub(crate) struct WiringTamper {
    sinks: TamperSinks,
    pending: BTreeMap<&'static str, Pending>,
    /// Per agent: at most one re-write per 10 s, one per 60 s once contested — the provider
    /// slots' pacing, so something fighting the wiring is not matched write for write.
    gates: BTreeMap<&'static str, ReapplyGate>,
}

impl WiringTamper {
    pub(crate) fn new(sinks: TamperSinks) -> Self {
        Self {
            sinks,
            pending: BTreeMap::new(),
            gates: BTreeMap::new(),
        }
    }

    /// For an agent the supervisor believes is wired: whether its wiring left the disk and may be
    /// re-written now. A drift is reported on first sight and never again until it is healed;
    /// the answer is `false` while the pacing gate holds, so the next pass re-asks.
    ///
    /// `false` for an agent this instance does not own the wiring of: it never wrote the file,
    /// so nothing in it is ours to lose.
    pub(crate) fn should_rewire(
        &mut self,
        config: &Config,
        binding: &dyn AgentBinding,
        port: u16,
    ) -> bool {
        if !super::owns_wiring_for(config, binding) {
            return false;
        }
        let (method, deltas) = match drift(binding, port) {
            WiringDrift::Intact | WiringDrift::Released => {
                // Put back while a re-write was paced, or by whoever took it: recovered.
                self.settle(binding.agent_type(), true);
                return false;
            }
            WiringDrift::Unwired(deltas) => (DETECTION_RELAY_UNWIRED, deltas),
            WiringDrift::Modified(deltas) => (DETECTION_RELAY_MODIFIED, deltas),
        };
        let agent = binding.agent_type();
        if !self.pending.contains_key(agent) {
            let event = wiring_event(binding, method, deltas);
            tracing::warn!(
                agent,
                method,
                "the model relay wiring was changed on disk — reporting it and re-wiring"
            );
            self.sinks.publish_detected(&event);
            self.pending.insert(
                agent,
                Pending {
                    event,
                    attempts: 0,
                    failure_reported: false,
                },
            );
        }
        let now = std::time::Instant::now();
        let gate = self.gates.entry(agent).or_default();
        if !gate.allows(now) {
            return false;
        }
        gate.note(now);
        true
    }

    /// Whether a drift of `agent` is waiting on its heal. The supervisor's failed-probe arm
    /// reads it: while it is, the file holds what somebody else put there, and the unwire that
    /// arm would run takes the ledger's record — the real prior — with it.
    pub(crate) fn is_pending(&self, agent: &str) -> bool {
        self.pending.contains_key(agent)
    }

    /// After a pass's probe and write for `binding`: close its pending detection. Healed only
    /// when the pass left it wired AND the disk reads back intact.
    pub(crate) fn settle_after_write(
        &mut self,
        binding: &dyn AgentBinding,
        port: u16,
        wired: bool,
    ) {
        let agent = binding.agent_type();
        if !self.pending.contains_key(agent) {
            return;
        }
        let healed = wired && matches!(drift(binding, port), WiringDrift::Intact);
        self.settle(agent, healed);
    }

    fn settle(&mut self, agent: &'static str, healed: bool) {
        let circuit = match self.gates.get(agent) {
            Some(gate) if gate.contested(std::time::Instant::now()) => "open",
            _ => "closed",
        };
        let Some(p) = self.pending.get_mut(agent) else {
            return;
        };
        p.attempts += 1;
        if healed {
            self.sinks.publish_healed(&TamperEvent::new_healed(
                &p.event,
                "succeeded",
                p.attempts,
                circuit,
            ));
            tracing::info!(agent, "the model relay wiring was put back");
            self.pending.remove(agent);
        } else if !p.failure_reported {
            self.sinks.publish_healed(&TamperEvent::new_healed(
                &p.event, "failed", p.attempts, circuit,
            ));
            p.failure_reported = true;
        }
    }

    /// Item 5: for each wired `ProxyEnv` agent whose editor trusts only the OS store, ask the store
    /// whether it still trusts the relay CA. One `ca_untrusted` per loss, one `tamper_healed` once
    /// trust is back and the agent is wired again (its proof passed) — which, after a
    /// `doctor --fix` restart, a new process sees: the pending detection's id is kept beside the
    /// CA ([`CA_PENDING_FILE`]).
    ///
    /// Read-only toward the store and every agent file. Runs at the top of a tick, BEFORE
    /// `react_to_trust_loss_in` can release the agent (that release is today's behaviour, kept),
    /// so the loss is still seen on a wired agent.
    pub(crate) async fn watch_ca_trust(
        &self,
        config: &Config,
        agents: &[crate::hooks::DetectedAgent],
        ca: Option<&crate::model_relay::ca::CaInfo>,
        wiring: &WiringState,
    ) {
        let Some(ca) = ca else {
            return;
        };
        let dir = crate::model_relay::ca::ca_dir(&crate::config::openlatch_dir());
        let watched: Vec<&'static str> = agents
            .iter()
            .filter(|a| {
                a.binding
                    .model_relay_wiring()
                    .is_some_and(|w| w.endpoint.is_proxy_env())
                    && super::owns_wiring_for(config, &*a.binding)
                    && !super::env_trust_only(config, &*a.binding)
                    && wiring.is_wired(a.agent_type())
            })
            .map(|a| a.agent_type())
            .collect();
        if watched.is_empty() {
            return;
        }
        let mut pending = load_ca_pending(&dir);
        let store = crate::model_relay::trust_store::store();
        let sha = ca.sha256_hex.clone();
        let trusted = match tokio::task::spawn_blocking(move || store.is_trusted(&sha)).await {
            Ok(Ok(trusted)) => trusted,
            // A store that cannot answer proves neither a loss nor a heal.
            _ => return,
        };
        let mut changed = false;
        for agent in watched {
            match (trusted, pending.get(agent)) {
                (false, None) => {
                    let event = ca_event(agent, &dir);
                    tracing::warn!(
                        agent,
                        code = crate::error::ERR_MODEL_RELAY_CA_UNTRUSTED,
                        "the user trust store no longer trusts the model relay's certificate \
                         authority — run `openlatch doctor --fix` to trust it again"
                    );
                    self.sinks.publish_detected(&event);
                    pending.insert(agent.to_string(), event.tamper.event_id.clone());
                    changed = true;
                }
                (true, Some(id)) => {
                    let detected = ca_event(agent, &dir).with_event_id(id);
                    self.sinks.publish_healed(&TamperEvent::new_healed(
                        &detected,
                        "succeeded",
                        1,
                        "closed",
                    ));
                    pending.remove(agent);
                    changed = true;
                }
                _ => {}
            }
        }
        if changed {
            store_ca_pending(&dir, &pending);
        }
    }
}

/// The files the wiring of `agents` lives in, for the supervisor's file watch: a change to one
/// is re-read after the quiet period, not at the next tick.
pub(crate) fn watch_files(config: &Config, agents: &[crate::hooks::DetectedAgent]) -> Vec<PathBuf> {
    let mut files = Vec::new();
    for agent in agents {
        let binding = &*agent.binding;
        let Some(w) = binding.model_relay_wiring() else {
            continue;
        };
        if !super::owns_wiring_for(config, binding) {
            continue;
        }
        match w.endpoint {
            EndpointConvention::EnvVars { .. } | EndpointConvention::TomlProvider { .. } => {
                files.extend(crate::hooks::model_relay_config_path(binding));
            }
            EndpointConvention::ProxyEnv { delivery, .. } => {
                for d in delivery {
                    match d {
                        ProxyDelivery::SettingsKey { file, .. } => files.extend(file()),
                        ProxyDelivery::EnvFile { .. } => {
                            let ol = crate::config::openlatch_dir();
                            files.push(crate::hooks::proxy_env_file::env_sh_path(&ol));
                            files.push(crate::hooks::proxy_env_file::live_marker_path(&ol));
                        }
                    }
                }
                if let Some((home, documents)) = profile_home(delivery) {
                    files.extend(crate::hooks::shell_profile::watch_files(
                        &home,
                        documents.as_deref(),
                    ));
                }
            }
        }
    }
    files
}

/// Where the shell-profile line lives, when this install writes it for `delivery`
/// (`super::profile_line_wanted`, the writer's own gate).
fn profile_home(delivery: &[ProxyDelivery]) -> Option<(PathBuf, Option<PathBuf>)> {
    if !super::profile_line_wanted(delivery, crate::supervision::owns_machine_supervision()) {
        return None;
    }
    let home = dirs::home_dir()?;
    let documents = if cfg!(windows) {
        dirs::document_dir()
    } else {
        None
    };
    Some((home, documents))
}

/// [`crate::hooks::wiring_on_disk`], plus the shell-profile line where this install writes one.
fn drift(binding: &dyn AgentBinding, port: u16) -> WiringDrift {
    let found = crate::hooks::wiring_on_disk(binding, port);
    let Some(EndpointConvention::ProxyEnv { delivery, .. }) =
        binding.model_relay_wiring().map(|w| w.endpoint)
    else {
        return found;
    };
    let Some((home, documents)) = profile_home(delivery) else {
        return found;
    };
    if crate::hooks::shell_profile::is_present(&home, documents.as_deref()) {
        return found;
    }
    let lost = FieldDelta {
        field: "shell_profile".to_string(),
        change: "removed".to_string(),
    };
    match found {
        WiringDrift::Released => WiringDrift::Released,
        WiringDrift::Intact => WiringDrift::Unwired(vec![lost]),
        WiringDrift::Unwired(mut d) => {
            d.push(lost);
            WiringDrift::Unwired(d)
        }
        WiringDrift::Modified(mut d) => {
            d.push(lost);
            WiringDrift::Modified(d)
        }
    }
}

fn wiring_event(binding: &dyn AgentBinding, method: &str, deltas: Vec<FieldDelta>) -> TamperEvent {
    let agent = binding.agent_type();
    let path = crate::hooks::model_relay_config_path(binding).unwrap_or_default();
    TamperEvent::new(
        agent.to_string(),
        agent.to_string(),
        crate::core::hook_state::hash_settings_path(&path),
        RELAY_WIRING_HOOK_EVENT.to_string(),
        method.to_string(),
    )
    .with_component(COMPONENT_MODEL_RELAY_WIRING)
    .with_field_deltas(deltas)
}

fn ca_event(agent: &str, ca_dir: &Path) -> TamperEvent {
    TamperEvent::new(
        COMPONENT_RELAY_CA.to_string(),
        agent.to_string(),
        crate::core::hook_state::hash_settings_path(&crate::model_relay::ca::ca_pem_path(ca_dir)),
        RELAY_WIRING_HOOK_EVENT.to_string(),
        DETECTION_CA_UNTRUSTED.to_string(),
    )
    .with_component(COMPONENT_RELAY_CA)
}

/// Agent → the id of its `ca_untrusted` event still waiting on a heal. Beside the CA, so
/// `uninstall`'s removal of the CA directory takes it too.
const CA_PENDING_FILE: &str = "untrusted-events.json";

fn load_ca_pending(ca_dir: &Path) -> BTreeMap<String, String> {
    std::fs::read_to_string(ca_dir.join(CA_PENDING_FILE))
        .ok()
        .and_then(|raw| serde_json::from_str(&raw).ok())
        .unwrap_or_default()
}

fn store_ca_pending(ca_dir: &Path, pending: &BTreeMap<String, String>) {
    let path = ca_dir.join(CA_PENDING_FILE);
    let result = if pending.is_empty() {
        match std::fs::remove_file(&path) {
            Err(e) if e.kind() != std::io::ErrorKind::NotFound => Err(e),
            _ => Ok(()),
        }
    } else {
        serde_json::to_string(pending)
            .map_err(std::io::Error::other)
            .and_then(|body| crate::core::fs_secure::write_readable(&path, &body))
    };
    if let Err(e) = result {
        tracing::warn!(error = %e, "could not record the relay CA's trust state — a heal after a restart is not linked to its detection");
    }
}