use std::collections::BTreeSet;
use std::path::Path;
use crate::error::OlError;
use crate::hooks::atomic::RewriteOutcome;
use crate::hooks::cline_providers::{Observation, SlotId};
use crate::hooks::model_relay_endpoints::{self, EndpointRecord, ReleasedBy, SlotState, SlotValue};
pub trait ProviderEndpoints: Send + Sync {
fn agent_type(&self) -> &'static str;
fn record_prefix(&self) -> &'static str;
fn observe(&self, recorded: &BTreeSet<String>) -> Observation;
fn write_slot(
&self,
file: &Path,
slot: &SlotId,
value: &SlotValue,
still: &dyn Fn(&SlotValue) -> bool,
) -> Result<RewriteOutcome, OlError>;
fn copies_of(&self, ports: &BTreeSet<u16>) -> Vec<(std::path::PathBuf, SlotId, u16)>;
fn watch_files(&self) -> Vec<std::path::PathBuf>;
fn reclaim_retired(&self, _main_port: u16) {}
}
pub fn now_unix() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
pub fn names_port(value: &SlotValue, port: u16) -> bool {
value.text().is_some_and(|t| {
reqwest::Url::parse(t.trim()).is_ok_and(|u| {
u.scheme() == "http" && u.host_str() == Some("127.0.0.1") && u.port() == Some(port)
})
})
}
const CONTENDED_ATTEMPTS: u32 = 5;
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct ReleaseSummary {
pub released: usize,
pub restored: usize,
pub left: usize,
pub failed: Vec<String>,
}
pub fn release_all(endpoints: &dyn ProviderEndpoints, by: ReleasedBy) -> ReleaseSummary {
let mut summary = ReleaseSummary::default();
let records = match model_relay_endpoints::endpoint_records(endpoints.record_prefix()) {
Ok(records) => records,
Err(e) => {
tracing::warn!(
agent = endpoints.agent_type(),
code = %e.code,
error = %e.message,
"endpoint records unreadable — no provider slot can be restored"
);
summary.failed.push(e.message);
return summary;
}
};
let live: Vec<(String, EndpointRecord)> =
records.into_iter().filter(|(_, r)| r.is_live()).collect();
let now = now_unix();
for (key, _) in &live {
match model_relay_endpoints::update_endpoint(key, |r| {
r.state = SlotState::Released;
r.released_by = Some(by);
r.changed_at = now;
}) {
Ok(_) => summary.released += 1,
Err(e) => summary.failed.push(format!("{key}: {}", e.message)),
}
}
for (key, rec) in &live {
let Some(slot) = SlotId::from_record_key(key) else {
continue;
};
restore(endpoints, &rec.file, &slot, rec, &mut summary, key);
}
let ports: BTreeSet<u16> = live.iter().map(|(_, r)| r.port).collect();
for (file, slot, port) in endpoints.copies_of(&ports) {
if let Some((key, rec)) = live.iter().find(|(_, r)| r.port == port) {
let prior = match rec.prior.text() {
Some(_) => rec.prior.clone(),
None => SlotValue::Absent,
};
let copy = EndpointRecord {
prior,
..rec.clone()
};
restore(
endpoints,
&file,
&slot,
©,
&mut summary,
&format!("{key} (copy)"),
);
}
}
summary
}
fn restore(
endpoints: &dyn ProviderEndpoints,
file: &Path,
slot: &SlotId,
rec: &EndpointRecord,
summary: &mut ReleaseSummary,
label: &str,
) {
let port = rec.port;
for attempt in 1..=CONTENDED_ATTEMPTS {
match endpoints.write_slot(file, slot, &rec.prior, &|current| names_port(current, port)) {
Ok(RewriteOutcome::Written) => {
summary.restored += 1;
return;
}
Ok(RewriteOutcome::Unchanged) | Ok(RewriteOutcome::Absent) => {
summary.left += 1;
return;
}
Ok(RewriteOutcome::Contended) if attempt < CONTENDED_ATTEMPTS => {
std::thread::sleep(std::time::Duration::from_millis(50));
}
Ok(RewriteOutcome::Contended) => {
summary
.failed
.push(format!("{label}: the file kept changing under the restore"));
return;
}
Err(e) => {
summary.failed.push(format!("{label}: {}", e.message));
return;
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::hooks::cline_providers::{
state_lanes_from, write_slot_in, ClineProviderEndpoints, LaneTag, StateLane,
};
use std::path::PathBuf;
use std::sync::Mutex;
fn with_openlatch_dir<T>(f: impl FnOnce(&Path) -> T) -> T {
let _guard = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let tmp = tempfile::tempdir().expect("tempdir");
let _env = crate::hooks::cline::EnvOverride::apply([(
"OPENLATCH_DIR",
Some(tmp.path().join("openlatch").into_os_string()),
)]);
f(tmp.path())
}
fn wired(port: u16, prior: SlotValue, file: &Path) -> EndpointRecord {
EndpointRecord {
port,
origin: "https://generativelanguage.googleapis.com/".into(),
prior,
last_written: Some(format!("http://127.0.0.1:{port}")),
file: file.to_path_buf(),
state: SlotState::Wired,
released_by: None,
changed_at: 1,
proven_at: None,
family: None,
pending_event: None,
proven_by: None,
misconfigured_event: None,
}
}
fn lane(root: &Path, body: &str) -> (Vec<StateLane>, PathBuf) {
let path = root.join("data").join("globalState.json");
std::fs::create_dir_all(path.parent().expect("parent")).expect("mkdir");
std::fs::write(&path, body).expect("write");
(
state_lanes_from(Some(path.clone()), Some(path.clone())),
path,
)
}
fn json(path: &Path) -> serde_json::Value {
serde_json::from_str(&std::fs::read_to_string(path).expect("read")).expect("json")
}
#[test]
fn uninstall_restores_every_representation_and_is_idempotent() {
with_openlatch_dir(|root| {
let (lanes, file) = lane(
root,
r#"{"geminiBaseUrl":"http://127.0.0.1:7601","ollamaBaseUrl":"http://127.0.0.1:7602","openAiBaseUrl":"http://127.0.0.1:7603/v1","keep":1}"#,
);
for (key, port, prior) in [
("geminiBaseUrl", 7601, SlotValue::Absent),
("ollamaBaseUrl", 7602, SlotValue::Text(String::new())),
(
"openAiBaseUrl",
7603,
SlotValue::Text("https://gw.corp/v1".into()),
),
] {
model_relay_endpoints::put_endpoint(
&format!("cline:gs:shared:{key}"),
wired(port, prior, &file),
)
.expect("put");
}
let endpoints = ClineProviderEndpoints::at(lanes, None);
let first = release_all(&endpoints, ReleasedBy::Teardown);
assert_eq!((first.released, first.restored), (3, 3), "{first:?}");
assert_eq!(
json(&file),
serde_json::json!({"ollamaBaseUrl":"","openAiBaseUrl":"https://gw.corp/v1","keep":1})
);
for (_, rec) in model_relay_endpoints::endpoint_records("cline:").expect("read") {
assert_eq!(rec.state, SlotState::Released);
assert_eq!(rec.released_by, Some(ReleasedBy::Teardown));
}
let before = std::fs::read(&file).expect("read");
let second = release_all(&endpoints, ReleasedBy::Teardown);
assert_eq!(second, ReleaseSummary::default(), "nothing live is left");
assert_eq!(std::fs::read(&file).expect("read"), before);
});
}
#[test]
fn a_slot_that_no_longer_names_our_port_is_left_alone() {
with_openlatch_dir(|root| {
let (lanes, file) = lane(root, r#"{"geminiBaseUrl":"https://gw.theirs/g"}"#);
model_relay_endpoints::put_endpoint(
"cline:gs:shared:geminiBaseUrl",
wired(7601, SlotValue::Absent, &file),
)
.expect("put");
let summary = release_all(
&ClineProviderEndpoints::at(lanes, None),
ReleasedBy::Teardown,
);
assert_eq!((summary.restored, summary.left), (0, 1));
assert_eq!(
json(&file),
serde_json::json!({"geminiBaseUrl":"https://gw.theirs/g"})
);
});
}
#[test]
fn a_copy_the_agent_made_of_our_value_is_restored_with_its_slot() {
with_openlatch_dir(|root| {
let (lanes, file) = lane(root, r#"{"geminiBaseUrl":"http://127.0.0.1:7601"}"#);
let pj = root.join("data").join("settings").join("providers.json");
std::fs::create_dir_all(pj.parent().expect("parent")).expect("mkdir");
std::fs::write(
&pj,
r#"{"providers":{"gemini":{"settings":{"provider":"gemini","baseUrl":"http://127.0.0.1:7601"},"tokenSource":"migration"},"ollama":{"settings":{"baseUrl":"http://127.0.0.1:11434"}}}}"#,
)
.expect("write");
model_relay_endpoints::put_endpoint(
"cline:gs:shared:geminiBaseUrl",
wired(7601, SlotValue::Absent, &file),
)
.expect("put");
let summary = release_all(
&ClineProviderEndpoints::at(lanes, Some(pj.clone())),
ReleasedBy::Teardown,
);
assert_eq!(summary.restored, 2, "{summary:?}");
let providers = json(&pj);
assert!(
providers["providers"]["gemini"]["settings"]
.get("baseUrl")
.is_none(),
"{providers}"
);
assert_eq!(
providers["providers"]["ollama"]["settings"]["baseUrl"], "http://127.0.0.1:11434",
"a customer's own loopback endpoint is not a copy of ours"
);
});
}
#[test]
fn restore_tombstones_before_it_touches_the_file() {
struct Spy {
inner: ClineProviderEndpoints,
states_at_write: Mutex<Vec<SlotState>>,
}
impl ProviderEndpoints for Spy {
fn agent_type(&self) -> &'static str {
"cline"
}
fn record_prefix(&self) -> &'static str {
"cline:"
}
fn observe(&self, recorded: &BTreeSet<String>) -> Observation {
self.inner.observe(recorded)
}
fn write_slot(
&self,
file: &Path,
slot: &SlotId,
value: &SlotValue,
still: &dyn Fn(&SlotValue) -> bool,
) -> Result<RewriteOutcome, OlError> {
for (_, rec) in model_relay_endpoints::endpoint_records("cline:").expect("read") {
self.states_at_write.lock().expect("lock").push(rec.state);
}
write_slot_in(file, slot, value, still)
}
fn copies_of(&self, _: &BTreeSet<u16>) -> Vec<(PathBuf, SlotId, u16)> {
Vec::new()
}
fn watch_files(&self) -> Vec<PathBuf> {
self.inner.watch_files()
}
}
with_openlatch_dir(|root| {
let (lanes, file) = lane(
root,
r#"{"geminiBaseUrl":"http://127.0.0.1:7601","ollamaBaseUrl":"http://127.0.0.1:7602"}"#,
);
for (key, port) in [("geminiBaseUrl", 7601), ("ollamaBaseUrl", 7602)] {
model_relay_endpoints::put_endpoint(
&format!("cline:gs:shared:{key}"),
wired(port, SlotValue::Absent, &file),
)
.expect("put");
}
let spy = Spy {
inner: ClineProviderEndpoints::at(lanes, None),
states_at_write: Mutex::new(Vec::new()),
};
release_all(&spy, ReleasedBy::Teardown);
let seen = spy.states_at_write.into_inner().expect("lock");
assert!(!seen.is_empty());
assert!(
seen.iter().all(|s| *s == SlotState::Released),
"every record was tombstoned before the first file write: {seen:?}"
);
});
}
#[test]
fn names_port_is_exact() {
assert!(names_port(
&SlotValue::Text("http://127.0.0.1:7601/v1".into()),
7601
));
assert!(!names_port(
&SlotValue::Text("http://127.0.0.1:7602".into()),
7601
));
assert!(!names_port(
&SlotValue::Text("http://localhost:7601".into()),
7601
));
assert!(!names_port(&SlotValue::Absent, 7601));
let _ = LaneTag::Shared;
}
}