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;
struct Pending {
event: TamperEvent,
attempts: u32,
failure_reported: bool,
}
pub(crate) struct WiringTamper {
sinks: TamperSinks,
pending: BTreeMap<&'static str, Pending>,
gates: BTreeMap<&'static str, ReapplyGate>,
}
impl WiringTamper {
pub(crate) fn new(sinks: TamperSinks) -> Self {
Self {
sinks,
pending: BTreeMap::new(),
gates: BTreeMap::new(),
}
}
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 => {
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
}
pub(crate) fn is_pending(&self, agent: &str) -> bool {
self.pending.contains_key(agent)
}
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;
}
}
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,
_ => 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);
}
}
}
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
}
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))
}
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)
}
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");
}
}