use std::collections::{BTreeMap, BTreeSet};
use std::sync::{Arc, Mutex};
use crate::config::Config;
use crate::core::cloud::tamper::{
FieldDelta, TamperEvent, DETECTION_RELAY_URL_REVERTED, DETECTION_RELAY_WIRING_MISCONFIGURED,
DETECTION_RELAY_WIRING_PENDING, RELAY_WIRING_HOOK_EVENT,
};
use crate::core::fs_secure::{fingerprint, FileFingerprint};
use crate::error::{
ERR_MODEL_RELAY_ENDPOINT_PORTS, ERR_MODEL_RELAY_PREFLIGHT_FAILED, ERR_MODEL_RELAY_STATE_FILE,
};
use crate::hooks::atomic::RewriteOutcome;
use crate::hooks::cline_providers::{
decide, relay_value, DecideCtx, Decision, SlotId, SlotObservation,
};
use crate::hooks::model_relay_endpoints::{
self as records, EndpointRecord, Proof, ReleasedBy, SlotState, SlotValue,
};
use crate::hooks::provider_endpoints::{names_port, now_unix, ProviderEndpoints};
use crate::model_relay::endpoints::{
endpoint_port_block, normalize_origin, EndpointListeners, EndpointSpec, Origin, RelayPorts,
};
use crate::model_relay::preflight::{self, EndpointVerdict, WiringState};
use crate::model_relay::session::SessionRegistry;
use crate::model_relay::wire_format::WireFormat;
use super::reconciler::TamperSinks;
const REAPPLY_MIN_GAP: std::time::Duration = std::time::Duration::from_secs(10);
const CONTESTED_GAP: std::time::Duration = std::time::Duration::from_secs(60);
const CONTESTED_REVERTS: usize = 20;
const CONTESTED_WINDOW: std::time::Duration = std::time::Duration::from_secs(600);
#[derive(Debug, Default)]
pub(crate) struct ReapplyGate {
recent: std::collections::VecDeque<std::time::Instant>,
}
impl ReapplyGate {
pub(crate) fn allows(&self, now: std::time::Instant) -> bool {
let gap = if self.contested(now) {
CONTESTED_GAP
} else {
REAPPLY_MIN_GAP
};
self.recent
.back()
.is_none_or(|last| now.saturating_duration_since(*last) >= gap)
}
pub(crate) fn note(&mut self, now: std::time::Instant) {
self.recent.push_back(now);
while self
.recent
.front()
.is_some_and(|t| now.saturating_duration_since(*t) > CONTESTED_WINDOW)
{
self.recent.pop_front();
}
}
pub(crate) fn contested(&self, now: std::time::Instant) -> bool {
self.recent
.iter()
.filter(|t| now.saturating_duration_since(**t) <= CONTESTED_WINDOW)
.count()
>= CONTESTED_REVERTS
}
}
pub(crate) struct EndpointWiring {
listeners: Arc<EndpointListeners>,
wiring: Arc<WiringState>,
ports: RelayPorts,
started_at: u64,
sinks: TamperSinks,
written: Mutex<BTreeMap<std::path::PathBuf, FileFingerprint>>,
squatted: Mutex<BTreeSet<u16>>,
reclaimed: std::sync::atomic::AtomicBool,
reapplies: Mutex<BTreeMap<String, ReapplyGate>>,
sightings: Arc<SessionRegistry>,
}
impl EndpointWiring {
pub(crate) fn new(
listeners: Arc<EndpointListeners>,
wiring: Arc<WiringState>,
ports: RelayPorts,
sinks: TamperSinks,
sightings: Arc<SessionRegistry>,
) -> Self {
Self {
listeners,
wiring,
ports,
started_at: now_unix(),
sinks,
written: Mutex::new(BTreeMap::new()),
squatted: Mutex::new(BTreeSet::from([ports.daemon])),
reclaimed: std::sync::atomic::AtomicBool::new(false),
reapplies: Mutex::new(BTreeMap::new()),
sightings,
}
}
fn written(&self) -> std::sync::MutexGuard<'_, BTreeMap<std::path::PathBuf, FileFingerprint>> {
self.written.lock().unwrap_or_else(|e| e.into_inner())
}
fn squatted(&self) -> BTreeSet<u16> {
self.squatted
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone()
}
fn pick_port(
&self,
key: &str,
block: &std::ops::RangeInclusive<u16>,
usable: &[(String, EndpointRecord)],
) -> Option<u16> {
let squatted = self.squatted();
records::allocate_port(key, block.clone(), usable)
.filter(|p| !squatted.contains(p))
.or_else(|| {
block.clone().find(|p| {
!squatted.contains(p)
&& !usable.iter().any(|(_, r)| r.port == *p && r.is_live())
})
})
}
pub(crate) async fn reconcile(
&self,
config: &Config,
agents: &[crate::hooks::DetectedAgent],
) -> bool {
let mut any_failed = false;
for agent in agents {
let Some(endpoints) = agent.binding.provider_endpoints() else {
continue;
};
if !super::owns_wiring_for(config, &*agent.binding) {
continue;
}
any_failed |= self.reconcile_agent(endpoints).await;
}
any_failed
}
pub(crate) fn watch_files(
&self,
config: &Config,
agents: &[crate::hooks::DetectedAgent],
) -> Vec<std::path::PathBuf> {
let mut files: Vec<std::path::PathBuf> = agents
.iter()
.filter(|agent| super::owns_wiring_for(config, &*agent.binding))
.filter_map(|agent| agent.binding.provider_endpoints())
.flat_map(|endpoints| endpoints.watch_files())
.collect();
files.sort();
files.dedup();
files
}
async fn reconcile_agent(&self, endpoints: &'static dyn ProviderEndpoints) -> bool {
if !self
.reclaimed
.swap(true, std::sync::atomic::Ordering::Relaxed)
{
endpoints.reclaim_retired(self.ports.main);
}
let prefix = endpoints.record_prefix();
let recorded = match records::endpoint_records(prefix) {
Ok(recorded) => recorded,
Err(e) => {
tracing::warn!(
agent = endpoints.agent_type(),
code = %e.code,
error = %e.message,
"endpoint records unreadable — provider slots are left as they are"
);
return true;
}
};
for (key, rec) in &recorded {
if rec.state == SlotState::Released && rec.released_by == Some(ReleasedBy::Teardown) {
continue;
}
self.serve(endpoints, key, rec).await;
}
let keys: BTreeSet<String> = recorded.iter().map(|(k, _)| k.clone()).collect();
let observation = endpoints.observe(&keys);
for problem in &observation.problems {
tracing::warn!(
?problem,
"a Cline state file cannot be read safely — its slots are left alone"
);
}
let ctx = DecideCtx {
ports: self.ports,
started_at: self.started_at,
};
let mut recorded = recorded;
let mut any_failed = false;
for obs in &observation.slots {
let key = obs.slot.record_key();
let rec = recorded
.iter()
.find(|(k, _)| *k == key)
.map(|(_, r)| r.clone());
let rec = rec.as_ref();
let failed = match decide(obs, rec, &ctx) {
Decision::Keep => {
if let Some(rec) = rec {
self.keep(endpoints, &key, rec, obs);
}
false
}
Decision::Skip => false,
Decision::Uncovered(_) => {
self.wiring.set_endpoint_verdict(&key, None);
false
}
Decision::Release { restore } => {
if let Some(rec) = rec {
self.release(endpoints, &key, rec, obs, restore).await;
}
false
}
Decision::Reapply => match rec {
Some(rec) => self.reapply(endpoints, &key, rec, obs),
None => false,
},
Decision::Wire {
path,
origin,
prior,
}
| Decision::NewPrior {
path,
origin,
prior,
} => {
self.wire(
endpoints,
&key,
rec,
&mut recorded,
obs,
&path,
origin,
prior,
)
.await
}
};
any_failed |= failed;
}
any_failed
}
async fn serve(&self, endpoints: &dyn ProviderEndpoints, key: &str, rec: &EndpointRecord) {
if self.squatted().contains(&rec.port) {
self.wiring.set_endpoint_verdict(
key,
Some(EndpointVerdict {
code: ERR_MODEL_RELAY_ENDPOINT_PORTS,
detail: format!("loopback port {} is held by another process", rec.port),
}),
);
return;
}
let Some(origin) = parse_origin(&rec.origin, &self.ports) else {
return;
};
let spec = EndpointSpec {
key: key.to_string(),
agent: endpoints.agent_type(),
family: rec.family.as_deref().and_then(family_of),
port: rec.port,
origin,
};
if let Err(e) = self.listeners.ensure(spec).await {
self.wiring.set_endpoint_verdict(
key,
Some(EndpointVerdict {
code: ERR_MODEL_RELAY_ENDPOINT_PORTS,
detail: e.message.clone(),
}),
);
self.squatted
.lock()
.unwrap_or_else(|e| e.into_inner())
.insert(rec.port);
if rec.is_live() {
let Some(slot) = SlotId::from_record_key(key) else {
return;
};
let _ = records::update_endpoint(key, |r| {
r.state = SlotState::Released;
r.released_by = Some(ReleasedBy::Wiring);
r.changed_at = now_unix();
});
let port = rec.port;
if let Ok(RewriteOutcome::Written) =
endpoints.write_slot(&rec.file, &slot, &rec.prior, &|cur| names_port(cur, port))
{
self.remember_write(&rec.file);
}
}
}
}
fn keep(
&self,
endpoints: &dyn ProviderEndpoints,
key: &str,
rec: &EndpointRecord,
obs: &SlotObservation,
) {
self.wiring.set_endpoint_verdict(key, None);
if rec.state == SlotState::Pending {
let _ = records::update_endpoint(key, |r| r.state = SlotState::Wired);
}
let now = now_unix();
if self
.listeners
.get(key)
.is_some_and(|ep| rec.served_since_written(ep.last_request_unix()))
{
if rec.proven_by != Some(Proof::Traffic) || rec.misconfigured_event.is_some() {
self.prove(endpoints, key, obs, Proof::Traffic, now);
}
return;
}
if rec.proven_at.is_none() {
if obs.traffic_only().is_some() {
return;
}
let saved_by_editor = {
let written = self.written();
match (written.get(&obs.file), fingerprint(&obs.file)) {
(Some(ours), Ok(current)) => *ours != current,
_ => false,
}
};
if saved_by_editor {
self.prove(endpoints, key, obs, Proof::EditorSave, now);
}
return;
}
if rec.proven_by != Some(Proof::EditorSave)
|| rec.misconfigured_event.is_some()
|| !obs.hooks_name_provider()
{
return;
}
let loaded_at = rec.proven_at.unwrap_or(now);
let sightings: Vec<u64> = obs
.provider_ids()
.iter()
.flat_map(|id| {
self.sightings
.provider_sightings(endpoints.agent_type(), id)
})
.collect();
let Some(seen) = misconfigured_since(loaded_at, &sightings, now) else {
return;
};
let event = self.event(
endpoints,
key,
&obs.file,
obs,
DETECTION_RELAY_WIRING_MISCONFIGURED,
);
self.sinks.publish_detected(&event);
let _ = records::update_endpoint(key, |r| {
r.misconfigured_event = Some(event.tamper.event_id.clone());
});
tracing::warn!(
agent = endpoints.agent_type(),
key,
loaded_at,
seen,
"the editor uses this provider and loaded its relay URL, yet no request reached the endpoint — something overrides it"
);
}
fn prove(
&self,
endpoints: &dyn ProviderEndpoints,
key: &str,
obs: &SlotObservation,
proof: Proof,
now: u64,
) {
let updated = records::update_endpoint(key, |r| {
r.proven_at.get_or_insert(now);
r.proven_by = Some(proof);
});
tracing::info!(
agent = endpoints.agent_type(),
key,
?proof,
"provider slot proven: the editor uses its relay endpoint"
);
let Ok(Some(rec)) = updated else {
return;
};
if let Some(id) = rec.pending_event.as_deref() {
self.heal(
endpoints,
key,
&rec,
obs,
DETECTION_RELAY_WIRING_PENDING,
id,
);
let _ = records::update_endpoint(key, |r| r.pending_event = None);
}
if proof == Proof::Traffic {
if let Some(id) = rec.misconfigured_event.as_deref() {
self.heal(
endpoints,
key,
&rec,
obs,
DETECTION_RELAY_WIRING_MISCONFIGURED,
id,
);
let _ = records::update_endpoint(key, |r| r.misconfigured_event = None);
}
}
}
fn heal(
&self,
endpoints: &dyn ProviderEndpoints,
key: &str,
rec: &EndpointRecord,
obs: &SlotObservation,
method: &str,
id: &str,
) {
let detected = self
.event(endpoints, key, &rec.file, obs, method)
.with_event_id(id);
self.sinks.publish_healed(&TamperEvent::new_healed(
&detected,
"succeeded",
1,
"closed",
));
}
async fn release(
&self,
endpoints: &dyn ProviderEndpoints,
key: &str,
rec: &EndpointRecord,
obs: &SlotObservation,
restore: Option<SlotValue>,
) {
if !obs.entry_present {
self.listeners.release(key).await;
}
let _ = records::update_endpoint(key, |r| {
r.state = SlotState::Released;
r.released_by = Some(ReleasedBy::Wiring);
r.changed_at = now_unix();
});
if let Some(value) = restore {
let port = rec.port;
match endpoints.write_slot(&obs.file, &obs.slot, &value, &|cur| names_port(cur, port)) {
Ok(RewriteOutcome::Written) => self.remember_write(&obs.file),
Ok(_) => {}
Err(e) => {
tracing::warn!(key, code = %e.code, error = %e.message, "could not hand a provider slot back")
}
}
}
self.wiring.set_endpoint_verdict(key, None);
tracing::info!(
agent = endpoints.agent_type(),
key,
"provider slot handed back"
);
}
fn reapply(
&self,
endpoints: &dyn ProviderEndpoints,
key: &str,
rec: &EndpointRecord,
obs: &SlotObservation,
) -> bool {
let Some(ours) = rec.last_written.clone() else {
return false;
};
let now = std::time::Instant::now();
{
let mut gates = self.reapplies.lock().unwrap_or_else(|e| e.into_inner());
let gate = gates.entry(key.to_string()).or_default();
if !gate.allows(now) {
return false;
}
let was_contested = gate.contested(now);
gate.note(now);
let contested = gate.contested(now);
if let Some(endpoint) = self.listeners.get(key) {
endpoint.set_contested(contested);
}
if contested && !was_contested {
tracing::warn!(
agent = endpoints.agent_type(),
key,
"a provider slot keeps being reverted — re-applying it at most once a minute"
);
}
}
let observed = obs.value.clone();
let detected = (rec.state == SlotState::Wired)
.then(|| self.event(endpoints, key, &obs.file, obs, DETECTION_RELAY_URL_REVERTED));
if let Some(event) = &detected {
self.sinks.publish_detected(event);
}
let written_at = now_unix();
match endpoints.write_slot(&obs.file, &obs.slot, &SlotValue::Text(ours), &|cur| {
*cur == observed
}) {
Ok(RewriteOutcome::Written) => {
self.remember_write(&obs.file);
let _ = records::update_endpoint(key, |r| {
r.state = SlotState::Wired;
r.changed_at = written_at;
r.proven_at = None;
r.proven_by = None;
});
if let Some(event) = &detected {
self.sinks.publish_healed(&TamperEvent::new_healed(
event,
"succeeded",
1,
"closed",
));
}
tracing::info!(
agent = endpoints.agent_type(),
key,
"relay URL re-applied to a reverted provider slot"
);
false
}
Ok(_) => false,
Err(e) => {
self.state_file_verdict(key, &e);
if let Some(event) = &detected {
self.sinks
.publish_healed(&TamperEvent::new_healed(event, "failed", 1, "closed"));
}
true
}
}
}
#[allow(clippy::too_many_arguments)]
async fn wire(
&self,
endpoints: &dyn ProviderEndpoints,
key: &str,
rec: Option<&EndpointRecord>,
recorded: &mut Vec<(String, EndpointRecord)>,
obs: &SlotObservation,
path: &str,
origin: Origin,
prior: SlotValue,
) -> bool {
let block = endpoint_port_block(self.ports.main);
let squatted = self.squatted();
let usable: Vec<(String, EndpointRecord)> = recorded
.iter()
.filter(|(k, r)| k != key || !squatted.contains(&r.port))
.cloned()
.collect();
let Some(port) = self.pick_port(key, &block, &usable) else {
self.wiring.set_endpoint_verdict(
key,
Some(EndpointVerdict {
code: ERR_MODEL_RELAY_ENDPOINT_PORTS,
detail: format!(
"every endpoint port in {}–{} is in use",
block.start(),
block.end()
),
}),
);
return true;
};
let family = obs.row.and_then(|r| r.family);
let mut value = relay_value(port, path);
let record = EndpointRecord {
port,
origin: origin.to_string(),
prior,
last_written: Some(value.clone()),
file: obs.file.clone(),
state: SlotState::Pending,
released_by: None,
changed_at: now_unix(),
proven_at: None,
family: family.map(|f| f.as_str().to_string()),
proven_by: None,
pending_event: rec.and_then(|r| r.pending_event.clone()),
misconfigured_event: rec.and_then(|r| r.misconfigured_event.clone()),
};
if let Err(e) = records::put_endpoint(key, record.clone()) {
self.state_file_verdict(key, &e);
return true;
}
match recorded.iter_mut().find(|(k, _)| k == key) {
Some((_, existing)) => *existing = record.clone(),
None => recorded.push((key.to_string(), record.clone())),
}
let mut port = port;
let spec = EndpointSpec {
key: key.to_string(),
agent: endpoints.agent_type(),
family,
port,
origin: origin.clone(),
};
if let Err(first_err) = self.listeners.ensure(spec).await {
self.squatted
.lock()
.unwrap_or_else(|e| e.into_inner())
.insert(port);
let Some(retry_port) = self.pick_port(key, &block, &usable) else {
self.abandon(key, ERR_MODEL_RELAY_ENDPOINT_PORTS, first_err.message);
return true;
};
port = retry_port;
value = relay_value(port, path);
let mut record = record;
record.port = port;
record.last_written = Some(value.clone());
if let Err(e) = records::put_endpoint(key, record.clone()) {
self.state_file_verdict(key, &e);
return true;
}
if let Some((_, existing)) = recorded.iter_mut().find(|(k, _)| k == key) {
*existing = record;
}
let spec = EndpointSpec {
key: key.to_string(),
agent: endpoints.agent_type(),
family,
port,
origin: origin.clone(),
};
if let Err(e) = self.listeners.ensure(spec).await {
self.squatted
.lock()
.unwrap_or_else(|e| e.into_inner())
.insert(port);
self.abandon(key, ERR_MODEL_RELAY_ENDPOINT_PORTS, e.message);
return true;
}
}
let probe_format = family.unwrap_or(WireFormat::Unknown);
if let Err(why) = preflight::probe(
port,
probe_format,
origin.as_url().as_str(),
preflight::PREFLIGHT_TIMEOUT,
)
.await
{
self.abandon(key, ERR_MODEL_RELAY_PREFLIGHT_FAILED, why);
return true;
}
let observed = obs.value.clone();
let written_at = now_unix();
match endpoints.write_slot(&obs.file, &obs.slot, &SlotValue::Text(value), &|cur| {
*cur == observed
}) {
Ok(RewriteOutcome::Written) => {
self.remember_write(&obs.file);
self.wiring.set_endpoint_verdict(key, None);
let event = self.event(
endpoints,
key,
&obs.file,
obs,
DETECTION_RELAY_WIRING_PENDING,
);
self.sinks.publish_detected(&event);
let _ = records::update_endpoint(key, |r| {
r.state = SlotState::Wired;
r.changed_at = written_at;
r.pending_event = Some(event.tamper.event_id.clone());
});
tracing::info!(
agent = endpoints.agent_type(),
key,
port,
origin = %origin,
"provider slot wired to its relay endpoint — the editor picks it up at its next start"
);
false
}
Ok(RewriteOutcome::Contended) => {
self.wiring.set_endpoint_verdict(
key,
Some(EndpointVerdict {
code: ERR_MODEL_RELAY_STATE_FILE,
detail:
"a save landed between the read and the write; retried on the next pass"
.into(),
}),
);
false
}
Ok(RewriteOutcome::Unchanged) => false,
Ok(RewriteOutcome::Absent) => {
self.abandon(
key,
ERR_MODEL_RELAY_STATE_FILE,
"the settings file disappeared".into(),
);
false
}
Err(e) => {
self.abandon(key, e.code, e.message);
true
}
}
}
fn abandon(&self, key: &str, code: &'static str, detail: String) {
let _ = records::update_endpoint(key, |r| {
r.state = SlotState::Released;
r.released_by = Some(ReleasedBy::Wiring);
r.changed_at = now_unix();
});
tracing::warn!(key, code, detail = %detail, "provider slot not wired");
self.wiring
.set_endpoint_verdict(key, Some(EndpointVerdict { code, detail }));
}
fn state_file_verdict(&self, key: &str, e: &crate::error::OlError) {
self.wiring.set_endpoint_verdict(
key,
Some(EndpointVerdict {
code: ERR_MODEL_RELAY_STATE_FILE,
detail: e.message.clone(),
}),
);
}
fn remember_write(&self, file: &std::path::Path) {
if let Ok(now) = fingerprint(file) {
self.written().insert(file.to_path_buf(), now);
}
}
fn event(
&self,
endpoints: &dyn ProviderEndpoints,
key: &str,
file: &std::path::Path,
obs: &SlotObservation,
method: &str,
) -> TamperEvent {
let field = match &obs.slot {
SlotId::GlobalState { key, .. } => key.clone(),
SlotId::ProvidersJson { id } => format!("providers.{id}.settings.baseUrl"),
};
TamperEvent::new(
key.to_string(),
endpoints.agent_type().to_string(),
crate::core::hook_state::hash_settings_path(file),
RELAY_WIRING_HOOK_EVENT.to_string(),
method.to_string(),
)
.with_field_deltas(vec![FieldDelta {
field,
change: "modified".to_string(),
}])
}
}
const MISCONFIGURED_AFTER_SECS: u64 = 120;
fn misconfigured_since(loaded_at: u64, sightings: &[u64], now: u64) -> Option<u64> {
sightings
.iter()
.copied()
.filter(|seen| *seen >= loaded_at && now.saturating_sub(*seen) >= MISCONFIGURED_AFTER_SECS)
.min()
}
fn parse_origin(raw: &str, ports: &RelayPorts) -> Option<Origin> {
normalize_origin(&reqwest::Url::parse(raw).ok()?, ports).ok()
}
fn family_of(name: &str) -> Option<WireFormat> {
WireFormat::ALL.into_iter().find(|f| f.as_str() == name)
}
#[cfg(test)]
mod tests {
use super::*;
fn ephemeral_ports() -> RelayPorts {
let free = || {
std::net::TcpListener::bind(("127.0.0.1", 0))
.and_then(|l| l.local_addr())
.expect("ephemeral port")
.port()
};
RelayPorts {
daemon: free(),
main: free(),
}
}
use std::time::{Duration, Instant};
#[test]
fn a_reapply_is_paced_and_a_contested_slot_slows_down() {
let mut gate = ReapplyGate::default();
let start = Instant::now();
assert!(gate.allows(start));
gate.note(start);
assert!(
!gate.allows(start + Duration::from_secs(5)),
"within the gap"
);
assert!(gate.allows(start + Duration::from_secs(10)));
let mut now = start;
for _ in 1..CONTESTED_REVERTS {
now += Duration::from_secs(10);
gate.note(now);
}
assert!(gate.contested(now));
assert!(!gate.allows(now + Duration::from_secs(30)));
assert!(gate.allows(now + Duration::from_secs(60)));
let later = now + CONTESTED_WINDOW + Duration::from_secs(1);
assert!(!gate.contested(later));
}
#[test]
fn keep_reports_an_editor_routing_around_its_loaded_url_once() {
use crate::hooks::cline_providers::{state_lanes_from, ClineProviderEndpoints};
use crate::hooks::model_relay_endpoints::Proof;
let _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().join("openlatch").into_os_string()),
)]);
let ports = RelayPorts {
daemon: 7500,
main: 7600,
};
let gs = dir.path().join("data").join("globalState.json");
std::fs::create_dir_all(gs.parent().expect("parent")).expect("mkdir");
std::fs::write(
&gs,
r#"{"actModeApiProvider":"gemini","geminiBaseUrl":"http://127.0.0.1:7601"}"#,
)
.expect("write");
let endpoints: &'static ClineProviderEndpoints = Box::leak(Box::new(
ClineProviderEndpoints::at(state_lanes_from(Some(gs.clone()), Some(gs.clone())), None),
));
let key = "cline:gs:shared:geminiBaseUrl";
let now = now_unix();
let record = EndpointRecord {
port: 7601,
origin: "https://generativelanguage.googleapis.com/".into(),
prior: SlotValue::Absent,
last_written: Some("http://127.0.0.1:7601".into()),
file: gs.clone(),
state: SlotState::Wired,
released_by: None,
changed_at: now - 1_000,
proven_at: Some(now - 1_000),
family: None,
pending_event: None,
proven_by: Some(Proof::EditorSave),
misconfigured_event: None,
};
records::put_endpoint(key, record).expect("put");
let (logger, mut events) = crate::logging::tamper_log::TamperLogger::channel();
let registry = Arc::new(SessionRegistry::default());
let factory: crate::model_relay::endpoints::StateFactory =
Arc::new(|_| -> crate::model_relay::ModelRelayState { unreachable!("not served") });
let wiring = EndpointWiring::new(
Arc::new(EndpointListeners::new(factory)),
Arc::new(WiringState::default()),
ports,
TamperSinks {
logger: Some(logger),
cloud_tx: None,
agent_id: String::new(),
client_version: String::new(),
},
registry.clone(),
);
let keep = || {
let rec = records::endpoint_records("cline:")
.expect("records")
.into_iter()
.find(|(k, _)| k == key)
.expect("record")
.1;
let observation = endpoints.observe(&BTreeSet::from([key.to_string()]));
let obs = observation
.slots
.iter()
.find(|s| s.slot.record_key() == key)
.expect("slot");
wiring.keep(endpoints, key, &rec, obs);
};
registry.note_provider("cline", "gemini", now - 2_000);
keep();
assert!(events.try_recv().is_err());
registry.note_provider("cline", "gemini", now - 600);
keep();
let event = events.try_recv().expect("a misconfigured event");
assert_eq!(
event.tamper.detection_method,
DETECTION_RELAY_WIRING_MISCONFIGURED
);
let recorded = records::endpoint_records("cline:").expect("records");
assert_eq!(
recorded[0].1.misconfigured_event.as_deref(),
Some(event.tamper.event_id.as_str())
);
keep();
assert!(events.try_recv().is_err(), "reported once");
}
#[test]
fn a_squatted_port_is_re_picked_within_the_same_pass() {
use crate::hooks::cline_providers::{state_lanes_from, ClineProviderEndpoints};
let _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().join("openlatch").into_os_string()),
)]);
let ports = ephemeral_ports();
let block = endpoint_port_block(ports.main);
let squat_port = *block.start();
let squatter = std::net::TcpListener::bind(("127.0.0.1", squat_port)).expect("squat");
let gs = dir.path().join("data").join("globalState.json");
std::fs::create_dir_all(gs.parent().expect("parent")).expect("mkdir");
std::fs::write(
&gs,
r#"{"actModeApiProvider":"ollama","ollamaBaseUrl":"http://127.0.0.1:9"}"#,
)
.expect("write");
let endpoints: &'static ClineProviderEndpoints = Box::leak(Box::new(
ClineProviderEndpoints::at(state_lanes_from(Some(gs.clone()), Some(gs.clone())), None),
));
let factory: crate::model_relay::endpoints::StateFactory = Arc::new(|_| {
crate::model_relay::ModelRelayState::new(
reqwest::Url::parse("http://127.0.0.1:9").expect("url"),
0,
1,
&[],
)
});
let (logger, _events) = crate::logging::tamper_log::TamperLogger::channel();
let wiring = EndpointWiring::new(
Arc::new(EndpointListeners::new(factory)),
Arc::new(WiringState::default()),
ports,
TamperSinks {
logger: Some(logger),
cloud_tx: None,
agent_id: String::new(),
client_version: String::new(),
},
Arc::new(SessionRegistry::default()),
);
tokio::runtime::Runtime::new()
.expect("runtime")
.block_on(wiring.reconcile_agent(endpoints));
drop(squatter);
let recorded = records::endpoint_records("cline:").expect("records");
let (_, rec) = recorded
.iter()
.find(|(k, _)| k == "cline:gs:shared:ollamaBaseUrl")
.expect("a record for the ollama slot");
assert_ne!(
rec.port, squat_port,
"the squatted port must never be the one recorded as wired"
);
assert!(
block.contains(&rec.port),
"the re-picked port must still be in this slot's own endpoint block, got {}",
rec.port
);
}
#[test]
fn a_squatted_port_is_never_retried_by_serve() {
let _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().join("openlatch").into_os_string()),
)]);
let ports = ephemeral_ports();
let port = *endpoint_port_block(ports.main).start();
let key = "cline:gs:shared:ollamaBaseUrl";
let record = EndpointRecord {
port,
origin: "http://127.0.0.1:9".into(),
prior: SlotValue::Absent,
last_written: Some(format!("http://127.0.0.1:{port}")),
file: dir.path().join("data").join("globalState.json"),
state: SlotState::Released,
released_by: Some(ReleasedBy::Wiring),
changed_at: now_unix(),
proven_at: None,
family: None,
proven_by: None,
pending_event: None,
misconfigured_event: None,
};
records::put_endpoint(key, record.clone()).expect("put");
let factory: crate::model_relay::endpoints::StateFactory =
Arc::new(|_| -> crate::model_relay::ModelRelayState { unreachable!("never served") });
let (logger, _events) = crate::logging::tamper_log::TamperLogger::channel();
let wiring = EndpointWiring::new(
Arc::new(EndpointListeners::new(factory)),
Arc::new(WiringState::default()),
ports,
TamperSinks {
logger: Some(logger),
cloud_tx: None,
agent_id: String::new(),
client_version: String::new(),
},
Arc::new(SessionRegistry::default()),
);
wiring
.squatted
.lock()
.unwrap_or_else(|e| e.into_inner())
.insert(port);
tokio::runtime::Runtime::new()
.expect("runtime")
.block_on(wiring.serve(
&crate::hooks::cline_providers::ClineProviderEndpoints::RESOLVED,
key,
&record,
));
assert!(
wiring.listeners.get(key).is_none(),
"serve() must not bind a port it already knows is squatted"
);
}
#[test]
fn a_provider_in_use_with_a_silent_endpoint_is_misconfigured_only_after_the_window() {
let loaded_at = 10_000;
assert_eq!(misconfigured_since(loaded_at, &[9_000], 20_000), None);
assert_eq!(
misconfigured_since(loaded_at, &[10_050], 10_050 + MISCONFIGURED_AFTER_SECS - 1),
None
);
assert_eq!(
misconfigured_since(loaded_at, &[10_900, 10_050, 9_000], 11_000),
Some(10_050)
);
assert_eq!(misconfigured_since(loaded_at, &[], 99_999), None);
}
}