pub mod admin;
pub mod auth;
#[cfg(feature = "model-relay")]
pub(crate) mod cline_hub;
#[cfg(feature = "model-relay")]
pub(crate) use cline_hub::release_hub as release_cline_hub_after_death;
pub mod cline_prompt_backfill;
pub mod config_monitor;
pub mod dedup;
#[cfg(feature = "model-relay")]
pub(crate) mod endpoint_wiring;
pub mod fallback_replay;
pub mod handlers;
pub mod hold_queue;
pub mod identity;
pub mod optimize;
pub mod policy_poller;
pub mod reconciler;
pub mod reinforcement;
pub mod session_state;
pub mod spend;
pub mod watcher;
use std::path::Path;
use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU32, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use axum::{extract::DefaultBodyLimit, middleware, routing::get, routing::post, Router};
use secrecy::SecretString;
use tokio::net::TcpListener;
use crate::cloud::{CloudState, CredentialProvider};
use crate::config::Config;
use crate::core::logging::tamper_log::{TamperLogger, TamperLoggerHandle};
use crate::core::supervision::task::{spawn_supervised, HealthRegistry, RestartPolicy, TaskSpec};
use crate::logging::{EventLogger, EventLoggerHandle};
use crate::privacy::PrivacyFilter;
use crate::update;
mod subsystem {
pub const MODEL_RELAY: &str = "model_relay";
pub const MODEL_RELAY_WIRING: &str = "model-relay-wiring";
pub const CLOUD_WORKER: &str = "cloud-worker";
pub const OUTBOX_DRAIN: &str = "outbox-drain";
pub const ALERTS_LONG_POLL: &str = "alerts-long-poll";
pub const POLICY_POLLER: &str = "policy-poller";
pub const RECONCILER: &str = "reconciler";
pub const DEDUP_EVICTOR: &str = "dedup-evictor";
pub const LOG_CLEANUP: &str = "log-cleanup";
pub const EVENT_LOG_WRITER: &str = "event-log-writer";
pub const TAMPER_LOG_WRITER: &str = "tamper-log-writer";
pub const FALLBACK_REPLAY: &str = "fallback-replay";
pub const UPDATE_CHECK: &str = "update-check";
pub const AUTO_UPDATE_WORKER: &str = "auto-update-worker";
pub const UPDATE_SENTINEL: &str = "update-sentinel";
pub const EGRESS_MONITOR: &str = "egress-monitor";
pub const CONFIG_MONITOR: &str = "config-monitor";
}
struct CredentialStoreAdapter {
store: Arc<dyn crate::auth::CredentialStore>,
}
impl CredentialProvider for CredentialStoreAdapter {
fn retrieve(&self) -> Option<SecretString> {
match self.store.retrieve() {
Ok(key) => Some(key),
Err(e) => {
tracing::debug!(
code = e.code,
"credential provider: retrieve failed — returning None to worker"
);
None
}
}
}
fn invalidate(&self) {
self.store.invalidate();
}
}
pub struct PolicyRuntime {
pub handle: crate::core::policy::PolicyHandle,
pub last_fetch_ok: Arc<AtomicBool>,
pub last_poll_ok_at: Arc<AtomicI64>,
pub floor: Arc<std::sync::RwLock<Option<ClientFloorState>>>,
pub sessions: Arc<session_state::SessionStateManager>,
pub dispatch_sessions: Arc<session_state::DispatchStateManager>,
pub holds: Arc<hold_queue::HoldQueue>,
pub reinforcement: Arc<reinforcement::ReinforcementStore>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ClientFloorState {
pub installed_version: String,
pub minimum_version: String,
}
impl PolicyRuntime {
pub fn load_from_disk(base_dir: &Path) -> Self {
use crate::core::policy::{project_document, store};
let cached = match store::load(base_dir) {
Ok(cached) => cached,
Err(e) => {
tracing::warn!(
target: "policy",
code = e.code(),
error = %e,
"cached policy bundle rejected at startup; running with no policy until the next successful fetch"
);
None
}
};
let last_fetch_ok = Arc::new(AtomicBool::new(true));
let last_poll_ok_at = Arc::new(AtomicI64::new(0));
let floor = Arc::new(std::sync::RwLock::new(None));
let resident = cached.and_then(|c| {
last_fetch_ok.store(c.meta.last_fetch_ok, Ordering::Relaxed);
if let Some(secs) = c
.meta
.last_poll_ok_at
.as_deref()
.and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok())
.map(|dt| dt.timestamp())
{
last_poll_ok_at.store(secs, Ordering::Relaxed);
}
if c.meta.state.as_deref() == Some("floor_blocked") {
if let Some(minimum_version) = c.meta.min_client_version.clone() {
if let Ok(mut guard) = floor.write() {
*guard = Some(ClientFloorState {
installed_version: env!("CARGO_PKG_VERSION").to_string(),
minimum_version,
});
}
}
}
match project_document(c.document) {
Ok((_, mut resident)) => {
resident.digest = Some(c.meta.digest);
Some(resident)
}
Err(error) => {
store::discard(base_dir);
tracing::warn!(
target: "policy",
code = crate::error::ERR_BUNDLE_INVALID,
%error,
"cached policy bundle projection failed at startup; running with no policy until the next successful fetch"
);
None
}
}
});
Self {
handle: crate::core::policy::new_handle(resident),
last_fetch_ok,
last_poll_ok_at,
floor,
sessions: Arc::new(session_state::SessionStateManager::default()),
dispatch_sessions: Arc::new(session_state::DispatchStateManager::default()),
holds: Arc::new(hold_queue::HoldQueue::default()),
reinforcement: Arc::new(reinforcement::ReinforcementStore::default()),
}
}
fn log_startup_state(&self) {
match self.handle.load().as_ref() {
Some(b) => tracing::info!(
target: "policy",
revision = b.revision,
rules = b.command_rules.len(),
request_rules = b.request_rules.len(),
enforcement_enabled = b.enforcement_enabled,
organization_id = %b.organization_id,
"policy engine enabled; enforcing the cached bundle"
),
None => tracing::info!(
target: "policy",
"policy engine enabled; no bundle — allowing everything until the first successful fetch"
),
}
}
}
pub struct AppState {
pub config: Arc<Config>,
pub token: String,
pub dedup: dedup::DedupStore,
pub cline_prompts: cline_prompt_backfill::PromptBackfill,
pub event_logger: EventLogger,
pub privacy_filter: PrivacyFilter,
pub event_counter: AtomicU64,
pub shutdown_tx: tokio::sync::Mutex<Option<tokio::sync::oneshot::Sender<()>>>,
pub started_at: std::time::Instant,
pub available_update: Mutex<Option<String>>,
pub cloud_tx: Option<tokio::sync::mpsc::Sender<crate::cloud::CloudEvent>>,
pub cloud_state: Option<CloudState>,
pub credential_provider: Option<Arc<dyn CredentialProvider>>,
pub local_ipv4: Option<std::net::Ipv4Addr>,
pub local_ipv6: Option<std::net::Ipv6Addr>,
pub public_ipv4: Option<std::net::Ipv4Addr>,
pub public_ipv6: Option<std::net::Ipv6Addr>,
pub tamper_logger: Option<TamperLogger>,
pub outbox: Option<Arc<crate::cloud::outbox::Outbox>>,
pub update_in_progress: Arc<AtomicBool>,
pub update_status: Arc<Mutex<update::UpdateStatusSnapshot>>,
pub update_progress: Arc<update::ApplyProgress>,
pub admin_shutdown_request: Arc<tokio::sync::Notify>,
pub restart_into: std::sync::OnceLock<std::path::PathBuf>,
pub last_hook_at_unix_secs: Arc<AtomicU64>,
pub hooks_in_flight: Arc<AtomicU32>,
pub content_hash_cache: Arc<config_monitor::ContentHashCache>,
pub config_monitor: Arc<config_monitor::ConfigMonitorRuntime>,
pub pending_alerts: Arc<config_monitor::PendingAlerts>,
pub policy: Option<PolicyRuntime>,
pub registry: Arc<crate::model_relay::session::SessionRegistry>,
pub health: Arc<HealthRegistry>,
pub egress: crate::egress::EgressState,
}
impl AppState {
pub fn set_available_update(&self, version: String) {
if let Ok(mut guard) = self.available_update.lock() {
*guard = Some(version);
}
}
pub fn get_available_update(&self) -> Option<String> {
self.available_update.lock().ok().and_then(|g| g.clone())
}
}
pub async fn start_server(
config: Config,
token: String,
credential_store: Option<Arc<dyn crate::auth::CredentialStore>>,
spawn_model_relay: bool,
) -> anyhow::Result<Served> {
let bind_host = match std::env::var("OPENLATCH_BIND_ALL").as_deref() {
Ok("true") | Ok("1") => "0.0.0.0",
_ => "127.0.0.1",
};
let bind_addr = format!("{}:{}", bind_host, config.port);
let listener = TcpListener::bind(&bind_addr).await?;
tracing::info!(
port = config.port,
addr = %bind_addr,
"daemon listening"
);
if let Err(e) = crate::config::write_port_file(config.port) {
tracing::warn!(error = %e, "failed to write daemon.port file");
}
serve_with_listener(
listener,
config,
token,
credential_store,
spawn_model_relay,
true,
)
.await
}
pub async fn start_server_with_listener(
listener: TcpListener,
config: Config,
token: String,
credential_store: Option<Arc<dyn crate::auth::CredentialStore>>,
spawn_model_relay: bool,
) -> anyhow::Result<Served> {
serve_with_listener(
listener,
config,
token,
credential_store,
spawn_model_relay,
false,
)
.await
}
async fn serve_with_listener(
listener: TcpListener,
config: Config,
token: String,
credential_store: Option<Arc<dyn crate::auth::CredentialStore>>,
spawn_model_relay: bool,
reconcile_wiring: bool,
) -> anyhow::Result<Served> {
let health = Arc::new(HealthRegistry::new());
let (tasks_shutdown_tx, tasks_shutdown_rx) = tokio::sync::watch::channel(false);
let (_logs_shutdown_tx, logs_shutdown_rx) = tokio::sync::watch::channel(false);
let log_dir = config.log_dir.clone();
let (event_logger, event_log_rx) = EventLogger::channel();
let event_log_rx = Arc::new(tokio::sync::Mutex::new(event_log_rx));
let logger_handle = {
let log_dir = log_dir.clone();
EventLoggerHandle::from_task(spawn_supervised(
&health,
TaskSpec::new(subsystem::EVENT_LOG_WRITER, RestartPolicy::OnFailure),
logs_shutdown_rx.clone(),
move || {
let rx = event_log_rx.clone();
let log_dir = log_dir.clone();
async move {
let mut rx = rx.lock_owned().await;
crate::logging::run_event_writer(log_dir, &mut rx).await
}
},
))
};
let privacy_filter = PrivacyFilter::new(&config.extra_patterns);
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let pending_alerts = Arc::new(config_monitor::PendingAlerts::new());
let policy_runtime = if config.policy.enabled {
let runtime = PolicyRuntime::load_from_disk(&crate::config::openlatch_dir());
runtime.log_startup_state();
Some(runtime)
} else {
tracing::info!(
target: "policy",
"policy engine disabled by config ([policy] enabled = false); no bundle is fetched and no resident bundle is consulted"
);
None
};
let egress_state = crate::egress::EgressState::new(&config.egress);
for warning in egress_state.warnings() {
tracing::warn!(target: "egress", "{warning}");
}
let cloud_timeouts = crate::egress::Timeouts {
connect: Some(std::time::Duration::from_millis(
config.cloud.timeout_connect_ms,
)),
total: Some(std::time::Duration::from_millis(
config.cloud.timeout_total_ms,
)),
};
let poll_timeouts = crate::egress::Timeouts::total(std::time::Duration::from_millis(
config.cloud.timeout_total_ms,
));
let egress_clients = Arc::new(crate::egress::EgressClients::new(
cloud_timeouts,
poll_timeouts,
));
if let Err(e) = egress_clients.apply(&config.egress) {
tracing::error!(
target: "egress",
code = %e.code,
error = %e.message,
"no egress client could be built; outbound traffic is refused until the route is fixed"
);
egress_state.record_failure(e.code, e.message.clone());
egress_state.record_failure(e.code, e.message);
}
crate::egress::emit_proxy_shape_if_changed(
&crate::config::openlatch_dir(),
&egress_state.snapshot(),
0,
);
{
let monitor_state = egress_state.clone();
let monitor_shutdown = tasks_shutdown_rx.clone();
let heal = crate::egress::SelfHeal {
resolver: None,
clients: egress_clients.clone(),
base: Arc::new(config.egress.clone()),
api_url: config.cloud.api_url.clone(),
openlatch_dir: crate::config::openlatch_dir(),
};
spawn_supervised(
&health,
TaskSpec::new(subsystem::EGRESS_MONITOR, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
move || {
crate::egress::run_egress_monitor(
monitor_state.clone(),
monitor_shutdown.clone(),
Some(heal.clone()),
)
},
);
}
let registry = Arc::new(crate::model_relay::session::SessionRegistry::default());
let source_formats: crate::cloud::worker::SourceFormats = crate::hooks::detect_agents()
.iter()
.filter_map(|a| {
a.binding
.model_relay_wiring()
.map(|w| (a.binding.agent_type(), w.wire_format))
})
.collect();
#[allow(clippy::type_complexity)]
let (cloud_tx, cloud_state_opt, outbox_opt, cloud_worker_task, credential_provider_opt): (
_,
_,
_,
Option<tokio::task::JoinHandle<()>>,
Option<Arc<dyn CredentialProvider>>,
) = {
if let Some(store) = credential_store {
let (tx, rx) = tokio::sync::mpsc::channel(config.cloud.channel_size);
let cloud_state = CloudState::new();
let host_key = identity::host_key(config.agent_id.as_deref().unwrap_or_default())
.filter(|key| {
let usable = reqwest::header::HeaderValue::from_str(key).is_ok();
if !usable {
tracing::warn!(
"host key is not a usable header value — this host will not be \
identified to the platform; check agent_id in config.toml"
);
}
usable
});
let cloud_config = crate::cloud::CloudConfig {
api_url: config.cloud.api_url.clone(),
timeout_connect_ms: config.cloud.timeout_connect_ms,
timeout_total_ms: config.cloud.timeout_total_ms,
retry_delay_ms: config.cloud.retry_delay_ms,
channel_size: config.cloud.channel_size,
rate_limit_default_secs: 30,
credential_poll_interval_ms: config.cloud.credential_poll_interval_ms,
fallback_max_bytes: config.cloud.fallback_max_bytes,
batch_max_events: config.cloud.batch_max_events,
batch_max_wait_ms: config.cloud.batch_max_wait_ms,
host_key: host_key.clone(),
host_id: identity::host_id(),
agent_id: config.agent_id.clone(),
};
let openlatch_dir = crate::config::openlatch_dir();
let provider: Arc<dyn CredentialProvider> = Arc::new(CredentialStoreAdapter { store });
let worker_state = cloud_state.clone();
let alerts_provider = provider.clone();
let alerts_state = cloud_state.clone();
let policy_provider = provider.clone();
let policy_state = cloud_state.clone();
let policy_dir = openlatch_dir.clone();
let admin_provider = provider.clone();
let worker_egress =
crate::egress::EgressReporter::recording(&config.egress, egress_state.clone());
let outbox = if config.cloud.outbox_enabled {
Some(Arc::new(crate::cloud::outbox::Outbox::new(
&openlatch_dir,
config.cloud.outbox_max_bytes,
)))
} else {
None
};
let cloud_rx = Arc::new(tokio::sync::Mutex::new(rx));
let worker_outbox = outbox.clone();
let worker_sessions = registry.clone();
let worker_client = egress_clients.cloud.clone();
let drain_provider = provider.clone();
let drain_config = cloud_config.clone();
let drain_state = cloud_state.clone();
let drain_dir = openlatch_dir.clone();
let drain_formats = source_formats.clone();
let drain_sessions = registry.clone();
let drain_egress = worker_egress.clone();
let drain_client = worker_client.clone();
let drain_shutdown = tasks_shutdown_rx.clone();
let cloud_worker = spawn_supervised(
&health,
TaskSpec::new(subsystem::CLOUD_WORKER, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
{
let shutdown_rx = tasks_shutdown_rx.clone();
let worker_formats = source_formats.clone();
move || {
let rx = cloud_rx.clone();
let provider = provider.clone();
let cloud_config = cloud_config.clone();
let worker_egress = worker_egress.clone();
let worker_state = worker_state.clone();
let openlatch_dir = openlatch_dir.clone();
let outbox = worker_outbox.clone();
let shutdown_rx = shutdown_rx.clone();
let worker_client = worker_client.clone();
let source_formats = worker_formats.clone();
let sessions = worker_sessions.clone();
async move {
let mut rx = rx.lock_owned().await;
crate::cloud::worker::run_cloud_worker_on(
&mut rx,
provider,
cloud_config,
worker_egress,
worker_state,
openlatch_dir,
outbox,
Some(shutdown_rx),
worker_client,
source_formats,
sessions,
)
.await
}
}
},
);
tracing::info!(
api_url = %config.cloud.api_url,
channel_size = config.cloud.channel_size,
"cloud forwarding worker started"
);
if let Some(drain_outbox) = outbox.clone() {
spawn_supervised(
&health,
TaskSpec::new(subsystem::OUTBOX_DRAIN, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
move || {
crate::cloud::worker::run_outbox_drain_on(
drain_outbox.clone(),
drain_client.clone(),
drain_config.clone(),
drain_provider.clone(),
drain_state.clone(),
drain_dir.clone(),
drain_formats.clone(),
drain_sessions.clone(),
drain_egress.clone(),
drain_shutdown.clone(),
#[cfg(test)]
None,
)
},
);
}
let alerts_url = config.cloud.api_url.clone();
let alerts_target = pending_alerts.clone();
let alerts_egress =
crate::egress::EgressReporter::recording(&config.egress, egress_state.clone());
let alerts_client = egress_clients.alerts.if_routed();
match (alerts_client, config.agent_id.clone()) {
(Some(client), Some(machine_id)) => {
spawn_supervised(
&health,
TaskSpec::new(subsystem::ALERTS_LONG_POLL, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
move || {
config_monitor::run_long_poll(
alerts_target.clone(),
alerts_state.clone(),
alerts_provider.clone(),
alerts_url.clone(),
machine_id.clone(),
client.clone(),
alerts_egress.clone(),
)
},
);
}
(Some(_), None) => {
tracing::warn!(
"alerts long-poll skipped: agent_id missing from config.toml; run `openlatch init` to provision"
);
}
_ => {}
}
if let Some(policy) = policy_runtime.as_ref() {
match egress_clients.poller.if_routed() {
Some(client) => {
let policy_handle = policy.handle.clone();
let policy_last_fetch_ok = policy.last_fetch_ok.clone();
let policy_last_poll_ok_at = policy.last_poll_ok_at.clone();
let policy_floor = policy.floor.clone();
let policy_api_url = config.cloud.api_url.clone();
let policy_cfg = config.policy.clone();
policy
.holds
.set_source(Arc::new(hold_queue::CloudHoldAnswerSource::new(
&policy_api_url,
policy_provider.clone(),
client.clone(),
)));
let policy_agent_id = config.agent_id.clone();
let policy_host_key = host_key.clone();
let policy_egress = crate::egress::EgressReporter::recording(
&config.egress,
egress_state.clone(),
);
spawn_supervised(
&health,
TaskSpec::new(subsystem::POLICY_POLLER, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
move || {
policy_poller::run_policy_poller(
policy_handle.clone(),
policy_last_fetch_ok.clone(),
policy_last_poll_ok_at.clone(),
policy_floor.clone(),
policy_state.clone(),
policy_provider.clone(),
policy_api_url.clone(),
policy_cfg.clone(),
policy_dir.clone(),
client.clone(),
policy_agent_id.clone(),
policy_host_key.clone(),
policy_egress.clone(),
)
},
);
tracing::info!(
target: "policy",
poll_interval_secs = config.policy.poll_interval_secs,
"policy bundle poller started"
);
}
None => {
tracing::warn!(
target: "policy",
code = crate::error::ERR_BUNDLE_FETCH_FAILED,
"no egress route is permitted, so the policy poller is not started; the resident bundle keeps enforcing but will not refresh"
);
}
}
}
(
Some(tx),
Some(cloud_state),
outbox,
Some(cloud_worker),
Some(admin_provider),
)
} else {
tracing::info!(
"no credential store provided — cloud forwarding cannot authenticate and is idle"
);
(None, None, None, None, None)
}
};
if policy_runtime.is_some() && cloud_state_opt.is_none() {
tracing::warn!(
target: "policy",
"policy is enabled but no credential provider is available; running disk-bundle-only — the resident bundle keeps enforcing and will not refresh"
);
}
let startup_started = std::time::Instant::now();
let host_ips = crate::net::HostIps::detect().await;
tracing::info!(
local_ipv4 = host_ips
.local_ipv4
.map(|a| a.to_string())
.as_deref()
.unwrap_or("none"),
local_ipv6 = host_ips
.local_ipv6
.map(|a| a.to_string())
.as_deref()
.unwrap_or("none"),
public_ipv4 = host_ips
.public_ipv4
.map(|a| a.to_string())
.as_deref()
.unwrap_or("none"),
public_ipv6 = host_ips
.public_ipv6
.map(|a| a.to_string())
.as_deref()
.unwrap_or("none"),
"host ips detected"
);
tokio::task::spawn_blocking(crate::daemon::identity::os_user);
let openlatch_dir_for_tamper = crate::config::openlatch_dir();
let (tamper_logger, tamper_rx) = TamperLogger::channel();
let tamper_rx = Arc::new(tokio::sync::Mutex::new(tamper_rx));
let _tamper_logger_handle: TamperLoggerHandle =
TamperLoggerHandle::from_task(spawn_supervised(
&health,
TaskSpec::new(subsystem::TAMPER_LOG_WRITER, RestartPolicy::OnFailure),
logs_shutdown_rx.clone(),
move || {
let rx = tamper_rx.clone();
let dir = openlatch_dir_for_tamper.clone();
async move {
let mut rx = rx.lock_owned().await;
crate::core::logging::tamper_log::run_tamper_writer(dir, &mut rx).await
}
},
));
let cache_max = config.inventory_monitor.cache_max_entries.max(64);
let content_hash_cache = Arc::new(config_monitor::ContentHashCache::new(cache_max));
let (config_monitor, config_monitor_request_rx, config_monitor_health) =
if config.inventory_monitor.enabled {
let (runtime, request_rx) = config_monitor::ConfigMonitorRuntime::pending();
let task_health = health.register(subsystem::CONFIG_MONITOR, RestartPolicy::Always);
(runtime, Some(request_rx), Some(task_health))
} else {
(config_monitor::ConfigMonitorRuntime::disabled(), None, None)
};
if !config.inventory_monitor.enabled {
tracing::info!("config monitor disabled by config");
} else {
tracing::info!("config monitor initialization pending");
}
#[cfg(feature = "model-relay")]
#[allow(clippy::type_complexity)]
let (model_relay_task, wiring_task, endpoint_listeners): (
Option<tokio::task::JoinHandle<()>>,
Option<tokio::task::JoinHandle<()>>,
Option<Arc<crate::model_relay::endpoints::EndpointListeners>>,
) = if spawn_model_relay {
use crate::model_relay;
use crate::model_relay::wire_format::WireFormat;
let model_relay_port = config.model_relay.port;
let first_listener = model_relay::bind_pinned(model_relay_port).await?;
tracing::info!(
port = model_relay_port,
"model relay listener bound (loopback only)"
);
let wiring = Arc::new(model_relay::preflight::WiringState::default());
let model_relay_patterns = config.extra_patterns.clone();
let model_relay_transforms_act = config.model_relay.transforms_act;
let model_relay_registry = registry.clone();
let model_relay_cloud_tx = cloud_tx.clone();
let model_relay_shutdown_rx = tasks_shutdown_rx.clone();
let pre_bound = Arc::new(tokio::sync::Mutex::new(Some(first_listener)));
let model_relay_policy = policy_runtime.as_ref().map(|p| p.handle.clone());
let model_relay_policy_sessions = policy_runtime
.as_ref()
.map(|p| p.sessions.clone())
.unwrap_or_else(|| Arc::new(session_state::SessionStateManager::default()));
let model_relay_wiring = wiring.clone();
let model_relay_upstream = reqwest::Url::parse(
&config
.model_relay
.upstream_for(WireFormat::AnthropicMessages),
)
.unwrap_or_else(|_| model_relay::default_upstream());
let model_relay_upstream_map: std::collections::BTreeMap<String, String> = WireFormat::ALL
.iter()
.map(|f| (f.as_str().to_string(), config.model_relay.upstream_for(*f)))
.collect();
let model_relay_explicit: std::collections::BTreeSet<&'static str> = WireFormat::ALL
.iter()
.map(|f| f.as_str())
.filter(|k| config.model_relay.upstream.contains_key(*k))
.collect();
let model_relay_egress = config.egress.clone();
let model_relay_client = egress_clients.model_relay.clone();
let model_relay_upstreams: crate::model_relay::UpstreamHandle =
Arc::new(std::sync::RwLock::new(std::collections::BTreeMap::new()));
let model_relay_upstreams_for_serve = model_relay_upstreams.clone();
let wiring_upstreams = model_relay_upstreams.clone();
let model_relay_reporter =
crate::egress::EgressReporter::recording(&config.egress, egress_state.clone());
let model_relay_install_id = config.agent_id.clone();
let endpoint_listeners = {
let upstream = model_relay_upstream.clone();
let patterns = model_relay_patterns.clone();
let egress = model_relay_egress.clone();
let registry = model_relay_registry.clone();
let cloud_tx = model_relay_cloud_tx.clone();
let policy = model_relay_policy.clone();
let policy_sessions = model_relay_policy_sessions.clone();
let wiring = model_relay_wiring.clone();
let client = model_relay_client.clone();
let reporter = model_relay_reporter.clone();
let trackers = model_relay::RelayTrackers::new(model_relay::DEFAULT_INFLIGHT);
let transforms_act = model_relay_transforms_act;
let install_id = model_relay_install_id.clone();
let factory: model_relay::endpoints::StateFactory = Arc::new(move |endpoint| {
model_relay::ModelRelayState::new_with_egress(
upstream.clone(),
endpoint.port(),
model_relay::DEFAULT_INFLIGHT,
&patterns,
&egress,
)
.with_trackers(&trackers)
.with_measurement(registry.clone(), cloud_tx.clone())
.with_transforms_act(transforms_act)
.with_policy(policy.clone())
.with_policy_sessions(policy_sessions.clone())
.with_wiring(wiring.clone())
.with_client_handle(client.clone())
.with_egress_reporter(reporter.clone())
.with_install_id(install_id.clone())
});
Arc::new(model_relay::endpoints::EndpointListeners::new(factory))
};
let endpoint_listeners_for_serve = endpoint_listeners.clone();
let endpoint_wiring = Arc::new(endpoint_wiring::EndpointWiring::new(
endpoint_listeners.clone(),
wiring.clone(),
model_relay::endpoints::RelayPorts {
daemon: config.port,
main: model_relay_port,
},
reconciler::TamperSinks {
logger: Some(tamper_logger.clone()),
cloud_tx: cloud_tx.clone(),
agent_id: config.agent_id.clone().unwrap_or_default(),
client_version: env!("OPENLATCH_VERSION").to_string(),
},
model_relay_registry.clone(),
));
let serve_task = spawn_supervised(
&health,
TaskSpec::new(subsystem::MODEL_RELAY, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
move || {
let bstate = Arc::new(
model_relay::ModelRelayState::new_with_egress(
model_relay_upstream.clone(),
model_relay_port,
model_relay::DEFAULT_INFLIGHT,
&model_relay_patterns,
&model_relay_egress,
)
.with_upstream_handle(model_relay_upstreams_for_serve.clone())
.with_upstream_map(model_relay_upstream_map.clone())
.with_measurement(model_relay_registry.clone(), model_relay_cloud_tx.clone())
.with_transforms_act(model_relay_transforms_act)
.with_policy(model_relay_policy.clone())
.with_policy_sessions(model_relay_policy_sessions.clone())
.with_wiring(model_relay_wiring.clone())
.with_client_handle(model_relay_client.clone())
.with_egress_reporter(model_relay_reporter.clone())
.with_install_id(model_relay_install_id.clone())
.with_explicit_upstreams(model_relay_explicit.clone())
.with_endpoints(endpoint_listeners_for_serve.clone())
.with_intercept(model_relay_wiring.intercept())
.with_refusals(model_relay_wiring.refusals()),
);
apply_learned_upstreams(
&model_relay_upstreams_for_serve,
&model_relay_upstream_map,
);
model_relay::serve_attempt(
pre_bound.clone(),
bstate,
model_relay_shutdown_rx.clone(),
)
},
);
let wiring_config = Arc::new(config.clone());
let wiring_state = wiring.clone();
let wiring_shutdown = tasks_shutdown_rx.clone();
let wiring_formats = source_formats.clone();
let wiring_endpoints = endpoint_wiring.clone();
let wiring_task = spawn_supervised(
&health,
TaskSpec::new(subsystem::MODEL_RELAY_WIRING, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
move || {
run_wiring_supervisor(
wiring_config.clone(),
model_relay_port,
wiring_state.clone(),
wiring_formats.clone(),
wiring_upstreams.clone(),
Some(wiring_endpoints.clone()),
wiring_shutdown.clone(),
)
},
);
(
Some(serve_task),
Some(wiring_task),
Some(endpoint_listeners),
)
} else {
if reconcile_wiring {
unwire_every_agent(&config);
let cfg = config.clone();
let _ = tokio::task::spawn_blocking(move || cline_hub::release_hub(&cfg)).await;
release_every_provider_endpoint(&config);
}
(None, None, None)
};
let state = Arc::new(AppState {
config: Arc::new(config.clone()),
token,
dedup: dedup::DedupStore::new(),
cline_prompts: cline_prompt_backfill::PromptBackfill::default(),
event_logger,
privacy_filter,
event_counter: AtomicU64::new(0),
shutdown_tx: tokio::sync::Mutex::new(Some(shutdown_tx)),
started_at: std::time::Instant::now(),
available_update: Mutex::new(None),
cloud_tx,
cloud_state: cloud_state_opt,
credential_provider: credential_provider_opt,
local_ipv4: host_ips.local_ipv4,
local_ipv6: host_ips.local_ipv6,
public_ipv4: host_ips.public_ipv4,
public_ipv6: host_ips.public_ipv6,
tamper_logger: Some(tamper_logger),
outbox: outbox_opt,
update_in_progress: Arc::new(AtomicBool::new(false)),
update_status: Arc::new(Mutex::new(update::UpdateStatusSnapshot::idle())),
update_progress: Arc::default(),
admin_shutdown_request: Arc::new(tokio::sync::Notify::new()),
restart_into: std::sync::OnceLock::new(),
last_hook_at_unix_secs: Arc::new(AtomicU64::new(0)),
hooks_in_flight: Arc::new(AtomicU32::new(0)),
content_hash_cache,
config_monitor: config_monitor.clone(),
pending_alerts,
policy: policy_runtime,
registry,
health: health.clone(),
egress: egress_state,
});
crate::telemetry::mark_served_as_daemon();
crate::telemetry::capture_global(crate::telemetry::Event::daemon_started(
state.config.port,
startup_started
.elapsed()
.as_millis()
.min(u128::from(u64::MAX)) as u64,
state.config.cloud.enabled,
));
{
let state_for_replay = state.clone();
spawn_supervised(
&health,
TaskSpec::new(subsystem::FALLBACK_REPLAY, RestartPolicy::OnFailure),
tasks_shutdown_rx.clone(),
move || fallback_replay::run(state_for_replay.clone()),
);
}
if config.update.check {
let current = env!("CARGO_PKG_VERSION").to_string();
let update_egress = config.egress.clone();
let update_registry = config.update.registry_origin.clone();
let state_for_update = state.clone();
spawn_supervised(
&health,
TaskSpec::new(subsystem::UPDATE_CHECK, RestartPolicy::OnFailure),
tasks_shutdown_rx.clone(),
move || {
let current = current.clone();
let update_egress = update_egress.clone();
let update_registry = update_registry.clone();
let state_for_update = state_for_update.clone();
async move {
if let Some(latest) =
update::check_for_update(¤t, &update_registry, &update_egress).await
{
tracing::warn!(code = crate::error::ERR_VERSION_OUTDATED, latest_version = %latest, "Update available: run `openlatch update`");
state_for_update.set_available_update(latest);
}
}
},
);
}
if config.update.auto_update {
let state_for_worker = state.clone();
spawn_supervised(
&health,
TaskSpec::new(subsystem::AUTO_UPDATE_WORKER, RestartPolicy::OnFailure),
tasks_shutdown_rx.clone(),
move || run_auto_update_worker(state_for_worker.clone()),
);
} else {
tracing::info!(target: "update", "auto-update worker disabled by config");
}
if let Some(sentinel) = update::read_sentinel() {
let port = config.port;
if sentinel.to == env!("CARGO_PKG_VERSION") {
*state.update_status.lock().expect("status mutex poisoned") =
update::UpdateStatusSnapshot::restarted(&sentinel);
}
let update_status = state.update_status.clone();
spawn_supervised(
&health,
TaskSpec::new(subsystem::UPDATE_SENTINEL, RestartPolicy::OnFailure),
tasks_shutdown_rx.clone(),
move || {
let sentinel = sentinel.clone();
let update_status = update_status.clone();
async move {
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
let healthy = probe_until_healthy(
|| probe_self_health(port),
&POST_RESTART_PROBE_BACKOFF,
)
.await;
settle_restarted_update(&update_status, &sentinel, healthy);
if !healthy {
tracing::warn!(target: "update", from = %sentinel.from, to = %sentinel.to, "post-restart /health probe failed; leaving sentinel + .bak in place for restart-loop rollback");
return;
}
tracing::info!(target: "update", from = %sentinel.from, to = %sentinel.to, "post-restart healthz probe succeeded; cleaning up");
if let Err(e) = update::cleanup_bak_files() {
tracing::warn!(target: "update", error = %e, "cleanup of .bak siblings failed (non-fatal)");
}
if let Err(e) = update::delete_sentinel() {
tracing::warn!(target: "update", error = %e, "delete of update sentinel failed (non-fatal)");
}
crate::install_state::InstallState::stamp_for_running_binary(env!(
"CARGO_PKG_VERSION"
));
}
},
);
}
let state_for_evict = state.clone();
spawn_supervised::<_, _, ()>(
&health,
TaskSpec::new(subsystem::DEDUP_EVICTOR, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
move || {
let state = state_for_evict.clone();
async move {
let mut interval = tokio::time::interval(std::time::Duration::from_secs(30));
loop {
interval.tick().await;
state.dedup.evict_expired();
}
}
},
);
let state_for_cleanup = state.clone();
spawn_supervised::<_, _, ()>(
&health,
TaskSpec::new(subsystem::LOG_CLEANUP, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
move || {
let state = state_for_cleanup.clone();
async move {
let mut interval = tokio::time::interval(std::time::Duration::from_secs(86_400));
loop {
interval.tick().await;
let log_dir = state.config.log_dir.clone();
let retention = state.config.retention_days;
match tokio::task::spawn_blocking(move || {
crate::logging::cleanup_old_logs(&log_dir, retention)
})
.await
{
Ok(Ok(deleted)) if deleted > 0 => {
tracing::info!(
deleted,
retention_days = retention,
"cleaned up old log files"
);
}
Ok(Ok(_)) => {}
Ok(Err(e)) => {
tracing::warn!(error = %e, "periodic log cleanup failed");
}
Err(e) => {
tracing::warn!(error = %e, "periodic log cleanup task join error");
}
}
}
}
},
);
let _watcher_guards;
let _poll_handle;
let detected = crate::hooks::detect_agents();
if !detected.is_empty() {
let targets: Vec<reconciler::AgentTarget> = detected
.iter()
.map(reconciler::AgentTarget::for_agent)
.collect();
let openlatch_dir = crate::config::openlatch_dir();
let token_file = openlatch_dir.join("daemon.token");
reconciler::run_startup_reconcile(targets.clone(), &openlatch_dir, config.port);
let (reconcile_tx, reconcile_rx) = tokio::sync::mpsc::channel(100);
let sinks = reconciler::TamperSinks {
logger: state.tamper_logger.clone(),
cloud_tx: state.cloud_tx.clone(),
agent_id: state.config.agent_id.clone().unwrap_or_default(),
client_version: env!("OPENLATCH_VERSION").to_string(),
};
let r = Arc::new(tokio::sync::Mutex::new(
reconciler::Reconciler::new_with_sinks(
reconcile_rx,
targets.clone(),
openlatch_dir,
config.port,
token_file,
sinks,
),
));
spawn_supervised(
&health,
TaskSpec::new(subsystem::RECONCILER, RestartPolicy::Always),
tasks_shutdown_rx.clone(),
move || {
let r = r.clone();
async move {
let mut guard = r.lock_owned().await;
guard.run().await
}
},
);
_watcher_guards = targets
.iter()
.filter_map(|target| {
match watcher::spawn_watcher(&target.settings_path, reconcile_tx.clone()) {
Ok(w) => {
tracing::info!(
agent = target.binding.agent_type(),
"filesystem watcher active"
);
Some(w)
}
Err(e) => {
tracing::warn!(
error = %e,
agent = target.binding.agent_type(),
"filesystem watcher failed — falling back to poll-only"
);
None
}
}
})
.collect();
_poll_handle = Some(watcher::spawn_poll_fallback(reconcile_tx));
tracing::info!(
agents = targets.len(),
"reconciler started (reactive watcher + 30s poll)"
);
} else {
_watcher_guards = Vec::new();
_poll_handle = None;
}
let hook_routes = Router::new()
.route("/hooks", post(handlers::ingest_cloudevent))
.route("/holds/{tool_use_id}", get(handlers::wait_for_hold))
.route_layer(middleware::from_fn_with_state(
state.clone(),
auth::bearer_auth,
));
let shutdown_route = Router::new()
.route("/shutdown", post(handlers::shutdown_handler))
.route_layer(middleware::from_fn_with_state(
state.clone(),
auth::bearer_auth,
));
let public_routes = Router::new()
.route("/health", get(handlers::health))
.route("/metrics", get(handlers::metrics));
let admin_routes = admin::router(state.clone());
let app = Router::new()
.merge(hook_routes)
.merge(shutdown_route)
.merge(public_routes)
.merge(admin_routes)
.layer(DefaultBodyLimit::max(1_048_576))
.with_state(state.clone());
let admin_shutdown_request = state.admin_shutdown_request.clone();
let (server_ready_tx, server_ready_rx) = tokio::sync::oneshot::channel();
let server_task = tokio::spawn(async move {
let _ = server_ready_tx.send(());
axum::serve(listener, app)
.with_graceful_shutdown(async move {
tokio::select! {
_ = signal_handler() => {
tracing::info!("received OS shutdown signal");
}
_ = shutdown_rx => {
tracing::info!("received shutdown via /shutdown endpoint");
}
_ = admin_shutdown_request.notified() => {
tracing::info!(target: "update", "received shutdown for in-flight auto-update");
}
}
})
.await
});
server_ready_rx
.await
.map_err(|_| anyhow::anyhow!("daemon server task exited before startup"))?;
let config_monitor_task = match (config_monitor_request_rx, config_monitor_health) {
(Some(request_rx), Some(monitor_health)) => {
let monitor_runtime = config_monitor.clone();
let runtime_for_init = monitor_runtime.clone();
let monitor_filter = state.privacy_filter.clone();
let filter_for_init = monitor_filter.clone();
let monitor_cache = state.content_hash_cache.clone();
let monitor_cloud_tx = state.cloud_tx.clone();
let monitor_event_logger = state.event_logger.clone();
let monitor_config = state.config.clone();
let monitor_shutdown = tasks_shutdown_rx.clone();
Some(tokio::spawn(config_monitor::run_startup(
monitor_runtime,
monitor_health,
monitor_filter,
monitor_shutdown,
move || async move {
let manifest =
tokio::task::spawn_blocking(config_monitor::manifest::load_embedded)
.await
.map_err(|error| {
format!("manifest initialization task failed: {error}")
})?
.map_err(|error| error.to_string())?;
runtime_for_init.mark_manifest_loaded();
let request_tx = runtime_for_init.request_tx().ok_or_else(|| {
"monitor request channel closed during startup".to_string()
})?;
config_monitor::ConfigMonitor::new(
Arc::new(manifest),
monitor_cache,
filter_for_init,
monitor_cloud_tx,
monitor_event_logger,
monitor_config,
)
.spawn_with_requests(request_tx, request_rx)
.await
.map_err(|error| error.to_string())
},
|handle, shutdown| handle.run_until_shutdown(shutdown),
)))
}
_ => None,
};
let server_result = match server_task.await {
Ok(result) => result.map_err(anyhow::Error::from),
Err(error) => Err(anyhow::anyhow!("daemon server task failed: {error}")),
};
let _ = tasks_shutdown_tx.send(true);
if let Some(config_monitor_task) = config_monitor_task {
if tokio::time::timeout(std::time::Duration::from_secs(5), config_monitor_task)
.await
.is_err()
{
tracing::warn!("config monitor did not stop within 5s — abandoning it");
}
}
if let Err(error) = server_result {
if state.restart_into.get().is_none() {
return Err(error);
}
tracing::error!(
target: "update",
%error,
"serve loop failed while draining for an update; handing over regardless"
);
}
if let Some(cloud_worker) = cloud_worker_task {
if tokio::time::timeout(std::time::Duration::from_secs(5), cloud_worker)
.await
.is_err()
{
tracing::warn!(
"cloud worker did not finish its shutdown flush within 5s — abandoning it"
);
}
}
#[cfg(feature = "model-relay")]
if let Some(model_relay_task) = model_relay_task {
if let Some(wiring_task) = wiring_task {
if tokio::time::timeout(std::time::Duration::from_secs(5), wiring_task)
.await
.is_err()
{
tracing::warn!(
"model relay wiring supervisor did not stop within 5s — abandoning it"
);
}
}
if tokio::time::timeout(std::time::Duration::from_secs(5), model_relay_task)
.await
.is_err()
{
tracing::warn!("model relay listener did not shut down within 5s — abandoning it");
}
if let Some(listeners) = endpoint_listeners {
listeners.shutdown_all().await;
}
unwire_every_agent(&state.config);
let cfg = state.config.clone();
let _ = tokio::task::spawn_blocking(move || cline_hub::release_hub(&cfg)).await;
}
let restart_into = state.restart_into.get().cloned();
let uptime_secs = state.started_at.elapsed().as_secs();
let events = state
.event_counter
.load(std::sync::atomic::Ordering::Relaxed);
crate::logging::daemon_log::log_shutdown(uptime_secs, events);
crate::telemetry::capture_global(crate::telemetry::Event::daemon_stopped(uptime_secs, events));
match Arc::try_unwrap(state) {
Ok(_state) => {
if tokio::time::timeout(std::time::Duration::from_secs(5), logger_handle.shutdown())
.await
.is_err()
{
tracing::warn!("event-log drain did not finish within 5s — abandoning it");
}
}
Err(arc) => {
tracing::warn!(
strong_refs = Arc::strong_count(&arc),
"AppState still has references at shutdown — final event-log batch may be dropped"
);
drop(arc);
}
}
if let Some(exe) = &restart_into {
crate::telemetry::flush_global().await;
tracing::info!(target: "update", exe = %exe.display(), "drained; handing over to the updated binary");
}
Ok(Served {
uptime_secs,
events,
restart_into,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Served {
pub uptime_secs: u64,
pub events: u64,
pub restart_into: Option<std::path::PathBuf>,
}
pub fn format_uptime(secs: u64) -> String {
let hours = secs / 3600;
let minutes = (secs % 3600) / 60;
let seconds = secs % 60;
if hours > 0 {
format!("{}h{}m", hours, minutes)
} else if minutes > 0 {
format!("{}m{}s", minutes, seconds)
} else {
format!("{}s", seconds)
}
}
async fn probe_self_health(port: u16) -> bool {
let Ok(client) = crate::egress::client_builder()
.timeout(std::time::Duration::from_secs(2))
.build()
else {
return false;
};
let url = format!("http://127.0.0.1:{port}/health");
match client.get(&url).send().await {
Ok(r) => r.status().is_success(),
Err(_) => false,
}
}
const POST_RESTART_PROBE_BACKOFF: [std::time::Duration; 4] = [
std::time::Duration::from_secs(1),
std::time::Duration::from_secs(2),
std::time::Duration::from_secs(4),
std::time::Duration::from_secs(8),
];
async fn probe_until_healthy<F, Fut>(mut probe: F, backoff: &[std::time::Duration]) -> bool
where
F: FnMut() -> Fut,
Fut: std::future::Future<Output = bool>,
{
if probe().await {
return true;
}
for pause in backoff {
tokio::time::sleep(*pause).await;
if probe().await {
return true;
}
}
false
}
fn settle_restarted_update(
status: &Mutex<update::UpdateStatusSnapshot>,
sentinel: &update::UpdateSentinel,
healthy: bool,
) {
let mut snap = status.lock().expect("status mutex poisoned");
let still_restarted = snap.status == update::UpdateStatusKind::InProgress
&& snap.stage == Some(update::ApplyStage::Healthz)
&& snap.started_at.as_deref() == Some(sentinel.applied_at.as_str());
if !still_restarted {
return;
}
if healthy {
snap.status = update::UpdateStatusKind::Completed;
snap.stage = None;
} else {
snap.status = update::UpdateStatusKind::Failed;
snap.error = Some("post-restart /health probe failed".to_string());
}
snap.ended_at = Some(crate::install_state::now_rfc3339());
}
async fn run_auto_update_worker(state: Arc<AppState>) {
use crate::install_state::{detect_install_method, InstallMethod};
if crate::telemetry::is_ci_environment() {
tracing::debug!(target: "update", "CI environment detected; auto-update worker disabled");
return;
}
if matches!(detect_install_method(), InstallMethod::CargoInstall) {
tracing::info!(target: "update", "cargo-install path detected; auto-update worker disabled — use `cargo install --force --locked openlatch-client`");
return;
}
tokio::time::sleep(std::time::Duration::from_secs(10)).await;
let normal_interval =
std::time::Duration::from_secs(state.config.update.check_interval_secs.max(1));
let critical_interval = std::time::Duration::from_secs(3600);
let defer_interval = std::time::Duration::from_secs(300);
let mut pending_since: Option<std::time::Instant> = None;
let current_version = env!("CARGO_PKG_VERSION").to_string();
loop {
let next_sleep = match worker_iteration(&state, ¤t_version, pending_since).await {
WorkerOutcome::Idle => {
pending_since = None;
normal_interval
}
WorkerOutcome::Deferred {
severity: update::Severity::Critical,
} => {
if pending_since.is_none() {
pending_since = Some(std::time::Instant::now());
}
critical_interval
}
WorkerOutcome::Deferred { .. } => {
if pending_since.is_none() {
pending_since = Some(std::time::Instant::now());
}
defer_interval
}
WorkerOutcome::Failed {
severity: update::Severity::Critical,
} => {
pending_since = None;
critical_interval
}
WorkerOutcome::Failed { .. } => {
pending_since = None;
normal_interval
}
};
tokio::time::sleep(next_sleep).await;
}
}
#[derive(Debug, Clone, Copy)]
enum WorkerOutcome {
Idle,
Deferred { severity: update::Severity },
Failed { severity: update::Severity },
}
async fn worker_iteration(
state: &Arc<AppState>,
current_version: &str,
pending_since: Option<std::time::Instant>,
) -> WorkerOutcome {
let registry_origin = state.config.update.registry_origin.clone();
let download_timeout =
std::time::Duration::from_secs(state.config.update.download_timeout_secs.max(1));
let result = update::check(current_version, ®istry_origin, &state.config.egress).await;
let mut install = crate::install_state::InstallState::load_or_default();
install.record_check();
if let Err(e) = install.save() {
tracing::warn!(target: "update", error = %e.message, "failed to write install-state.json");
}
let (latest, severity) = match result {
update::CheckResult::Available {
latest, severity, ..
} => (latest, severity),
update::CheckResult::UpToDate { .. } | update::CheckResult::Failed { .. } => {
return WorkerOutcome::Idle;
}
};
let pending_age = pending_since
.map(|t| std::time::Instant::now().saturating_duration_since(t))
.unwrap_or_default();
if !update::should_apply_now(
severity,
&state.last_hook_at_unix_secs,
&state.hooks_in_flight,
pending_age,
state.config.update.quiet_window_secs,
state.config.update.max_defer_secs,
) {
tracing::debug!(target: "update", latest = %latest, severity = %severity.as_str(), "deferring update — agent active or quiet window not met");
return WorkerOutcome::Deferred { severity };
}
if state
.update_in_progress
.compare_exchange(
false,
true,
std::sync::atomic::Ordering::AcqRel,
std::sync::atomic::Ordering::Acquire,
)
.is_err()
{
tracing::info!(target: "update", "auto-update worker yielding to in-flight manual apply");
return WorkerOutcome::Deferred { severity };
}
state.update_progress.begin();
*state.update_status.lock().expect("status mutex poisoned") =
update::UpdateStatusSnapshot::in_progress(current_version, &latest);
let opts = update::ApplyOptions {
current_version: current_version.to_string(),
registry_origin,
download_timeout,
force_cargo_install: false,
mode: update::ApplyMode::Rpc,
egress: state.config.egress.clone(),
progress: state.update_progress.clone(),
};
admin::run_apply_in_daemon(state.clone(), opts, severity).await;
WorkerOutcome::Failed { severity }
}
#[cfg(feature = "model-relay")]
pub(crate) fn owns_wiring_for(
config: &Config,
binding: &dyn crate::hooks::binding::AgentBinding,
) -> bool {
let on_default_port = config.model_relay.port == crate::model_relay::default_model_relay_port();
match config.model_relay.own_agent_wiring {
Some(false) => false,
None => on_default_port,
Some(true) => on_default_port || !binding.config_is_machine_global(),
}
}
#[cfg(feature = "model-relay")]
pub(crate) fn env_trust_only(
config: &Config,
binding: &dyn crate::hooks::binding::AgentBinding,
) -> bool {
env_trust_only_from(
owns_wiring_for(config, binding),
crate::supervision::owns_machine_supervision(),
config.model_relay.port == crate::model_relay::default_model_relay_port(),
)
}
#[cfg(feature = "model-relay")]
fn env_trust_only_from(owns_wiring: bool, owns_machine: bool, on_default_port: bool) -> bool {
owns_wiring && !owns_machine && !on_default_port
}
#[cfg(feature = "model-relay")]
const WIRING_TICK: std::time::Duration = std::time::Duration::from_secs(60);
#[cfg(feature = "model-relay")]
const WIRING_BACKOFF_MAX: std::time::Duration = std::time::Duration::from_secs(300);
#[cfg(feature = "model-relay")]
fn wiring_delay(base: std::time::Duration, consecutive_failures: u32) -> std::time::Duration {
if consecutive_failures == 0 {
return base;
}
let shift = (consecutive_failures - 1).min(8);
base.saturating_mul(1u32 << shift).min(WIRING_BACKOFF_MAX)
}
#[cfg(feature = "model-relay")]
fn wiring_tick() -> std::time::Duration {
match std::env::var("OPENLATCH_MODEL_RELAY_WIRING_TICK_MS") {
Ok(v) => match v.parse::<u64>() {
Ok(ms) => std::time::Duration::from_millis(ms.max(50)),
Err(_) => WIRING_TICK,
},
Err(_) => WIRING_TICK,
}
}
#[cfg(feature = "model-relay")]
fn apply_learned_upstreams(
upstreams: &crate::model_relay::UpstreamHandle,
configured: &std::collections::BTreeMap<String, String>,
) {
for fmt in crate::model_relay::wire_format::WireFormat::ALL {
let recorded =
crate::hooks::model_relay_endpoints::peek(&crate::hooks::upstream_record_key(fmt));
let learned = match recorded {
Some(Some(recorded)) => match reqwest::Url::parse(&recorded) {
Ok(url) => Some(url),
Err(_) => {
tracing::warn!(
format = fmt.as_str(),
value = %recorded,
"recorded upstream is not a valid URL — leaving this format on its default"
);
continue;
}
},
_ => None,
};
let mut map = upstreams.write().unwrap_or_else(|e| e.into_inner());
let target = match learned {
Some(url) => url,
None => {
let Some(url) = configured
.get(fmt.as_str())
.and_then(|v| reqwest::Url::parse(v).ok())
else {
continue;
};
if map.get(fmt.as_str()).is_none_or(|current| *current == url) {
continue;
}
url
}
};
if map.get(fmt.as_str()) == Some(&target) {
continue;
}
map.insert(fmt.as_str(), target.clone());
tracing::info!(
format = fmt.as_str(),
upstream = %target,
"model relay upstream for {} now {}",
fmt.as_str(),
target
);
}
}
#[cfg(feature = "model-relay")]
#[derive(Debug, Clone, Copy)]
struct TrustClock {
wall: time::OffsetDateTime,
mono: std::time::Instant,
}
#[cfg(feature = "model-relay")]
type ProxyEnvDeclarations =
std::collections::BTreeMap<&'static str, std::collections::BTreeSet<String>>;
#[cfg(all(test, feature = "model-relay"))]
fn react_to_trust_loss(
config: &Config,
agents: &[crate::hooks::DetectedAgent],
ca: Option<&crate::model_relay::ca::CaInfo>,
at: TrustClock,
refusals: &crate::model_relay::intercept::RefusalLedger,
holds: &mut std::collections::BTreeMap<String, std::time::Instant>,
wiring: &crate::model_relay::preflight::WiringState,
) -> std::collections::BTreeSet<&'static str> {
let declared = proxy_env_declarations(agents);
let proxy_env = proxy_env_agents(agents, &declared);
react_to_trust_loss_in(config, &proxy_env, ca, at, refusals, holds, wiring)
}
#[cfg(feature = "model-relay")]
fn proxy_env_agents<'a>(
agents: &'a [crate::hooks::DetectedAgent],
declared: &'a ProxyEnvDeclarations,
) -> Vec<(
&'a crate::hooks::DetectedAgent,
&'a std::collections::BTreeSet<String>,
)> {
agents
.iter()
.filter_map(|a| declared.get(a.agent_type()).map(|hosts| (a, hosts)))
.collect()
}
#[cfg(feature = "model-relay")]
fn react_to_trust_loss_in(
config: &Config,
proxy_env: &[(
&crate::hooks::DetectedAgent,
&std::collections::BTreeSet<String>,
)],
ca: Option<&crate::model_relay::ca::CaInfo>,
at: TrustClock,
refusals: &crate::model_relay::intercept::RefusalLedger,
holds: &mut std::collections::BTreeMap<String, std::time::Instant>, wiring: &crate::model_relay::preflight::WiringState,
) -> std::collections::BTreeSet<&'static str> {
use crate::model_relay::{ca as ca_mod, ca_lifecycle as life, preflight::Verdict};
let (now, clock) = (at.wall, at.mono);
let mut skip = std::collections::BTreeSet::new();
if proxy_env.is_empty() {
return skip;
}
if let Some(info) = ca {
if life::days_left(info, now) <= ca_mod::CA_RELEASE_DAYS {
let reason = life::expiring_reason(info.not_after);
for (a, _) in proxy_env {
if wiring.is_wired(a.agent_type())
|| wiring.verdict(a.agent_type()) != Verdict::Failed(reason.clone())
{
release_for_trust(config, a, reason.clone(), wiring);
}
skip.insert(a.agent_type());
}
return skip;
}
}
let mut fresh = Vec::new();
for host in refusals.tripped_hosts() {
if !holds.contains_key(&host) {
holds.insert(host.clone(), clock);
fresh.push(host);
}
}
for (a, hosts) in proxy_env {
let Some(host) = hosts.iter().find(|h| holds.contains_key(h.as_str())) else {
continue;
};
if fresh.iter().any(|f| f == host) || wiring.is_wired(a.agent_type()) {
release_for_trust(config, a, life::refused_reason(host), wiring);
}
skip.insert(a.agent_type());
}
skip
}
#[cfg(feature = "model-relay")]
fn release_for_trust(
config: &Config,
a: &crate::hooks::DetectedAgent,
reason: String,
wiring: &crate::model_relay::preflight::WiringState,
) {
use crate::model_relay::preflight::Verdict;
if owns_wiring_for(config, &*a.binding) {
release_proxy_delivery(
&*a.binding,
crate::hooks::model_relay_endpoints::ReleasedBy::Wiring,
);
}
tracing::warn!(
agent = a.agent_type(),
reason = %reason,
"model relay released the agent's proxy delivery — its calls go direct"
);
wiring.set_wired(a.agent_type(), false);
wiring.set_verdict(a.agent_type(), Verdict::Failed(reason)); }
async fn run_wiring_supervisor(
config: Arc<Config>,
port: u16,
wiring: Arc<crate::model_relay::preflight::WiringState>,
source_formats: crate::cloud::worker::SourceFormats,
upstreams: crate::model_relay::UpstreamHandle,
endpoints: Option<Arc<endpoint_wiring::EndpointWiring>>,
mut shutdown: tokio::sync::watch::Receiver<bool>,
) {
use crate::model_relay::proxy::upstream_failures;
let mut last_failures = upstream_failures();
let mut consecutive_failures: u32 = 0;
let tick = wiring_tick();
let configured_upstreams: std::collections::BTreeMap<String, String> =
crate::model_relay::wire_format::WireFormat::ALL
.iter()
.map(|f| (f.as_str().to_string(), config.model_relay.upstream_for(*f)))
.collect();
let file_changed = Arc::new(tokio::sync::Notify::new());
let mut file_trigger: Option<watcher::FileTrigger> = None;
let mut holds: std::collections::BTreeMap<String, std::time::Instant> =
std::collections::BTreeMap::new();
loop {
if *shutdown.borrow() {
return;
}
let failures = upstream_failures();
let forwarding_broke = failures > last_failures;
last_failures = failures;
apply_learned_upstreams(&upstreams, &configured_upstreams);
let agents = crate::hooks::detect_agents();
for a in &agents {
if a.binding.model_relay_wiring().is_some() {
wiring.seed(a.agent_type());
}
}
let any_proxy_env = agents.iter().any(|a| {
a.binding.model_relay_wiring().is_some_and(|w| {
matches!(
w.endpoint,
crate::hooks::binding::EndpointConvention::ProxyEnv { .. }
)
})
});
let ca_info = if any_proxy_env {
crate::model_relay::ca::inspect(&crate::model_relay::ca::ca_dir(
&crate::config::openlatch_dir(),
))
} else {
None
};
let trust_clock = TrustClock {
wall: time::OffsetDateTime::now_utc(),
mono: std::time::Instant::now(),
};
let declared = proxy_env_declarations(&agents);
let skip = react_to_trust_loss_in(
&config,
&proxy_env_agents(&agents, &declared),
ca_info.as_ref(),
trust_clock,
&wiring.refusals(),
&mut holds,
&wiring,
);
let wired_agents: Vec<crate::hooks::DetectedAgent> = agents
.iter()
.filter(|a| !skip.contains(a.agent_type()))
.cloned()
.collect();
let declared: ProxyEnvDeclarations = declared
.into_iter()
.filter(|(a, _)| !skip.contains(a))
.collect();
let mut any_failed = wire_agents_in(
wired_agents,
port,
forwarding_broke,
&config,
&wiring,
&source_formats,
declared,
)
.await;
if let Some(endpoints) = &endpoints {
any_failed |= endpoints.reconcile(&config, &agents).await;
}
cline_hub::hub_tick(&config, port, &wiring, &agents).await;
apply_learned_upstreams(&upstreams, &configured_upstreams);
if any_failed {
consecutive_failures = consecutive_failures.saturating_add(1);
} else {
consecutive_failures = 0;
}
if let Some(endpoints) = &endpoints {
let wanted: Vec<std::path::PathBuf> = endpoints
.watch_files(&config, &agents)
.into_iter()
.filter(|f| f.parent().is_some_and(std::path::Path::is_dir))
.collect();
let armed = file_trigger.as_ref().map(|t| t.watched().to_vec());
if armed.as_deref() != Some(&wanted[..]) {
let previous = file_trigger.take();
let notify = file_changed.clone();
let armed = tokio::task::spawn_blocking(move || {
drop(previous);
watcher::spawn_file_trigger(&wanted, notify)
})
.await;
file_trigger = match armed {
Ok(Ok(trigger)) => trigger,
Ok(Err(e)) => {
tracing::warn!(
error = %e,
"cannot watch provider settings files — changes are picked up on the next tick"
);
None
}
Err(e) => {
tracing::warn!(
error = %e,
"arming the provider settings watch failed — changes are picked up on the next tick"
);
None
}
};
}
}
let delay = wiring_delay(tick, consecutive_failures);
let tick_done = tokio::time::sleep(delay);
tokio::pin!(tick_done);
loop {
tokio::select! {
_ = &mut tick_done => break,
_ = shutdown.changed() => return,
_ = file_changed.notified(), if file_trigger.is_some() => {
tokio::select! {
_ = tokio::time::sleep(FILE_QUIET_PERIOD) => {}
_ = shutdown.changed() => return,
}
if let Some(endpoints) = &endpoints {
let agents = crate::hooks::detect_agents();
let kept: Vec<_> = agents
.iter()
.filter(|a| wiring.declares_interceptor(a.agent_type()))
.cloned()
.collect();
publish_intercept_hosts(&wiring, &kept);
endpoints.reconcile(&config, &agents).await;
}
}
}
}
}
}
#[cfg(feature = "model-relay")]
const FILE_QUIET_PERIOD: std::time::Duration = std::time::Duration::from_secs(2);
#[cfg(all(test, feature = "model-relay"))]
async fn wire_agents(
agents: Vec<crate::hooks::DetectedAgent>,
port: u16,
forwarding_broke: bool,
config: &Config,
wiring: &crate::model_relay::preflight::WiringState,
source_formats: &crate::cloud::worker::SourceFormats,
) -> bool {
let declared = proxy_env_declarations(&agents);
wire_agents_in(
agents,
port,
forwarding_broke,
config,
wiring,
source_formats,
declared,
)
.await
}
#[cfg(feature = "model-relay")]
async fn wire_agents_in(
agents: Vec<crate::hooks::DetectedAgent>,
port: u16,
forwarding_broke: bool,
config: &Config,
wiring: &crate::model_relay::preflight::WiringState,
source_formats: &crate::cloud::worker::SourceFormats,
declared: ProxyEnvDeclarations,
) -> bool {
use crate::model_relay::preflight::{self, Verdict};
let interceptor = publish_declarations(wiring, declared);
let mut any_failed = false;
for agent in agents {
let a = agent.agent_type();
let plane = agent.binding.model_relay_wiring();
let Some(w) = plane else {
wiring.set_wired(a, false);
continue;
};
if wiring.is_wired(a) && !forwarding_broke {
continue;
}
let upstream = config.model_relay.upstream_for(w.wire_format);
let verdict = match &w.endpoint {
crate::hooks::binding::EndpointConvention::ProxyEnv {
intercept_hosts, ..
} => {
prove_proxy_env(
config,
port,
w.wire_format,
&upstream,
intercept_hosts,
interceptor.clone(),
!env_trust_only(config, &*agent.binding),
)
.await
}
_ => {
preflight::probe(port, w.wire_format, &upstream, preflight::PREFLIGHT_TIMEOUT).await
}
};
match verdict {
Ok(()) => {
wiring.set_verdict(a, Verdict::Ok);
if !wiring.is_wired(a) {
wire_model_relay_config(config, &*agent.binding, port, wiring, source_formats);
}
}
Err(reason) => {
let ca_refusal =
w.endpoint.is_proxy_env() && reason.starts_with(preflight::CA_REASON_MARKER);
if !ca_refusal {
any_failed = true;
}
if wiring.is_wired(a) {
tracing::error!(
port,
agent = a,
reason = %reason,
"model_relay preflight failed on a wired listener — removing the agent \
wiring so sessions fall back to a direct provider connection"
);
} else {
tracing::warn!(
port,
agent = a,
reason = %reason,
"model_relay preflight failed — the agent stays unwired and model calls \
go straight to the provider (nothing is captured)"
);
}
wiring.set_verdict(a, Verdict::Failed(reason));
if owns_wiring_for(config, &*agent.binding) {
if matches!(
w.endpoint,
crate::hooks::binding::EndpointConvention::ProxyEnv { .. }
) {
release_proxy_delivery(
&*agent.binding,
crate::hooks::model_relay_endpoints::ReleasedBy::Wiring,
);
} else {
unwire_model_relay_config(config, &*agent.binding);
}
wiring.set_wired(a, false);
}
}
}
}
any_failed
}
#[cfg(feature = "model-relay")]
fn proxy_env_declarations(agents: &[crate::hooks::DetectedAgent]) -> ProxyEnvDeclarations {
agents
.iter()
.filter_map(|a| match a.binding.model_relay_wiring()?.endpoint {
crate::hooks::binding::EndpointConvention::ProxyEnv {
intercept_hosts, ..
} => {
let mut hosts: std::collections::BTreeSet<String> = intercept_hosts
.iter()
.map(|h| h.to_ascii_lowercase())
.collect();
if let Some(ep) = a.binding.provider_endpoints() {
hosts.extend(ep.intercepted_hosts());
}
Some((a.agent_type(), hosts))
}
_ => None,
})
.collect()
}
#[cfg(all(test, feature = "model-relay"))]
fn proxy_env_hosts(agents: &[crate::hooks::DetectedAgent]) -> std::collections::BTreeSet<String> {
proxy_env_declarations(agents)
.into_values()
.flatten()
.collect()
}
#[cfg(feature = "model-relay")]
fn publish_intercept_hosts(
wiring: &crate::model_relay::preflight::WiringState,
agents: &[crate::hooks::DetectedAgent],
) -> Option<Arc<crate::model_relay::ca::Interceptor>> {
publish_declarations(wiring, proxy_env_declarations(agents))
}
#[cfg(feature = "model-relay")]
fn publish_declarations(
wiring: &crate::model_relay::preflight::WiringState,
declared: ProxyEnvDeclarations,
) -> Option<Arc<crate::model_relay::ca::Interceptor>> {
let ic = refresh_interceptor(
&wiring.intercept(),
declared.values().flatten().cloned().collect(),
);
wiring.set_interceptors(declared);
ic
}
#[cfg(feature = "model-relay")]
fn refresh_interceptor(
slot: &crate::model_relay::InterceptSlot,
hosts: std::collections::BTreeSet<String>,
) -> Option<Arc<crate::model_relay::ca::Interceptor>> {
let filled = slot.read().map(|g| g.clone()).unwrap_or(None);
match filled {
Some(ic) => {
if let Err(e) = ic.set_hosts(hosts) {
tracing::warn!(code = %e.code, error = %e.message,
"could not refresh the model relay's intercepted host set");
}
Some(ic)
}
None if hosts.is_empty() => None,
None => {
let dir = crate::model_relay::ca::ca_dir(&crate::config::openlatch_dir());
match crate::model_relay::ca::Interceptor::new(&dir, hosts) {
Ok(ic) => {
let ic = Arc::new(ic);
match slot.write() {
Ok(mut g) => *g = Some(ic.clone()),
Err(poisoned) => *poisoned.into_inner() = Some(ic.clone()),
}
Some(ic)
}
Err(e) => {
tracing::error!(code = %e.code, error = %e.message,
"the model relay could not create its certificate authority — \
ProxyEnv agents stay unwired");
None
}
}
}
}
}
#[cfg(feature = "model-relay")]
async fn prove_proxy_env(
config: &Config,
port: u16,
fmt: crate::model_relay::wire_format::WireFormat,
upstream: &str,
hosts: &[&str],
ic: Option<Arc<crate::model_relay::ca::Interceptor>>,
store_required: bool,
) -> Result<(), String> {
use crate::model_relay::preflight;
let Some(host) = hosts.first().copied() else {
return Err("ProxyEnv declares no intercept host".to_string());
};
let Some(ic) = ic else {
return Err(format!(
"{}the relay has no certificate authority — see the daemon log",
preflight::CA_REASON_MARKER
));
};
let dir = crate::model_relay::ca::ca_dir(&crate::config::openlatch_dir());
let Some(info) = crate::model_relay::ca::inspect(&dir) else {
return Err(format!(
"{}{} is missing or unreadable",
preflight::CA_REASON_MARKER,
crate::model_relay::ca::ca_pem_path(&dir).display()
));
};
if ca_too_close_to_expiry(info.not_after, time::OffsetDateTime::now_utc()) {
let reason = format!(
"{}{}{}",
preflight::CA_REASON_MARKER,
preflight::CA_EXPIRING_PREFIX,
info.not_after.date()
);
tracing::warn!(code = crate::error::ERR_MODEL_RELAY_CA_UNTRUSTED, host, reason = %reason,
"the model relay's certificate authority is close to expiry");
return Err(reason);
}
if store_required {
let store = crate::model_relay::trust_store::store();
let host_owned = host.to_string();
let ic_for_blocking = ic.clone();
let store_result = tokio::task::spawn_blocking(move || {
preflight::prove_store(&*store, &ic_for_blocking, &host_owned)
})
.await
.map_err(|e| {
format!(
"{}the trust-store check did not complete: {e}",
preflight::CA_REASON_MARKER
)
})?;
if let Err(e) = store_result {
tracing::warn!(code = crate::error::ERR_MODEL_RELAY_CA_UNTRUSTED, host, reason = %e,
"the user trust store did not accept the relay's certificate authority");
return Err(e);
}
} else {
static LOGGED: std::sync::Once = std::sync::Once::new();
LOGGED.call_once(|| {
tracing::info!(host, "isolated instance — env-trusted deliveries only");
});
}
let ca_pem = crate::model_relay::ca::ca_pem_path(&dir);
if let Err(e) = preflight::probe_intercept(
config,
port,
fmt,
upstream,
&ca_pem,
host,
preflight::PREFLIGHT_TIMEOUT,
)
.await
{
if e.starts_with(preflight::CA_REASON_MARKER) {
tracing::warn!(code = crate::error::ERR_MODEL_RELAY_CA_UNTRUSTED, host, reason = %e,
"the model relay's intercepted tunnel could not be proven");
}
return Err(e);
}
if store_required && !crate::model_relay::trust_store::platform_proven() {
tracing::warn!(
code = crate::error::ERR_MODEL_RELAY_TRUST_UNPROVEN,
host,
"user trust store verified on a platform whose trust path is not yet proven end to end"
);
}
Ok(())
}
#[cfg(feature = "model-relay")]
fn ca_too_close_to_expiry(not_after: time::OffsetDateTime, now: time::OffsetDateTime) -> bool {
(not_after - now).whole_days() <= crate::model_relay::ca::CA_RELEASE_DAYS
}
#[cfg(feature = "model-relay")]
fn wire_model_relay_config(
config: &Config,
binding: &dyn crate::hooks::binding::AgentBinding,
port: u16,
wiring: &crate::model_relay::preflight::WiringState,
source_formats: &crate::cloud::worker::SourceFormats,
) {
let path = crate::hooks::model_relay_config_path(binding);
let path_label = path
.as_ref()
.map(|p| p.display().to_string())
.unwrap_or_else(|| "the agent config".into());
if !owns_wiring_for(config, binding) {
tracing::info!(
port,
agent = binding.agent_type(),
"isolated model_relay instance — {path_label} left untouched; route sessions through \
this listener explicitly ({})",
crate::hooks::isolated_wiring_hint(binding, port),
);
return;
}
let install_id = match config.agent_id.clone() {
Some(id) => id,
None => crate::config::ensure_agent_id(&crate::config::openlatch_dir().join("config.toml"))
.unwrap_or_default(),
};
match crate::hooks::write_model_relay_config_trusting(
binding,
port,
&install_id,
!env_trust_only(config, binding),
) {
Ok(()) => {
wiring.set_wired(binding.agent_type(), true);
if let Some(w) = binding.model_relay_wiring() {
if w.endpoint.is_proxy_env() {
wiring.clear_intercept_proof(binding.agent_type());
} else {
wiring.set_wired_format(binding.agent_type(), w.wire_format);
source_formats.set(binding.agent_type(), w.wire_format);
}
}
tracing::info!(
port,
agent = binding.agent_type(),
path = %path_label,
"agent wired to the model relay — model calls route via http://127.0.0.1:{port}"
);
if let Some(crate::hooks::binding::EndpointConvention::ProxyEnv { delivery, .. }) =
binding.model_relay_wiring().map(|w| w.endpoint)
{
if profile_line_wanted(delivery, crate::supervision::owns_machine_supervision()) {
if let Some(home) = dirs::home_dir() {
let ol = crate::config::openlatch_dir();
let documents = if cfg!(windows) {
dirs::document_dir()
} else {
None
};
match crate::hooks::shell_profile::ensure_line(
&home,
documents.as_deref(),
&crate::hooks::proxy_env_file::env_sh_path(&ol),
&crate::hooks::proxy_env_file::env_ps1_path(&ol),
) {
Ok(files) => tracing::info!(
agent = binding.agent_type(),
files = ?files,
"shell profile sources the model relay wrapper — open a new \
shell to pick it up"
),
Err(e) => tracing::warn!(
code = %e.code,
error = %e.message,
"could not update the shell profile — shell-launched sessions \
are not captured"
),
}
}
}
}
if matches!(
binding.model_relay_wiring().map(|w| w.endpoint),
Some(crate::hooks::binding::EndpointConvention::EnvVars { .. })
) {
tracing::info!(
"Claude Code disables Remote Control while ANTHROPIC_BASE_URL is set — \
stop the daemon (`openlatch stop`) to restore a direct connection"
);
}
}
Err(e) => {
if e.code == crate::error::ERR_MODEL_RELAY_SETTINGS {
wiring.set_verdict(
binding.agent_type(),
crate::model_relay::preflight::Verdict::Failed(e.message.clone()),
);
}
if matches!(
binding.model_relay_wiring().map(|w| w.endpoint),
Some(crate::hooks::binding::EndpointConvention::ProxyEnv { .. })
) {
release_proxy_delivery(
binding,
crate::hooks::model_relay_endpoints::ReleasedBy::Wiring,
);
}
tracing::error!(
code = %e.code,
error = %e.message,
agent = binding.agent_type(),
path = %path_label,
"failed to wire agent to the model relay — agents will talk to the provider directly"
);
}
}
}
#[cfg(feature = "model-relay")]
fn profile_line_wanted(
delivery: &[crate::hooks::binding::ProxyDelivery],
owns_machine: bool,
) -> bool {
owns_machine
&& delivery
.iter()
.any(|d| matches!(d, crate::hooks::binding::ProxyDelivery::EnvFile { .. }))
}
#[cfg(feature = "model-relay")]
fn unwire_model_relay_config(config: &Config, binding: &dyn crate::hooks::binding::AgentBinding) {
let path_label = crate::hooks::model_relay_config_path(binding)
.map(|p| p.display().to_string())
.unwrap_or_else(|| "the agent config".into());
match crate::hooks::remove_model_relay_config(binding) {
Ok(()) => tracing::info!(
agent = binding.agent_type(),
path = %path_label,
"agent model relay wiring removed — agents connect to the provider directly"
),
Err(e) => tracing::warn!(
code = %e.code,
error = %e.message,
agent = binding.agent_type(),
path = %path_label,
"failed to remove agent model_relay wiring — the agent may still point at \
127.0.0.1:{}, which nothing is listening on",
config.model_relay.port,
),
}
}
#[cfg(feature = "model-relay")]
pub(crate) fn release_proxy_delivery(
binding: &dyn crate::hooks::binding::AgentBinding,
by: crate::hooks::model_relay_endpoints::ReleasedBy,
) {
match crate::hooks::release_proxy_env(binding, by) {
Ok(()) => tracing::info!(
agent = binding.agent_type(),
?by,
"proxy delivery released — the agent's surfaces go direct"
),
Err(e) => tracing::warn!(
code = %e.code,
error = %e.message,
agent = binding.agent_type(),
?by,
"could not release every proxy value — the agent may still name 127.0.0.1"
),
}
}
#[cfg(feature = "model-relay")]
pub(crate) fn ensure_proxy_env_trust(
config: &Config,
agents: &[crate::hooks::DetectedAgent],
owns_machine: bool,
) -> Option<crate::model_relay::trust_store::InstallOutcome> {
if !owns_machine {
return None;
}
let needed = agents.iter().any(|a| {
matches!(
a.binding.model_relay_wiring().map(|w| w.endpoint),
Some(crate::hooks::binding::EndpointConvention::ProxyEnv { .. })
) && owns_wiring_for(config, &*a.binding)
});
if !needed {
return None;
}
let dir = crate::model_relay::ca::ca_dir(&crate::config::openlatch_dir());
let ca = match crate::model_relay::ca::LocalCa::load_or_generate(&dir) {
Ok(ca) => ca,
Err(e) => {
return Some(crate::model_relay::trust_store::InstallOutcome::Refused(
format!("the relay CA could not be created: {}", e.message),
))
}
};
Some(
crate::model_relay::trust_store::store()
.install(&crate::model_relay::ca::ca_pem_path(&dir), &ca.sha256_hex()),
)
}
#[cfg(feature = "model-relay")]
fn release_every_provider_endpoint(config: &Config) {
for agent in crate::hooks::detect_agents() {
let Some(endpoints) = agent.binding.provider_endpoints() else {
continue;
};
if !owns_wiring_for(config, &*agent.binding) {
continue;
}
let summary = crate::hooks::provider_endpoints::release_all(
endpoints,
crate::hooks::model_relay_endpoints::ReleasedBy::Teardown,
);
if summary.released > 0 || !summary.failed.is_empty() {
tracing::info!(
agent = agent.agent_type(),
released = summary.released,
restored = summary.restored,
failed = summary.failed.len(),
"model relay is off — provider slots handed back"
);
}
}
}
#[cfg(feature = "model-relay")]
fn unwire_every_agent(config: &Config) {
for agent in crate::hooks::detect_agents() {
if agent.binding.model_relay_wiring().is_none() {
continue;
}
if owns_wiring_for(config, &*agent.binding) {
unwire_model_relay_config(config, &*agent.binding);
}
}
}
async fn signal_handler() {
#[cfg(unix)]
{
use tokio::signal::unix::{signal, SignalKind};
let mut sigterm =
signal(SignalKind::terminate()).expect("failed to register SIGTERM handler");
let mut sighup = signal(SignalKind::hangup()).expect("failed to register SIGHUP handler");
let stopped_by = tokio::select! {
_ = tokio::signal::ctrl_c() => SIGNAL_OPERATOR,
_ = sigterm.recv() => SIGNAL_SUPERVISOR,
_ = sighup.recv() => {
tracing::info!("received SIGHUP — draining (no config-reload path exists)");
SIGNAL_SUPERVISOR
}
};
SIGNAL.store(stopped_by, std::sync::atomic::Ordering::SeqCst);
}
#[cfg(not(unix))]
{
use tokio::signal::windows;
let mut ctrl_break = windows::ctrl_break().expect("failed to register ctrl_break handler");
tokio::select! {
_ = tokio::signal::ctrl_c() => {}
_ = ctrl_break.recv() => {}
}
SIGNAL.store(SIGNAL_OPERATOR, std::sync::atomic::Ordering::SeqCst);
}
}
static SIGNAL: std::sync::atomic::AtomicU8 = std::sync::atomic::AtomicU8::new(SIGNAL_NONE);
const SIGNAL_NONE: u8 = 0;
const SIGNAL_OPERATOR: u8 = 1;
#[cfg_attr(not(unix), allow(dead_code))]
const SIGNAL_SUPERVISOR: u8 = 2;
pub fn interrupted() -> bool {
SIGNAL.load(std::sync::atomic::Ordering::SeqCst) != SIGNAL_NONE
}
pub fn stopped_by_operator() -> bool {
SIGNAL.load(std::sync::atomic::Ordering::SeqCst) == SIGNAL_OPERATOR
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_restarted_update_completes_only_on_a_healthy_probe() {
let sentinel = update::UpdateSentinel {
from: "0.1.0".into(),
to: "0.2.0".into(),
applied_at: "2026-09-23T00:00:00Z".into(),
};
let status = Mutex::new(update::UpdateStatusSnapshot::restarted(&sentinel));
assert_eq!(
status.lock().unwrap().status,
update::UpdateStatusKind::InProgress
);
settle_restarted_update(&status, &sentinel, true);
let snap = status.lock().unwrap().clone();
assert_eq!(snap.status, update::UpdateStatusKind::Completed);
assert_eq!(snap.stage, None);
assert!(snap.ended_at.is_some());
let status = Mutex::new(update::UpdateStatusSnapshot::restarted(&sentinel));
settle_restarted_update(&status, &sentinel, false);
let snap = status.lock().unwrap().clone();
assert_eq!(snap.status, update::UpdateStatusKind::Failed);
assert_eq!(snap.stage, Some(update::ApplyStage::Healthz));
}
#[test]
fn settling_leaves_a_newer_apply_and_an_unseeded_status_alone() {
let sentinel = update::UpdateSentinel {
from: "0.1.0".into(),
to: "0.2.0".into(),
applied_at: "2026-09-23T00:00:00Z".into(),
};
for healthy in [true, false] {
let status = Mutex::new(update::UpdateStatusSnapshot::in_progress("0.2.0", "0.3.0"));
settle_restarted_update(&status, &sentinel, healthy);
let snap = status.lock().unwrap().clone();
assert_eq!(snap.status, update::UpdateStatusKind::InProgress);
assert_eq!(snap.stage, Some(update::ApplyStage::Check));
assert_eq!(snap.to.as_deref(), Some("0.3.0"));
assert_eq!(snap.ended_at, None);
}
let status = Mutex::new(update::UpdateStatusSnapshot::idle());
settle_restarted_update(&status, &sentinel, true);
assert_eq!(
status.lock().unwrap().status,
update::UpdateStatusKind::Idle
);
}
#[tokio::test]
async fn the_post_restart_probe_retries_before_giving_up() {
use std::sync::atomic::AtomicUsize;
let no_pause = [std::time::Duration::ZERO; 4];
let calls = AtomicUsize::new(0);
let healthy = probe_until_healthy(
|| async { calls.fetch_add(1, Ordering::SeqCst) >= 2 },
&no_pause,
)
.await;
assert!(healthy, "a daemon that answers on the third try is healthy");
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"stops at the first success"
);
let calls = AtomicUsize::new(0);
let healthy = probe_until_healthy(
|| async {
calls.fetch_add(1, Ordering::SeqCst);
false
},
&no_pause,
)
.await;
assert!(!healthy);
assert_eq!(
calls.load(Ordering::SeqCst),
5,
"one try, then one per pause"
);
let budget: std::time::Duration = POST_RESTART_PROBE_BACKOFF.iter().sum();
assert_eq!(budget, std::time::Duration::from_secs(15));
}
#[test]
fn a_recorded_upstream_reaches_the_live_handle() {
use crate::model_relay::wire_format::WireFormat;
let _dir_lock = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().expect("tempdir");
let previous = std::env::var_os("OPENLATCH_DIR");
std::env::set_var("OPENLATCH_DIR", dir.path());
crate::hooks::model_relay_endpoints::record(
&crate::hooks::upstream_record_key(WireFormat::OpenAiChatCompletions),
Some("http://127.0.0.1:11434".to_string()),
)
.expect("record the endpoint the wiring replaced");
let handle: crate::model_relay::UpstreamHandle =
Arc::new(std::sync::RwLock::new(std::collections::BTreeMap::new()));
let configured = std::collections::BTreeMap::new();
apply_learned_upstreams(&handle, &configured);
let installed = handle
.read()
.unwrap_or_else(|e| e.into_inner())
.get(WireFormat::OpenAiChatCompletions.as_str())
.map(|u| u.to_string());
match previous {
Some(v) => std::env::set_var("OPENLATCH_DIR", v),
None => std::env::remove_var("OPENLATCH_DIR"),
}
assert_eq!(
installed.as_deref(),
Some("http://127.0.0.1:11434/"),
"the relay must forward this format to the provider it replaced"
);
}
#[test]
fn a_forgotten_upstream_record_reverts_the_live_handle() {
use crate::model_relay::wire_format::WireFormat;
let _dir_lock = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().expect("tempdir");
let _env = crate::hooks::cline::EnvOverride::apply([(
"OPENLATCH_DIR",
Some(dir.path().as_os_str().to_os_string()),
)]);
let key = crate::hooks::upstream_record_key(WireFormat::AnthropicMessages);
let configured: std::collections::BTreeMap<String, String> = [(
WireFormat::AnthropicMessages.as_str().to_string(),
"https://api.anthropic.com".to_string(),
)]
.into();
let handle: crate::model_relay::UpstreamHandle =
Arc::new(std::sync::RwLock::new(std::collections::BTreeMap::new()));
let current = || {
handle
.read()
.unwrap_or_else(|e| e.into_inner())
.get(WireFormat::AnthropicMessages.as_str())
.map(|u| u.to_string())
};
crate::hooks::model_relay_endpoints::record(&key, Some("https://gw.corp".into()))
.expect("record");
apply_learned_upstreams(&handle, &configured);
assert_eq!(current().as_deref(), Some("https://gw.corp/"));
crate::hooks::model_relay_endpoints::forget(&key);
apply_learned_upstreams(&handle, &configured);
assert_eq!(current().as_deref(), Some("https://api.anthropic.com/"));
}
#[cfg(feature = "model-relay")]
#[test]
fn wiring_delay_backs_off_and_caps() {
let t = WIRING_TICK;
assert_eq!(wiring_delay(t, 0), t, "a green gate idles at the tick");
assert_eq!(
wiring_delay(t, 1),
t,
"the first failure retries at the tick"
);
assert_eq!(wiring_delay(t, 2), t * 2);
assert_eq!(wiring_delay(t, 3), t * 4);
assert_eq!(wiring_delay(t, 4), WIRING_BACKOFF_MAX);
assert_eq!(wiring_delay(t, 50), WIRING_BACKOFF_MAX);
assert_eq!(wiring_delay(t, u32::MAX), WIRING_BACKOFF_MAX);
}
#[cfg(feature = "model-relay")]
struct WiringFixture {
port: u16,
_openlatch_dir: tempfile::TempDir,
agent_root: tempfile::TempDir,
_dir_lock: std::sync::MutexGuard<'static, ()>,
previous_dir: Option<std::ffi::OsString>,
}
#[cfg(feature = "model-relay")]
impl Drop for WiringFixture {
fn drop(&mut self) {
match self.previous_dir.take() {
Some(v) => std::env::set_var("OPENLATCH_DIR", v),
None => std::env::remove_var("OPENLATCH_DIR"),
}
}
}
#[cfg(feature = "model-relay")]
impl WiringFixture {
async fn new() -> Self {
use crate::model_relay::wire_format::WireFormat;
use crate::model_relay::{mock, serve_ephemeral, ModelRelayState};
let upstream = mock::spawn_always_200().await;
let dead = mock::closed_port().await;
let map: std::collections::BTreeMap<String, String> = [
(
WireFormat::AnthropicMessages.as_str().to_string(),
format!("http://127.0.0.1:{upstream}"),
),
(
WireFormat::OpenAiResponses.as_str().to_string(),
format!("http://127.0.0.1:{dead}"),
),
]
.into_iter()
.collect();
let state = Arc::new(
ModelRelayState::new(
reqwest::Url::parse(&format!("http://127.0.0.1:{upstream}")).unwrap(),
0,
8,
&[],
)
.with_upstream_map(map),
);
let port = serve_ephemeral(state).await;
let _dir_lock = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let _openlatch_dir = tempfile::tempdir().expect("tempdir");
let previous_dir = std::env::var_os("OPENLATCH_DIR");
std::env::set_var("OPENLATCH_DIR", _openlatch_dir.path());
Self {
port,
_openlatch_dir,
agent_root: tempfile::tempdir().expect("tempdir"),
_dir_lock,
previous_dir,
}
}
fn config(&self) -> Config {
let mut cfg = Config::defaults();
cfg.agent_id = Some("agt_test".to_string());
cfg.model_relay.own_agent_wiring = Some(true);
cfg.model_relay.port = crate::model_relay::default_model_relay_port();
cfg
}
fn agent(
&self,
agent_type: &'static str,
fmt: crate::model_relay::wire_format::WireFormat,
machine_global: bool,
) -> crate::hooks::DetectedAgent {
use crate::hooks::binding::{EndpointConvention, ModelRelayWiring};
let dir = self.agent_root.path().join(agent_type);
std::fs::create_dir_all(&dir).expect("agent dir");
crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: Arc::new(crate::hooks::binding::test_support::FakeBinding {
agent_type,
display_name: agent_type,
config_dir: dir,
model_relay_wiring: Some(ModelRelayWiring {
wire_format: fmt,
endpoint: EndpointConvention::EnvVars {
base_url: "ANTHROPIC_BASE_URL",
headers: "ANTHROPIC_CUSTOM_HEADERS",
},
install_id_header: "x-openlatch-install-id",
}),
config_is_machine_global: machine_global,
..Default::default()
}),
}
}
fn is_written(&self, agent_type: &str) -> bool {
let path = self
.agent_root
.path()
.join(agent_type)
.join("settings.json");
std::fs::read_to_string(path)
.map(|raw| raw.contains("ANTHROPIC_BASE_URL"))
.unwrap_or(false)
}
}
#[cfg(feature = "model-relay")]
#[tokio::test(flavor = "multi_thread")]
async fn wiring_loop_continues_past_a_failed_agent() {
use crate::model_relay::preflight::{Verdict, WiringState};
use crate::model_relay::wire_format::WireFormat;
let fx = WiringFixture::new().await;
let cfg = fx.config();
let wiring = WiringState::default();
let formats = crate::cloud::worker::SourceFormats::new();
let agents = vec![
fx.agent("codex-cli", WireFormat::OpenAiResponses, false),
fx.agent("claude-code", WireFormat::AnthropicMessages, false),
];
let any_failed = wire_agents(agents, fx.port, false, &cfg, &wiring, &formats).await;
assert!(
any_failed,
"the Codex probe failed and the tick must say so"
);
assert!(
wiring.is_wired("claude-code"),
"the healthy agent is still wired after the failing one"
);
assert_eq!(wiring.verdict("claude-code"), Verdict::Ok);
assert!(fx.is_written("claude-code"));
assert!(
!wiring.is_wired("codex-cli"),
"the failing agent is not wired"
);
assert!(matches!(wiring.verdict("codex-cli"), Verdict::Failed(_)));
}
#[cfg(feature = "model-relay")]
#[tokio::test(flavor = "multi_thread")]
async fn probe_failure_blocks_only_that_agents_write() {
use crate::model_relay::preflight::WiringState;
use crate::model_relay::wire_format::WireFormat;
let fx = WiringFixture::new().await;
let cfg = fx.config();
let wiring = WiringState::default();
let formats = crate::cloud::worker::SourceFormats::new();
let agents = vec![
fx.agent("claude-code", WireFormat::AnthropicMessages, false),
fx.agent("codex-cli", WireFormat::OpenAiResponses, false),
];
wire_agents(agents, fx.port, false, &cfg, &wiring, &formats).await;
assert!(
fx.is_written("claude-code"),
"the agent whose probe passed is written"
);
assert!(
!fx.is_written("codex-cli"),
"the agent whose probe failed must be left untouched — the endpoint is \
written only after a proven round trip IN THAT AGENT'S FORMAT"
);
}
#[cfg(feature = "model-relay")]
#[tokio::test(flavor = "multi_thread")]
async fn failed_probe_clears_wired_so_the_next_tick_reprobes() {
use crate::model_relay::preflight::{Verdict, WiringState};
use crate::model_relay::wire_format::WireFormat;
let fx = WiringFixture::new().await;
let cfg = fx.config();
let wiring = WiringState::default();
let formats = crate::cloud::worker::SourceFormats::new();
wiring.set_wired("codex-cli", true);
wiring.set_verdict("codex-cli", Verdict::Ok);
let first = wire_agents(
vec![fx.agent("codex-cli", WireFormat::OpenAiResponses, false)],
fx.port,
true,
&cfg,
&wiring,
&formats,
)
.await;
assert!(first, "the probe failed");
assert!(
!wiring.is_wired("codex-cli"),
"a failed probe must clear the in-memory gate as well as the file"
);
wiring.set_verdict("codex-cli", Verdict::Pending);
let second = wire_agents(
vec![fx.agent("codex-cli", WireFormat::OpenAiResponses, false)],
fx.port,
false,
&cfg,
&wiring,
&formats,
)
.await;
assert!(second, "the next tick re-probes it");
assert!(
matches!(wiring.verdict("codex-cli"), Verdict::Failed(_)),
"a re-probed agent gets a fresh verdict; a skipped one keeps Pending"
);
}
#[cfg(feature = "model-relay")]
#[tokio::test(flavor = "multi_thread")]
async fn isolated_daemon_never_writes_a_machine_global_agent() {
use crate::model_relay::preflight::{Verdict, WiringState};
use crate::model_relay::wire_format::WireFormat;
let fx = WiringFixture::new().await;
let mut cfg = fx.config();
cfg.model_relay.port = crate::model_relay::default_model_relay_port() + 99;
let wiring = WiringState::default();
let formats = crate::cloud::worker::SourceFormats::new();
let agents = vec![
fx.agent("relocated", WireFormat::AnthropicMessages, false),
fx.agent("machine-global", WireFormat::AnthropicMessages, true),
];
let any_failed = wire_agents(agents, fx.port, false, &cfg, &wiring, &formats).await;
assert!(!any_failed, "both probes pass; only the guard differs");
assert_eq!(wiring.verdict("relocated"), Verdict::Ok);
assert_eq!(wiring.verdict("machine-global"), Verdict::Ok);
assert!(
fx.is_written("relocated"),
"an agent whose config this sandbox owns is wired"
);
assert!(
wiring.is_wired("relocated"),
"and its wiring flag follows the write"
);
assert!(
!fx.is_written("machine-global"),
"a machine-global agent's config is NOT this instance's to write"
);
assert!(
!wiring.is_wired("machine-global"),
"and it must not be reported as wired either"
);
}
#[cfg(feature = "model-relay")]
#[test]
fn owns_wiring_for_asks_the_binding_on_a_non_default_port() {
use crate::hooks::binding::test_support::FakeBinding;
let relocated = FakeBinding {
config_is_machine_global: false,
..Default::default()
};
let global = FakeBinding {
config_is_machine_global: true,
..Default::default()
};
let mut cfg = Config::defaults();
cfg.model_relay.port = crate::model_relay::default_model_relay_port();
cfg.model_relay.own_agent_wiring = None;
assert!(owns_wiring_for(&cfg, &relocated));
assert!(owns_wiring_for(&cfg, &global));
cfg.model_relay.own_agent_wiring = Some(false);
assert!(!owns_wiring_for(&cfg, &relocated));
assert!(!owns_wiring_for(&cfg, &global));
cfg.model_relay.port = crate::model_relay::default_model_relay_port() + 99;
cfg.model_relay.own_agent_wiring = None;
assert!(!owns_wiring_for(&cfg, &relocated));
assert!(!owns_wiring_for(&cfg, &global));
cfg.model_relay.own_agent_wiring = Some(true);
assert!(owns_wiring_for(&cfg, &relocated));
assert!(
!owns_wiring_for(&cfg, &global),
"opting in is a decision about your own sandbox, not a way to seize a shared file"
);
}
#[cfg(feature = "model-relay")]
struct TrustLossFixture {
_env: crate::hooks::cline::EnvOverride,
_dir_lock: std::sync::MutexGuard<'static, ()>,
ol_dir: tempfile::TempDir,
}
#[cfg(feature = "model-relay")]
impl TrustLossFixture {
fn new() -> Self {
let _dir_lock = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let ol_dir = tempfile::tempdir().expect("tempdir");
let _env = crate::hooks::cline::EnvOverride::apply([(
"OPENLATCH_DIR",
Some(ol_dir.path().as_os_str().to_os_string()),
)]);
Self {
_env,
_dir_lock,
ol_dir,
}
}
fn write_wrapper(&self, agent: &str, command: &str) {
let ca_pem = self
.ol_dir
.path()
.join("model-relay")
.join("ca")
.join("ca.pem");
crate::hooks::proxy_env_file::write_block(
agent,
&[command],
"http://127.0.0.1:7600",
&ca_pem,
"127.0.0.1:7600",
)
.expect("write wrapper block");
}
fn env_sh(&self) -> String {
std::fs::read_to_string(crate::hooks::proxy_env_file::env_sh_path(
self.ol_dir.path(),
))
.unwrap_or_default()
}
fn proxy_env_agent(
&self,
agent_type: &'static str,
delivery: &'static [crate::hooks::binding::ProxyDelivery],
hosts: &'static [&'static str],
) -> crate::hooks::DetectedAgent {
use crate::hooks::binding::test_support::{proxy_env_wiring, FakeBinding};
crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: Arc::new(FakeBinding {
agent_type,
model_relay_wiring: Some(proxy_env_wiring(hosts, delivery)),
..Default::default()
}),
}
}
}
#[test]
#[cfg(feature = "model-relay")]
fn seven_days_before_expiry_every_proxy_env_delivery_is_released() {
use crate::hooks::binding::{EndpointConvention, ModelRelayWiring, ProxyDelivery};
use crate::model_relay::ca::CaInfo;
use crate::model_relay::ca_lifecycle::expiring_reason;
use crate::model_relay::intercept::RefusalLedger;
use crate::model_relay::preflight::{Verdict, WiringState};
use crate::model_relay::wire_format::WireFormat;
const DELIVERY_A: &[ProxyDelivery] = &[ProxyDelivery::EnvFile {
commands: &["fake-a-cli"],
}];
const DELIVERY_B: &[ProxyDelivery] = &[ProxyDelivery::EnvFile {
commands: &["fake-b-cli"],
}];
let cfg = Config::defaults();
let at = TrustClock {
wall: time::OffsetDateTime::now_utc(),
mono: std::time::Instant::now(),
};
{
let fx = TrustLossFixture::new();
fx.write_wrapper("fake-a", "fake-a-cli");
fx.write_wrapper("fake-b", "fake-b-cli");
let agent_a = fx.proxy_env_agent("fake-a", DELIVERY_A, &["a.test"]);
let agent_b = fx.proxy_env_agent("fake-b", DELIVERY_B, &["b.test"]);
let envvars_agent = crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: Arc::new(crate::hooks::binding::test_support::FakeBinding {
agent_type: "claude-shaped",
model_relay_wiring: Some(ModelRelayWiring {
wire_format: WireFormat::AnthropicMessages,
endpoint: EndpointConvention::EnvVars {
base_url: "ANTHROPIC_BASE_URL",
headers: "ANTHROPIC_CUSTOM_HEADERS",
},
install_id_header: "x-openlatch-install-id",
}),
..Default::default()
}),
};
let agents = vec![agent_a, agent_b, envvars_agent];
let wiring = WiringState::default();
wiring.set_wired("fake-a", true);
wiring.set_wired("fake-b", true);
let ca = CaInfo {
sha256_hex: "deadbeef".to_string(),
not_after: time::OffsetDateTime::now_utc() + time::Duration::days(6),
};
let refusals = RefusalLedger::default();
let mut holds = std::collections::BTreeMap::new();
let skip =
react_to_trust_loss(&cfg, &agents, Some(&ca), at, &refusals, &mut holds, &wiring);
assert!(skip.contains("fake-a"));
assert!(skip.contains("fake-b"));
assert!(!skip.contains("claude-shaped"));
let content = fx.env_sh();
assert!(
!content.contains("fake-a"),
"fake-a's block must be gone: {content}"
);
assert!(
!content.contains("fake-b"),
"fake-b's block must be gone: {content}"
);
let reason = expiring_reason(ca.not_after);
assert_eq!(wiring.verdict("fake-a"), Verdict::Failed(reason.clone()));
assert_eq!(wiring.verdict("fake-b"), Verdict::Failed(reason));
assert!(!wiring.is_wired("fake-a"));
assert!(!wiring.is_wired("fake-b"));
}
{
let fx2 = TrustLossFixture::new();
fx2.write_wrapper("fake-a", "fake-a-cli");
let agent_a2 = fx2.proxy_env_agent("fake-a", DELIVERY_A, &["a.test"]);
let agents2 = vec![agent_a2];
let wiring2 = WiringState::default();
wiring2.set_wired("fake-a", true);
let ca8 = CaInfo {
sha256_hex: "deadbeef".to_string(),
not_after: time::OffsetDateTime::now_utc() + time::Duration::days(8),
};
let refusals2 = RefusalLedger::default();
let mut holds2 = std::collections::BTreeMap::new();
let skip2 = react_to_trust_loss(
&cfg,
&agents2,
Some(&ca8),
at,
&refusals2,
&mut holds2,
&wiring2,
);
assert!(
skip2.is_empty(),
"outside CA_RELEASE_DAYS, nothing is released"
);
assert!(fx2.env_sh().contains("fake-a"), "the block must survive");
}
}
#[test]
#[cfg(feature = "model-relay")]
fn a_tripped_host_releases_only_the_bindings_that_name_it() {
use crate::hooks::binding::ProxyDelivery;
use crate::model_relay::ca_lifecycle::refused_reason;
use crate::model_relay::intercept::{RefusalKind, RefusalLedger};
use crate::model_relay::preflight::{Verdict, WiringState};
const DELIVERY_A: &[ProxyDelivery] = &[ProxyDelivery::EnvFile {
commands: &["fake-a-cli"],
}];
const DELIVERY_B: &[ProxyDelivery] = &[ProxyDelivery::EnvFile {
commands: &["fake-b-cli"],
}];
let fx = TrustLossFixture::new();
fx.write_wrapper("fake-a", "fake-a-cli");
fx.write_wrapper("fake-b", "fake-b-cli");
let agent_a = fx.proxy_env_agent("fake-a", DELIVERY_A, &["a.test"]);
let agent_b = fx.proxy_env_agent("fake-b", DELIVERY_B, &["b.test"]);
let agents = vec![agent_a, agent_b];
let cfg = Config::defaults();
let wiring = WiringState::default();
wiring.set_wired("fake-a", true);
wiring.set_wired("fake-b", true);
let refusals = RefusalLedger::default();
for _ in 0..3 {
refusals.record_refusal("a.test", RefusalKind::HandshakeEof);
}
let mut holds = std::collections::BTreeMap::new();
let at = TrustClock {
wall: time::OffsetDateTime::now_utc(),
mono: std::time::Instant::now(),
};
let skip = react_to_trust_loss(&cfg, &agents, None, at, &refusals, &mut holds, &wiring);
assert!(skip.contains("fake-a"));
assert!(!skip.contains("fake-b"));
assert_eq!(
wiring.verdict("fake-a"),
Verdict::Failed(refused_reason("a.test"))
);
assert!(!wiring.is_wired("fake-a"));
assert!(wiring.is_wired("fake-b"), "B is untouched");
assert_eq!(
wiring.verdict("fake-b"),
Verdict::Pending,
"B's verdict was never set by this reaction"
);
assert!(!fx.env_sh().contains("fake-a"), "A's block is gone");
assert!(fx.env_sh().contains("fake-b"), "B's block survives");
}
#[test]
#[cfg(feature = "model-relay")]
fn a_hold_blocks_rearming_until_restart() {
use crate::hooks::binding::ProxyDelivery;
use crate::model_relay::intercept::{RefusalKind, RefusalLedger};
use crate::model_relay::preflight::WiringState;
const DELIVERY_A: &[ProxyDelivery] = &[ProxyDelivery::EnvFile {
commands: &["fake-a-cli"],
}];
let fx = TrustLossFixture::new();
fx.write_wrapper("fake-a", "fake-a-cli");
let agent_a = fx.proxy_env_agent("fake-a", DELIVERY_A, &["a.test"]);
let agents = vec![agent_a];
let cfg = Config::defaults();
let wiring = WiringState::default();
wiring.set_wired("fake-a", true);
let refusals = RefusalLedger::default();
for _ in 0..3 {
refusals.record_refusal("a.test", RefusalKind::HandshakeEof);
}
let mut holds = std::collections::BTreeMap::new();
let base = std::time::Instant::now();
let at1 = TrustClock {
wall: time::OffsetDateTime::now_utc(),
mono: base,
};
let skip1 = react_to_trust_loss(&cfg, &agents, None, at1, &refusals, &mut holds, &wiring);
assert!(skip1.contains("fake-a"));
assert!(holds.contains_key("a.test"));
assert!(!wiring.is_wired("fake-a"));
let at2 = TrustClock {
wall: time::OffsetDateTime::now_utc(),
mono: base + std::time::Duration::from_secs(11 * 60),
};
let skip2 = react_to_trust_loss(&cfg, &agents, None, at2, &refusals, &mut holds, &wiring);
assert!(skip2.contains("fake-a"));
assert!(holds.contains_key("a.test"));
let at3 = TrustClock {
wall: time::OffsetDateTime::now_utc(),
mono: base + std::time::Duration::from_secs(2 * 60 * 60),
};
let skip3 = react_to_trust_loss(&cfg, &agents, None, at3, &refusals, &mut holds, &wiring);
assert!(skip3.contains("fake-a"));
assert!(holds.contains_key("a.test"));
let fresh_refusals = RefusalLedger::default();
let mut fresh_holds = std::collections::BTreeMap::new();
let skip4 = react_to_trust_loss(
&cfg,
&agents,
None,
at3,
&fresh_refusals,
&mut fresh_holds,
&wiring,
);
assert!(!skip4.contains("fake-a"), "a restart re-arms");
}
#[cfg(feature = "model-relay")]
#[test]
fn wiring_tick_defaults_to_production_and_never_busy_loops() {
assert_eq!(wiring_tick(), WIRING_TICK, "unset must mean the real tick");
assert_eq!(
wiring_delay(std::time::Duration::from_millis(50), 0),
std::time::Duration::from_millis(50)
);
}
#[test]
fn test_openlatch_marker_detected_in_settings() {
let with_hooks = r#"{"hooks": {"_openlatch": true, "preToolUse": []}}"#;
assert!(with_hooks.contains("\"_openlatch\""));
let without_hooks = r#"{"hooks": {"preToolUse": []}}"#;
assert!(!without_hooks.contains("\"_openlatch\""));
}
#[test]
fn test_format_uptime_seconds_only() {
assert_eq!(format_uptime(0), "0s");
assert_eq!(format_uptime(45), "45s");
assert_eq!(format_uptime(59), "59s");
}
#[test]
fn test_format_uptime_minutes_and_seconds() {
assert_eq!(format_uptime(60), "1m0s");
assert_eq!(format_uptime(192), "3m12s");
assert_eq!(format_uptime(3599), "59m59s");
}
#[test]
fn test_format_uptime_hours_and_minutes() {
assert_eq!(format_uptime(3600), "1h0m");
assert_eq!(format_uptime(8094), "2h14m");
assert_eq!(format_uptime(7200), "2h0m");
}
mod policy_startup {
use super::*;
use crate::core::policy::store::{self, BundleMeta};
use crate::generated::types::PolicyBundle;
const BODY: &str = r#"{"schema_version":1,"revision":42,"organization_id":"0192f8a1-4c3b-7e2a-9f10-5d8c3b1a7e42","built_at":"2026-07-21T09:00:00Z","enforcement_enabled":true,"signature":null,"rules":[{"rule_id":"OL-CMD-ENF","kind":"command","match_pattern":"*olcanary-enforce*","action":"deny","mode":"enforce","severity":"high","reason":"Canary enforce"}]}"#;
fn seed(base: &std::path::Path, last_poll_ok_at: Option<&str>) {
let body = BODY.as_bytes();
let bundle: PolicyBundle = serde_json::from_slice(body).expect("fixture parses");
let mut meta = BundleMeta::activated(
&bundle,
store::digest_of(body),
Some("\"sha256:deadbeef\"".to_string()),
);
meta.last_poll_ok_at = last_poll_ok_at.map(str::to_string);
store::store(base, body, &meta).expect("cache writes");
}
#[test]
fn cached_bundle_is_resident_before_anything_can_serve() {
let tmp = tempfile::tempdir().expect("tempdir");
seed(tmp.path(), None);
let runtime = PolicyRuntime::load_from_disk(tmp.path());
let guard = runtime.handle.load();
let bundle = guard.as_ref().as_ref().expect("bundle resident at startup");
assert_eq!(bundle.schema_version, 1);
assert_eq!(
bundle.digest.as_deref(),
Some(store::digest_of(BODY.as_bytes()).as_str())
);
assert_eq!(bundle.revision, 42);
assert_eq!(bundle.command_rules.len(), 1);
assert_eq!(bundle.command_rules[0].rule_id, "OL-CMD-ENF");
assert!(bundle.enforcement_enabled);
}
#[test]
fn tampered_bundle_is_rejected_and_the_daemon_starts_with_no_policy() {
let tmp = tempfile::tempdir().expect("tempdir");
seed(tmp.path(), None);
std::fs::write(
store::bundle_path(tmp.path()),
BODY.replace(r#""rules":[{"#, r#""rules":[{"x":1,"#),
)
.expect("tamper writes");
let runtime = PolicyRuntime::load_from_disk(tmp.path());
assert!(
runtime.handle.load().is_none(),
"a tampered bundle must never activate"
);
}
#[test]
fn digest_valid_but_unsupported_cache_is_discarded_before_polling() {
let tmp = tempfile::tempdir().expect("tempdir");
let valid: PolicyBundle = serde_json::from_str(BODY).expect("fixture parses");
let unsupported = br#"{"schema_version":3}"#;
let meta = BundleMeta::activated(
&valid,
store::digest_of(unsupported),
Some("\"sha256:stale\"".to_string()),
);
store::store(tmp.path(), unsupported, &meta).expect("cache writes");
let runtime = PolicyRuntime::load_from_disk(tmp.path());
assert!(runtime.handle.load().is_none());
assert!(!store::bundle_path(tmp.path()).exists());
assert!(!store::meta_path(tmp.path()).exists());
}
#[test]
fn no_cache_starts_with_no_policy() {
let tmp = tempfile::tempdir().expect("tempdir");
let runtime = PolicyRuntime::load_from_disk(tmp.path());
assert!(runtime.handle.load().is_none());
assert_eq!(runtime.last_poll_ok_at.load(Ordering::Relaxed), 0);
}
#[test]
fn poll_clock_is_seeded_from_the_meta_file() {
let tmp = tempfile::tempdir().expect("tempdir");
seed(tmp.path(), Some("2026-07-21T09:00:00Z"));
let runtime = PolicyRuntime::load_from_disk(tmp.path());
assert_eq!(
runtime.last_poll_ok_at.load(Ordering::Relaxed),
1_784_624_400
);
assert!(runtime.last_fetch_ok.load(Ordering::Relaxed));
}
}
#[cfg(feature = "model-relay")]
mod proxy_env {
use super::*;
use crate::hooks::binding::test_support::{proxy_env_wiring, FakeBinding};
use crate::hooks::binding::{EndpointConvention, ModelRelayWiring, ProxyDelivery};
use crate::hooks::model_relay_endpoints::ReleasedBy;
use crate::hooks::{proxy_env_file, proxy_settings};
use crate::model_relay::preflight::{Verdict, WiringState};
use crate::model_relay::trust_store::{
self, test_support::FakeStore, InstallOutcome, TrustStore,
};
use crate::model_relay::wire_format::WireFormat;
use crate::model_relay::{mock, serve_ephemeral, ModelRelayState};
const HOST: &str = "a.test";
const AGENT: &str = "fake-a";
const KEY: &str = "fake.proxy";
const SEED: &str = "{\n \"a\": 1,\n}\n";
const FAKE_DELIVERY: &[ProxyDelivery] = &[
ProxyDelivery::SettingsKey {
file: fake_settings_file,
key: KEY,
},
ProxyDelivery::EnvFile {
commands: &["fakecli"],
},
];
const FOREIGN_ONLY_DELIVERY: &[ProxyDelivery] = &[ProxyDelivery::SettingsKey {
file: fake_settings_file,
key: KEY,
}];
const ENV_FILE_ONLY_DELIVERY: &[ProxyDelivery] = &[ProxyDelivery::EnvFile {
commands: &["fakecli"],
}];
fn fake_settings_file() -> Option<std::path::PathBuf> {
Some(
crate::config::openlatch_dir()
.join("fake-agent")
.join("settings.json"),
)
}
fn settings_path() -> std::path::PathBuf {
fake_settings_file().expect("fake_settings_file always answers Some")
}
fn seed_settings(content: &str) {
let path = settings_path();
std::fs::create_dir_all(path.parent().expect("parent")).expect("settings dir");
std::fs::write(&path, content).expect("seed settings");
}
fn proxy_env_test_config(agent_id: &str) -> Config {
let mut cfg = Config::defaults();
cfg.agent_id = Some(agent_id.to_string());
cfg.model_relay.own_agent_wiring = Some(true);
cfg.model_relay.port = crate::model_relay::default_model_relay_port();
cfg
}
fn env_sh_exists() -> bool {
proxy_env_file::env_sh_path(&crate::config::openlatch_dir()).exists()
}
fn marker_exists() -> bool {
proxy_env_file::live_marker_path(&crate::config::openlatch_dir()).exists()
}
fn settings_record(
agent: &str,
key: &str,
) -> Option<crate::hooks::model_relay_endpoints::EndpointRecord> {
let prefix = format!("{}{agent}:", proxy_settings::RECORD_PREFIX);
let recs = crate::hooks::model_relay_endpoints::endpoint_records(&prefix)
.expect("read the model-relay-endpoints ledger");
recs.into_iter()
.find(|(k, _)| k.starts_with(&format!("{prefix}{key}@")))
.map(|(_, r)| r)
}
struct ProxyEnvFixture {
port: u16,
_openlatch_dir: tempfile::TempDir,
home_dir: tempfile::TempDir,
wiring: Arc<WiringState>,
fake_store: Arc<FakeStore>,
previous_dir: Option<std::ffi::OsString>,
previous_home: Option<std::ffi::OsString>,
_trust_guard: trust_store::TrustStoreGuard,
_home_lock: std::sync::MutexGuard<'static, ()>,
_dir_lock: std::sync::MutexGuard<'static, ()>,
}
impl Drop for ProxyEnvFixture {
fn drop(&mut self) {
match self.previous_home.take() {
Some(v) => std::env::set_var("HOME", v),
None => std::env::remove_var("HOME"),
}
match self.previous_dir.take() {
Some(v) => std::env::set_var("OPENLATCH_DIR", v),
None => std::env::remove_var("OPENLATCH_DIR"),
}
}
}
impl ProxyEnvFixture {
async fn new() -> Self {
Self::build(mock::spawn_always_200().await).await
}
async fn build(upstream: u16) -> Self {
let wiring = Arc::new(WiringState::default());
let map: std::collections::BTreeMap<String, String> = [(
WireFormat::AnthropicMessages.as_str().to_string(),
format!("http://127.0.0.1:{upstream}"),
)]
.into_iter()
.collect();
let state = Arc::new(
ModelRelayState::new(
reqwest::Url::parse(&format!("http://127.0.0.1:{upstream}")).unwrap(),
0,
8,
&[],
)
.with_upstream_map(map)
.with_explicit_upstreams([WireFormat::AnthropicMessages.as_str()].into())
.with_wiring(wiring.clone())
.with_intercept(wiring.intercept())
.with_refusals(wiring.refusals()),
);
let port = serve_ephemeral(state).await;
let _dir_lock = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let _home_lock = crate::hooks::claude_code::CONFIG_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let _openlatch_dir = tempfile::tempdir().expect("tempdir");
let previous_dir = std::env::var_os("OPENLATCH_DIR");
std::env::set_var("OPENLATCH_DIR", _openlatch_dir.path());
let home_dir = tempfile::tempdir().expect("tempdir");
let previous_home = std::env::var_os("HOME");
std::env::set_var("HOME", home_dir.path());
let fake_store = Arc::new(FakeStore::default());
let _trust_guard =
trust_store::install_for_tests(fake_store.clone() as Arc<dyn TrustStore>);
seed_settings(SEED);
Self {
port,
_openlatch_dir,
home_dir,
wiring,
fake_store,
previous_dir,
previous_home,
_trust_guard,
_home_lock,
_dir_lock,
}
}
fn config(&self) -> Config {
proxy_env_test_config("agt_test")
}
fn agent(
&self,
hosts: &'static [&'static str],
delivery: &'static [ProxyDelivery],
) -> crate::hooks::DetectedAgent {
crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: Arc::new(FakeBinding {
agent_type: AGENT,
display_name: "Fake",
config_dir: crate::config::openlatch_dir().join("fake-agent"),
model_relay_wiring: Some(proxy_env_wiring(hosts, delivery)),
..Default::default()
}),
}
}
}
fn env_vars_agent(agent_type: &'static str) -> crate::hooks::DetectedAgent {
crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: Arc::new(FakeBinding {
agent_type,
model_relay_wiring: Some(ModelRelayWiring {
wire_format: WireFormat::AnthropicMessages,
endpoint: EndpointConvention::EnvVars {
base_url: "ANTHROPIC_BASE_URL",
headers: "ANTHROPIC_CUSTOM_HEADERS",
},
install_id_header: "x-openlatch-install-id",
}),
..Default::default()
}),
}
}
fn toml_provider_agent(agent_type: &'static str) -> crate::hooks::DetectedAgent {
crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: Arc::new(FakeBinding {
agent_type,
model_relay_wiring: Some(ModelRelayWiring {
wire_format: WireFormat::OpenAiResponses,
endpoint: EndpointConvention::TomlProvider {
provider_name: "fake-provider",
wire_api: "chat",
},
install_id_header: "x-openlatch-install-id",
}),
..Default::default()
}),
}
}
#[test]
fn the_interceptor_is_built_only_for_proxy_env() {
let non_proxy_env = vec![
env_vars_agent("envvars-agent"),
toml_provider_agent("toml-agent"),
];
assert!(proxy_env_hosts(&non_proxy_env).is_empty());
let proxy_env_agent = crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: Arc::new(FakeBinding {
agent_type: AGENT,
model_relay_wiring: Some(proxy_env_wiring(&[HOST], FAKE_DELIVERY)),
..Default::default()
}),
};
let hosts = proxy_env_hosts(&[proxy_env_agent]);
assert_eq!(hosts, [HOST.to_string()].into_iter().collect());
let _dir_lock = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().expect("tempdir");
let previous = std::env::var_os("OPENLATCH_DIR");
std::env::set_var("OPENLATCH_DIR", dir.path());
let slot: crate::model_relay::InterceptSlot = Default::default();
let ic = refresh_interceptor(&slot, Default::default());
assert!(ic.is_none(), "an empty host set must leave the slot empty");
assert!(
!crate::model_relay::ca::ca_dir(&crate::config::openlatch_dir()).exists(),
"no ProxyEnv host must mean no CA is ever generated"
);
match previous {
Some(v) => std::env::set_var("OPENLATCH_DIR", v),
None => std::env::remove_var("OPENLATCH_DIR"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn proxy_env_agent_is_wired_only_after_both_proofs() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert_eq!(fx.wiring.verdict(AGENT), Verdict::Ok);
assert!(fx.wiring.is_wired(AGENT));
let settings = std::fs::read_to_string(settings_path()).expect("settings file");
assert!(settings.contains(&format!("http://127.0.0.1:{}", fx.port)));
assert!(env_sh_exists());
assert!(marker_exists());
}
#[tokio::test(flavor = "multi_thread")]
async fn an_untrusted_store_leaves_every_delivery_unwritten() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(false, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
match fx.wiring.verdict(AGENT) {
Verdict::Failed(reason) => assert!(
reason.starts_with(crate::model_relay::preflight::CA_REASON_MARKER),
"reason must carry the CA marker: {reason}"
),
other => panic!("expected Failed, got {other:?}"),
}
assert!(!fx.wiring.is_wired(AGENT));
assert_eq!(std::fs::read_to_string(settings_path()).unwrap(), SEED);
assert!(!env_sh_exists());
assert!(!marker_exists());
}
#[tokio::test(flavor = "multi_thread")]
async fn a_lone_foreign_settings_key_fails_the_verdict_with_the_marker() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let foreign_seed = "{\n \"fake.proxy\": \"http://corp:3128\"\n}\n";
seed_settings(foreign_seed);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
wire_agents(
vec![fx.agent(&[HOST], FOREIGN_ONLY_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
match fx.wiring.verdict(AGENT) {
Verdict::Failed(reason) => {
assert!(
reason.starts_with(crate::model_relay::preflight::SETTINGS_REASON_MARKER)
);
assert!(reason.contains(crate::model_relay::preflight::SETTINGS_FOREIGN));
}
other => panic!("expected Failed, got {other:?}"),
}
assert!(!fx.wiring.is_wired(AGENT));
assert_eq!(
std::fs::read_to_string(settings_path()).unwrap(),
foreign_seed
);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_garbage_ca_releases_what_was_written() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert_eq!(fx.wiring.verdict(AGENT), Verdict::Ok);
assert!(fx.wiring.is_wired(AGENT));
let ca_pem = crate::model_relay::ca::ca_pem_path(&crate::model_relay::ca::ca_dir(
&crate::config::openlatch_dir(),
));
std::fs::write(&ca_pem, "garbage").expect("corrupt ca.pem");
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
true,
&cfg,
&fx.wiring,
&formats,
)
.await;
match fx.wiring.verdict(AGENT) {
Verdict::Failed(reason) => assert!(
reason.starts_with(crate::model_relay::preflight::CA_REASON_MARKER),
"reason must carry the CA marker: {reason}"
),
other => panic!("expected Failed, got {other:?}"),
}
assert!(!fx.wiring.is_wired(AGENT));
assert_eq!(std::fs::read_to_string(settings_path()).unwrap(), SEED);
assert!(!env_sh_exists());
assert!(!marker_exists());
let rec = settings_record(AGENT, KEY)
.expect("the settings record must survive as a tombstone");
assert_eq!(rec.released_by, Some(ReleasedBy::Wiring));
}
#[tokio::test(flavor = "multi_thread")]
async fn a_failed_later_delivery_releases_the_earlier_ones() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
let env_sh = proxy_env_file::env_sh_path(&crate::config::openlatch_dir());
std::fs::create_dir_all(&env_sh).expect("plant a directory at env.sh's path");
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert!(!fx.wiring.is_wired(AGENT));
assert_eq!(
std::fs::read_to_string(settings_path()).unwrap(),
SEED,
"the settings key a later delivery's hard failure left behind must be released"
);
let rec = settings_record(AGENT, KEY)
.expect("the settings record must survive as a tombstone");
assert_eq!(rec.released_by, Some(ReleasedBy::Wiring));
assert!(!marker_exists());
}
#[tokio::test(flavor = "multi_thread")]
async fn teardown_releases_every_proxy_value() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
let agent = fx.agent(&[HOST], FAKE_DELIVERY);
let binding = agent.binding.clone();
wire_agents(vec![agent], fx.port, false, &cfg, &fx.wiring, &formats).await;
assert!(fx.wiring.is_wired(AGENT));
unwire_model_relay_config(&cfg, &*binding);
assert_eq!(std::fs::read_to_string(settings_path()).unwrap(), SEED);
assert!(!env_sh_exists());
assert!(!marker_exists());
let rec = settings_record(AGENT, KEY)
.expect("the settings record must survive as a tombstone");
assert_eq!(rec.released_by, Some(ReleasedBy::Teardown));
}
#[tokio::test(flavor = "multi_thread")]
async fn release_proxy_delivery_records_who_released() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
let agent = fx.agent(&[HOST], FAKE_DELIVERY);
let binding = agent.binding.clone();
wire_agents(vec![agent], fx.port, false, &cfg, &fx.wiring, &formats).await;
release_proxy_delivery(&*binding, ReleasedBy::Wiring);
let rec = settings_record(AGENT, KEY).expect("record");
assert_eq!(rec.released_by, Some(ReleasedBy::Wiring));
let ledger_path = crate::config::openlatch_dir().join("model-relay-endpoints.json");
let bytes_before = std::fs::read(&ledger_path).expect("ledger bytes");
release_proxy_delivery(&*binding, ReleasedBy::Wiring);
let bytes_after = std::fs::read(&ledger_path).expect("ledger bytes");
assert_eq!(bytes_before, bytes_after, "a second release is a no-op");
}
#[tokio::test(flavor = "multi_thread")]
async fn no_machine_global_write_from_a_relocated_instance() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert!(fx.wiring.is_wired(AGENT));
for rc in [".zshrc", ".bashrc", ".bash_profile", ".profile"] {
assert!(
!fx.home_dir.path().join(rc).exists(),
"{rc} must not be created by a relocated (non-machine-owning) instance"
);
}
assert_eq!(fx.fake_store.install_calls.load(Ordering::SeqCst), 0);
}
#[test]
fn machine_global_gates_need_both_predicates() {
assert!(profile_line_wanted(ENV_FILE_ONLY_DELIVERY, true));
assert!(!profile_line_wanted(ENV_FILE_ONLY_DELIVERY, false));
assert!(!profile_line_wanted(FOREIGN_ONLY_DELIVERY, true));
assert!(!profile_line_wanted(FOREIGN_ONLY_DELIVERY, false));
}
#[test]
fn env_trust_only_is_an_isolated_instance_that_owns_the_wiring() {
let rows = [
((true, false, false), true),
((true, true, false), false),
((true, true, true), false),
((true, false, true), false),
((false, false, false), false),
((false, true, false), false),
((false, false, true), false),
((false, true, true), false),
];
for ((w, m, d), want) in rows {
assert_eq!(
env_trust_only_from(w, m, d),
want,
"wiring={w} machine={m} default_port={d}"
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn an_isolated_instance_wires_the_wrapper_without_the_store() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(false, Ordering::SeqCst);
let mut cfg = fx.config();
cfg.model_relay.port = fx.port;
assert!(!crate::supervision::owns_machine_supervision());
let formats = crate::cloud::worker::SourceFormats::new();
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert_eq!(fx.wiring.verdict(AGENT), Verdict::Ok);
assert!(fx.wiring.is_wired(AGENT));
assert!(env_sh_exists());
assert!(marker_exists());
assert_eq!(
std::fs::read_to_string(settings_path()).unwrap(),
SEED,
"the store-trusted settings key is never delivered from an isolated instance"
);
assert_eq!(fx.fake_store.install_calls.load(Ordering::SeqCst), 0);
}
#[tokio::test(flavor = "multi_thread")]
async fn ensure_proxy_env_trust_runs_only_when_needed() {
let fx = ProxyEnvFixture::new().await;
let cfg = fx.config();
let envvars_only = vec![env_vars_agent("envvars-agent")];
assert!(ensure_proxy_env_trust(&cfg, &envvars_only, true).is_none());
assert_eq!(fx.fake_store.install_calls.load(Ordering::SeqCst), 0);
let proxy_agents = vec![fx.agent(&[HOST], FAKE_DELIVERY)];
assert!(ensure_proxy_env_trust(&cfg, &proxy_agents, false).is_none());
assert_eq!(fx.fake_store.install_calls.load(Ordering::SeqCst), 0);
let outcome = ensure_proxy_env_trust(&cfg, &proxy_agents, true);
assert_eq!(outcome, Some(InstallOutcome::Installed));
assert_eq!(fx.fake_store.install_calls.load(Ordering::SeqCst), 1);
let dir = crate::model_relay::ca::ca_dir(&crate::config::openlatch_dir());
assert!(crate::model_relay::ca::ca_pem_path(&dir).is_file());
let sha = crate::model_relay::ca::LocalCa::load_or_generate(&dir)
.expect("the CA now exists")
.sha256_hex();
assert!(fx.fake_store.is_trusted(&sha).expect("is_trusted"));
let second = ensure_proxy_env_trust(&cfg, &proxy_agents, true);
assert_eq!(second, Some(InstallOutcome::AlreadyTrusted));
}
#[test]
fn ca_expiry_gate_is_seven_days() {
let now = time::OffsetDateTime::now_utc();
assert!(
ca_too_close_to_expiry(now + time::Duration::days(6), now),
"6 days left must be too close"
);
assert!(
ca_too_close_to_expiry(
now + time::Duration::days(7) + time::Duration::hours(23),
now
),
"7 days 23 hours left is still whole_days() == 7 — 04's own boundary"
);
assert!(
!ca_too_close_to_expiry(now + time::Duration::days(8), now),
"8 days left must not be too close"
);
}
#[test]
#[ignore]
#[cfg(target_os = "macos")]
fn owner_acceptance_real_store_on_real_home() {
fn run(program: &str, args: &[&str]) -> (bool, String) {
let out = std::process::Command::new(program)
.args(args)
.output()
.unwrap_or_else(|e| panic!("could not run {program}: {e}"));
(
out.status.success(),
String::from_utf8_lossy(&out.stdout).into_owned(),
)
}
fn acceptance_settings_file() -> Option<std::path::PathBuf> {
std::env::var("OPENLATCH_ACCEPT_SETTINGS")
.ok()
.map(std::path::PathBuf::from)
}
if std::env::var("OPENLATCH_OWNER_ACCEPTANCE").as_deref() != Ok("1") {
panic!(
"refuses: set OPENLATCH_OWNER_ACCEPTANCE=1 to run this owner acceptance check"
);
}
let (ok, keychains) = run("security", &["list-keychains", "-d", "user"]);
if !ok || !keychains.contains("login.keychain-db") {
panic!(
"refuses: `security list-keychains -d user` does not name login.keychain-db \
— this looks like a sandbox HOME"
);
}
let (ok, manager) = run("launchctl", &["managername"]);
if !ok || manager.trim() != "Aqua" {
panic!("refuses: no desktop session (`launchctl managername` != Aqua)");
}
let settings_path = acceptance_settings_file().unwrap_or_else(|| {
panic!("refuses: set OPENLATCH_ACCEPT_SETTINGS to an absolute settings.json path")
});
if !settings_path.is_absolute() {
panic!("refuses: OPENLATCH_ACCEPT_SETTINGS must be an absolute path");
}
let snapshot = std::fs::read_to_string(&settings_path).unwrap_or_else(|e| {
panic!("refuses: could not read {}: {e}", settings_path.display())
});
let parsed = jsonc_parser::parse_to_ast(
&snapshot,
&jsonc_parser::CollectOptions {
comments: jsonc_parser::CommentCollectionStrategy::Off,
tokens: false,
},
&jsonc_parser::ParseOptions::default(),
)
.unwrap_or_else(|e| {
panic!(
"refuses: {} does not parse as JSONC: {e}",
settings_path.display()
)
});
match parsed.value {
Some(jsonc_parser::ast::Value::Object(o)) if !o.properties.is_empty() => {}
_ => panic!(
"refuses: {} must parse as a non-empty JSON object",
settings_path.display()
),
}
let _dir_lock = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let openlatch_dir = tempfile::tempdir().expect("tempdir");
let previous_dir = std::env::var_os("OPENLATCH_DIR");
std::env::set_var("OPENLATCH_DIR", openlatch_dir.path());
let _trust_guard = trust_store::install_for_tests(trust_store::platform());
let rt = tokio::runtime::Runtime::new().expect("tokio runtime");
rt.block_on(async {
let host = "marketplace.visualstudio.com";
let wiring = Arc::new(WiringState::default());
let state = Arc::new(
ModelRelayState::new(
reqwest::Url::parse(&format!("https://{host}")).unwrap(),
0,
8,
&[],
)
.with_wiring(wiring.clone())
.with_intercept(wiring.intercept())
.with_refusals(wiring.refusals()),
);
let port = serve_ephemeral(state).await;
let binding: Arc<dyn crate::hooks::binding::AgentBinding> = Arc::new(FakeBinding {
agent_type: "vscode-acceptance",
model_relay_wiring: Some(proxy_env_wiring(
&["marketplace.visualstudio.com"],
&[ProxyDelivery::SettingsKey {
file: acceptance_settings_file,
key: "http.proxy",
}],
)),
..Default::default()
});
let agents = vec![crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: binding.clone(),
}];
let cfg = proxy_env_test_config("agt_owner_acceptance");
let dir = crate::model_relay::ca::ca_dir(&crate::config::openlatch_dir());
let sha = crate::model_relay::ca::LocalCa::load_or_generate(&dir)
.expect("load or generate the CA")
.sha256_hex();
match ensure_proxy_env_trust(&cfg, &agents, true) {
Some(InstallOutcome::Installed) | Some(InstallOutcome::AlreadyTrusted) => {}
other => panic!("A-1: expected Installed or AlreadyTrusted, got {other:?}"),
}
assert!(
trust_store::store().is_trusted(&sha).expect("is_trusted"),
"A-1 read-back"
);
println!("A-1: the relay CA is trusted in the login keychain");
let formats = crate::cloud::worker::SourceFormats::new();
wire_agents(agents, port, false, &cfg, &wiring, &formats).await;
assert_eq!(
wiring.verdict("vscode-acceptance"),
Verdict::Ok,
"A-2: the tick must prove green"
);
let written = std::fs::read_to_string(&settings_path).expect("read settings after A-2");
assert!(
written.contains(&format!("\"http.proxy\": \"http://127.0.0.1:{port}\"")),
"A-2: the settings file must now carry our proxy URL"
);
println!("A-2: {} now names http://127.0.0.1:{port}", settings_path.display());
let hold_secs: u64 = std::env::var("OPENLATCH_ACCEPT_HOLD_SECS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(180);
println!(
"Launch VS Code now: code --user-data-dir \"$X\" --extensions-dir \"$Y\" --new-window"
);
println!(
"Open the Extensions view and search. Holding {hold_secs}s (relay log with \
RUST_LOG=info --nocapture should show intercepted CONNECTs for {host} with \
completed handshakes and no \"tls handshake eof\")."
);
std::thread::sleep(std::time::Duration::from_secs(hold_secs));
release_proxy_delivery(&*binding, ReleasedBy::Teardown);
let after_release =
std::fs::read_to_string(&settings_path).expect("read settings after A-3");
assert_eq!(
after_release, snapshot,
"A-3: release must restore the pre-test bytes exactly"
);
println!(
"A-3: released at {:?}. Holding 60s more with the relay still up — search \
again WITHOUT reloading and record whether new CONNECT lines keep arriving \
(reload needed) or not (live pickup).",
std::time::SystemTime::now()
);
std::thread::sleep(std::time::Duration::from_secs(60));
assert!(
trust_store::store().remove(&sha).expect("remove"),
"A-4: remove"
);
assert!(
!trust_store::store().is_trusted(&sha).expect("is_trusted"),
"A-4: read-back"
);
println!("A-4: the relay CA has been removed from the login keychain");
});
match previous_dir {
Some(v) => std::env::set_var("OPENLATCH_DIR", v),
None => std::env::remove_var("OPENLATCH_DIR"),
}
}
impl ProxyEnvFixture {
async fn with_upstream(upstream: u16) -> Self {
Self::build(upstream).await
}
}
fn lane_endpoints(
dir: &std::path::Path,
url: &str,
) -> (
std::path::PathBuf,
&'static dyn crate::hooks::provider_endpoints::ProviderEndpoints,
) {
let path = dir.join("globalState.json");
write_lane(&path, url);
let lanes = crate::hooks::cline_providers::state_lanes_from(
Some(path.clone()),
Some(path.clone()),
);
let ep: &'static crate::hooks::cline_providers::ClineProviderEndpoints =
Box::leak(Box::new(
crate::hooks::cline_providers::ClineProviderEndpoints::at(lanes, None),
));
(
path,
ep as &'static dyn crate::hooks::provider_endpoints::ProviderEndpoints,
)
}
fn write_lane(path: &std::path::Path, url: &str) {
std::fs::write(path, format!("{{\"openAiBaseUrl\": \"{url}\"}}\n")).expect("lane");
}
fn agent_with_endpoints(
hosts: &'static [&'static str],
delivery: &'static [ProxyDelivery],
endpoints: &'static dyn crate::hooks::provider_endpoints::ProviderEndpoints,
) -> crate::hooks::DetectedAgent {
crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: Arc::new(FakeBinding {
agent_type: AGENT,
display_name: "Fake",
config_dir: crate::config::openlatch_dir().join("fake-agent"),
model_relay_wiring: Some(proxy_env_wiring(hosts, delivery)),
provider_endpoints: Some(endpoints),
..Default::default()
}),
}
}
fn interceptor_of(wiring: &WiringState) -> Arc<crate::model_relay::ca::Interceptor> {
wiring
.intercept()
.read()
.expect("slot")
.clone()
.expect("the tick filled the interceptor")
}
#[tokio::test(flavor = "multi_thread")]
async fn a_proxy_env_agent_records_no_wire_format() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
fx.wiring.set_wired("claude-code", true);
fx.wiring
.set_wired_format("claude-code", WireFormat::AnthropicMessages);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert_eq!(
fx.wiring.verdict(AGENT),
Verdict::Ok,
"premise: a green tick"
);
assert!(fx.wiring.is_wired(AGENT));
assert_eq!(fx.wiring.wired_format(AGENT), None);
assert_eq!(
fx.wiring
.sole_wired_agent_for(WireFormat::AnthropicMessages),
Some("claude-code")
);
assert!(!fx.wiring.intercept_proven(AGENT));
}
#[tokio::test(flavor = "multi_thread")]
async fn an_unchanged_tick_keeps_the_intercept_proof() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
async fn tick(
fx: &ProxyEnvFixture,
cfg: &Config,
formats: &crate::cloud::worker::SourceFormats,
broke: bool,
) {
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
broke,
cfg,
&fx.wiring,
formats,
)
.await;
}
tick(&fx, &cfg, &formats, false).await;
assert!(fx.wiring.is_wired(AGENT), "premise: tick 1 is green");
fx.wiring.note_intercepted(HOST);
assert!(fx.wiring.intercept_proven(AGENT));
let written = std::fs::read(settings_path()).expect("settings");
tick(&fx, &cfg, &formats, false).await;
assert_eq!(std::fs::read(settings_path()).expect("settings"), written);
assert!(fx.wiring.intercept_proven(AGENT), "tick 2: nothing changed");
tick(&fx, &cfg, &formats, true).await;
assert_eq!(fx.wiring.verdict(AGENT), Verdict::Ok);
assert!(
fx.wiring.intercept_proven(AGENT),
"tick 3: a green re-probe of a wired agent writes nothing and keeps the proof"
);
fx.wiring.set_wired(AGENT, false);
tick(&fx, &cfg, &formats, false).await;
assert!(
fx.wiring.is_wired(AGENT),
"tick 4 re-establishes the wiring"
);
assert!(
!fx.wiring.intercept_proven(AGENT),
"a new wiring starts unproven"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn the_tick_publishes_static_and_dynamic_hosts() {
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
let lane_dir = tempfile::tempdir().expect("lane");
let (lane, ep) = lane_endpoints(lane_dir.path(), "https://qwen.gpu.test/v1");
let agent = agent_with_endpoints(&[HOST], FAKE_DELIVERY, ep);
wire_agents(
vec![agent.clone()],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert_eq!(fx.wiring.sole_interceptor_for(HOST), Some(AGENT));
assert_eq!(fx.wiring.sole_interceptor_for("qwen.gpu.test"), Some(AGENT));
assert!(interceptor_of(&fx.wiring).intercepts("qwen.gpu.test"));
write_lane(&lane, "http://127.0.0.1:8123/v1");
publish_intercept_hosts(&fx.wiring, std::slice::from_ref(&agent));
assert_eq!(fx.wiring.sole_interceptor_for("qwen.gpu.test"), None);
assert!(!interceptor_of(&fx.wiring).intercepts("qwen.gpu.test"));
assert_eq!(fx.wiring.sole_interceptor_for(HOST), Some(AGENT));
assert!(interceptor_of(&fx.wiring).intercepts(HOST));
wire_agents(
vec![env_vars_agent("envvars-agent")],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert_eq!(fx.wiring.sole_interceptor_for(HOST), None);
}
#[tokio::test(flavor = "multi_thread")]
async fn the_store_is_asked_before_any_network_probe() {
let dead = mock::closed_port().await;
let fx = ProxyEnvFixture::with_upstream(dead).await;
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
fx.fake_store.verify.store(false, Ordering::SeqCst);
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
match fx.wiring.verdict(AGENT) {
Verdict::Failed(reason) => {
assert!(
reason.starts_with(crate::model_relay::preflight::CA_REASON_MARKER),
"{reason}"
);
assert!(
reason.contains("does not trust the relay's CA"),
"the store answered first, before any probe: {reason}"
);
}
other => panic!("expected Failed, got {other:?}"),
}
fx.fake_store.verify.store(true, Ordering::SeqCst);
wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
match fx.wiring.verdict(AGENT) {
Verdict::Failed(reason) => {
assert!(
!reason.contains("does not trust"),
"the store is not the failure: {reason}"
);
assert!(
!reason.starts_with(crate::model_relay::preflight::CA_REASON_MARKER),
"an unreachable upstream is not a CA refusal: {reason}"
);
}
other => panic!("expected Failed, got {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn a_store_refusal_is_not_a_network_failure() {
let dead = mock::closed_port().await;
let fx = ProxyEnvFixture::with_upstream(dead).await;
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
fx.fake_store.verify.store(false, Ordering::SeqCst);
let failed = wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert!(!failed, "a CA refusal is local: no backoff");
fx.fake_store.verify.store(true, Ordering::SeqCst);
let failed = wire_agents(
vec![fx.agent(&[HOST], FAKE_DELIVERY)],
fx.port,
false,
&cfg,
&fx.wiring,
&formats,
)
.await;
assert!(failed, "control: a dead upstream still backs off");
}
fn fake_b_settings_file() -> Option<std::path::PathBuf> {
Some(
crate::config::openlatch_dir()
.join("fake-b")
.join("settings.json"),
)
}
#[tokio::test(flavor = "multi_thread")]
async fn a_tripped_dynamic_host_releases_its_agent() {
use crate::model_relay::ca_lifecycle::refused_reason;
use crate::model_relay::intercept::RefusalKind;
const B_DELIVERY: &[ProxyDelivery] = &[ProxyDelivery::SettingsKey {
file: fake_b_settings_file,
key: KEY,
}];
let fx = ProxyEnvFixture::new().await;
fx.fake_store.verify.store(true, Ordering::SeqCst);
let cfg = fx.config();
let formats = crate::cloud::worker::SourceFormats::new();
let lane_dir = tempfile::tempdir().expect("lane");
let (_lane, ep) = lane_endpoints(lane_dir.path(), "https://qwen.gpu.test/v1");
let a = agent_with_endpoints(&[HOST], FAKE_DELIVERY, ep);
let b_file = fake_b_settings_file().expect("path");
std::fs::create_dir_all(b_file.parent().expect("parent")).expect("mkdir");
std::fs::write(&b_file, SEED).expect("seed b");
let b = crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::ClaudeCode,
binding: Arc::new(FakeBinding {
agent_type: "fake-b",
display_name: "Fake B",
config_dir: crate::config::openlatch_dir().join("fake-b"),
model_relay_wiring: Some(proxy_env_wiring(&["b.test"], B_DELIVERY)),
..Default::default()
}),
};
let agents = vec![a, b];
wire_agents(agents.clone(), fx.port, false, &cfg, &fx.wiring, &formats).await;
assert!(fx.wiring.is_wired(AGENT), "premise: fake-a is wired");
assert!(fx.wiring.is_wired("fake-b"), "premise: fake-b is wired");
assert!(env_sh_exists());
let refusals = fx.wiring.refusals();
for _ in 0..3 {
refusals.record_refusal("qwen.gpu.test", RefusalKind::HandshakeEof);
}
let mut holds = std::collections::BTreeMap::new();
let at = TrustClock {
wall: time::OffsetDateTime::now_utc(),
mono: std::time::Instant::now(),
};
let skip =
react_to_trust_loss(&cfg, &agents, None, at, &refusals, &mut holds, &fx.wiring);
assert!(skip.contains(AGENT), "the dynamic host names fake-a");
assert!(!fx.wiring.is_wired(AGENT));
assert_eq!(
fx.wiring.verdict(AGENT),
Verdict::Failed(refused_reason("qwen.gpu.test"))
);
assert_eq!(std::fs::read_to_string(settings_path()).unwrap(), SEED);
assert!(!env_sh_exists());
assert!(!skip.contains("fake-b"));
assert!(fx.wiring.is_wired("fake-b"), "fake-b is untouched");
}
}
}