use std::path::Path;
use jsonc_parser::ast::Value;
use jsonc_parser::common::Ranged;
use jsonc_parser::tokens::{Token, TokenAndRange};
use jsonc_parser::{parse_to_ast, CollectOptions, CommentCollectionStrategy, ParseOptions};
use crate::core::path_compat;
use crate::error::{OlError, ERR_MODEL_RELAY_IO, ERR_MODEL_RELAY_SETTINGS};
use crate::hooks::atomic::{self, RewriteLimits, RewriteOutcome};
use crate::hooks::model_relay_endpoints::{self, EndpointRecord, ReleasedBy, SlotState, SlotValue};
use crate::hooks::provider_endpoints::now_unix;
use crate::model_relay::preflight;
pub const RECORD_PREFIX: &str = "proxy:";
const LIMITS: RewriteLimits = RewriteLimits {
max_bytes: crate::hooks::cline_providers::MAX_STATE_FILE_BYTES,
};
const LAYOUT_PROBE_URL: &str = "http://127.0.0.1:1";
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum KeyWrite {
Written,
AlreadyOurs,
Foreign(String),
Unparseable(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum Observed {
Unset(SlotValue),
Ours(String),
Foreign(String),
Unparseable(String),
NoFile,
}
fn collect_options() -> CollectOptions {
CollectOptions {
comments: CommentCollectionStrategy::Off,
tokens: true,
}
}
fn parse_options() -> ParseOptions {
ParseOptions::default()
}
fn record_key(agent: &str, key: &str, file: &Path) -> String {
format!(
"{RECORD_PREFIX}{agent}:{key}@{}",
path_compat::display_path(file)
)
}
fn find_record(rk: &str) -> Result<Option<EndpointRecord>, OlError> {
Ok(model_relay_endpoints::endpoint_records(rk)?
.into_iter()
.find(|(k, _)| k == rk)
.map(|(_, r)| r))
}
fn reads_back(new: &str, key: &str, url: &str) -> bool {
let Ok(parsed) = parse_to_ast(new, &collect_options(), &parse_options()) else {
return false;
};
match parsed.value {
Some(Value::Object(obj)) => matches!(
obj.properties.iter().rev().find(|p| p.name.as_str() == key).map(|p| &p.value),
Some(Value::StringLit(s)) if s.value == url
),
_ => false,
}
}
pub(crate) fn observe_key(file: &Path, key: &str, agent: &str) -> Result<Observed, OlError> {
let bytes = match std::fs::read(file) {
Ok(b) => b,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Observed::NoFile),
Err(e) => {
return Err(OlError::new(
ERR_MODEL_RELAY_IO,
format!("cannot read {}: {e}", file.display()),
))
}
};
if bytes.len() as u64 > LIMITS.max_bytes {
return Ok(Observed::Unparseable(format!(
"the settings file is over the {} byte limit for a file OpenLatch edits",
LIMITS.max_bytes
)));
}
let raw = match String::from_utf8(bytes) {
Ok(s) => s,
Err(e) => return Ok(Observed::Unparseable(e.to_string())),
};
let parsed = match parse_to_ast(&raw, &collect_options(), &parse_options()) {
Ok(p) => p,
Err(e) => return Ok(Observed::Unparseable(e.to_string())),
};
let obj = match parsed.value {
None => return Ok(Observed::Unset(SlotValue::Absent)),
Some(Value::Object(obj)) => obj,
Some(_) => {
return Ok(Observed::Unparseable(
"the settings file is not a JSON object".to_string(),
))
}
};
let current = match obj.properties.iter().rev().find(|p| p.name.as_str() == key) {
None => SlotValue::Absent,
Some(p) => match &p.value {
Value::NullKeyword(_) => SlotValue::Null,
Value::StringLit(s) => SlotValue::Text(s.value.to_string()),
other => return Ok(Observed::Foreign(other.text(&raw).to_string())),
},
};
if matches!(current, SlotValue::Absent) {
if let Some(new) = splice_set(&raw, key, LAYOUT_PROBE_URL, &SlotValue::Absent)? {
if !reads_back(&new, key, LAYOUT_PROBE_URL) {
return Ok(Observed::Unparseable(
preflight::SETTINGS_LAYOUT.to_string(),
));
}
}
}
match current {
SlotValue::Text(s) => {
let rk = record_key(agent, key, file);
let rec = find_record(&rk)?;
if rec.as_ref().and_then(|r| r.last_written.as_deref()) == Some(s.as_str()) {
Ok(Observed::Ours(s))
} else if s.trim().is_empty() {
Ok(Observed::Unset(SlotValue::Text(s)))
} else {
Ok(Observed::Foreign(s))
}
}
SlotValue::Absent | SlotValue::Null => Ok(Observed::Unset(current)),
}
}
fn parse_port(url: &str) -> Result<u16, OlError> {
reqwest::Url::parse(url)
.ok()
.and_then(|u| u.port())
.ok_or_else(|| {
OlError::new(
ERR_MODEL_RELAY_IO,
format!("{url} is not a URL with a port"),
)
})
}
enum CreateOutcome {
Written,
AlreadyExists,
Err(OlError),
}
fn create_new_file(file: &Path, key: &str, url: &str) -> CreateOutcome {
let json_url = match serde_json::to_string(url) {
Ok(j) => j,
Err(e) => return CreateOutcome::Err(OlError::new(ERR_MODEL_RELAY_IO, e.to_string())),
};
let content = format!("{{\n \"{key}\": {json_url}\n}}\n");
use std::io::Write;
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(file)
{
Ok(mut f) => match f.write_all(content.as_bytes()) {
Ok(()) => CreateOutcome::Written,
Err(e) => CreateOutcome::Err(OlError::new(ERR_MODEL_RELAY_IO, e.to_string())),
},
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => CreateOutcome::AlreadyExists,
Err(e) => CreateOutcome::Err(OlError::new(
ERR_MODEL_RELAY_IO,
format!("cannot create {}: {e}", file.display()),
)),
}
}
pub fn write_proxy_key(
file: &Path,
key: &str,
url: &str,
agent: &str,
) -> Result<KeyWrite, OlError> {
if !file.parent().is_some_and(std::path::Path::is_dir) {
return Err(OlError::new(
ERR_MODEL_RELAY_SETTINGS,
preflight::SETTINGS_NO_FILE.to_string(),
));
}
if key.contains('@') {
return Err(OlError::new(
ERR_MODEL_RELAY_IO,
format!("proxy settings key {key:?} must not contain '@'"),
));
}
let rk = record_key(agent, key, file);
let port = parse_port(url)?;
for _ in 0..crate::hooks::provider_endpoints::CONTENDED_ATTEMPTS {
let observed = observe_key(file, key, agent)?;
match &observed {
Observed::Foreign(v) => return Ok(KeyWrite::Foreign(v.clone())),
Observed::Unparseable(why) => return Ok(KeyWrite::Unparseable(why.clone())),
Observed::Ours(v) if v == url => {
if find_record(&rk)?.is_some_and(|r| r.state != SlotState::Wired) {
model_relay_endpoints::update_endpoint(&rk, |r| {
r.state = SlotState::Wired;
r.changed_at = now_unix();
})?;
}
return Ok(KeyWrite::AlreadyOurs);
}
Observed::Ours(_) | Observed::Unset(_) | Observed::NoFile => {}
}
let prior = match find_record(&rk)? {
Some(rec) => rec.prior,
None => match &observed {
Observed::NoFile => SlotValue::Absent,
Observed::Unset(v) => v.clone(),
Observed::Ours(v) => SlotValue::Text(v.clone()),
Observed::Foreign(_) | Observed::Unparseable(_) => unreachable!("handled above"),
},
};
model_relay_endpoints::put_endpoint(
&rk,
EndpointRecord {
port,
origin: String::new(),
prior,
last_written: Some(url.to_string()),
file: file.to_path_buf(),
state: SlotState::Pending,
released_by: None,
changed_at: now_unix(),
proven_at: None,
proven_by: None,
family: None,
pending_event: None,
misconfigured_event: None,
},
)?;
if matches!(observed, Observed::NoFile) {
match create_new_file(file, key, url) {
CreateOutcome::Written => {
model_relay_endpoints::update_endpoint(&rk, |r| r.state = SlotState::Wired)?;
return Ok(KeyWrite::Written);
}
CreateOutcome::AlreadyExists => continue,
CreateOutcome::Err(e) => return Err(e),
}
}
let expect = match &observed {
Observed::Unset(v) => v.clone(),
Observed::Ours(v) => SlotValue::Text(v.clone()),
_ => SlotValue::Absent,
};
let mut refused = false;
let outcome = atomic::rewrite_existing_text(file, LIMITS, |raw| {
match splice_set(raw, key, url, &expect)? {
Some(new) if !reads_back(&new, key, url) => {
refused = true;
Ok(None)
}
other => Ok(other),
}
})?;
if refused {
model_relay_endpoints::update_endpoint(&rk, |r| {
r.state = SlotState::Released;
r.released_by = Some(ReleasedBy::Wiring);
r.changed_at = now_unix();
})?;
return Ok(KeyWrite::Unparseable(
preflight::SETTINGS_LAYOUT.to_string(),
));
}
match outcome {
RewriteOutcome::Written => {
model_relay_endpoints::update_endpoint(&rk, |r| r.state = SlotState::Wired)?;
return Ok(KeyWrite::Written);
}
RewriteOutcome::Unchanged | RewriteOutcome::Contended | RewriteOutcome::Absent => {
std::thread::sleep(std::time::Duration::from_millis(50));
}
}
}
model_relay_endpoints::update_endpoint(&rk, |r| {
r.state = SlotState::Released;
r.released_by = Some(ReleasedBy::Wiring);
r.changed_at = now_unix();
})?;
Err(OlError::new(
ERR_MODEL_RELAY_SETTINGS,
preflight::SETTINGS_CONTENDED.to_string(),
))
}
pub fn release_proxy_key(
file: &Path,
key: &str,
agent: &str,
by: ReleasedBy,
) -> Result<bool, OlError> {
let rk = record_key(agent, key, file);
let Some(rec) = find_record(&rk)? else {
return Ok(false);
};
if rec.state != SlotState::Released {
model_relay_endpoints::update_endpoint(&rk, |r| {
r.state = SlotState::Released;
r.released_by = Some(by);
r.changed_at = now_unix();
})?;
}
for _ in 0..crate::hooks::provider_endpoints::CONTENDED_ATTEMPTS {
let outcome = atomic::rewrite_existing_text(file, LIMITS, |raw| {
splice_restore(raw, key, rec.last_written.as_deref(), &rec.prior)
})?;
match outcome {
RewriteOutcome::Written => return Ok(true),
RewriteOutcome::Unchanged | RewriteOutcome::Absent => return Ok(false),
RewriteOutcome::Contended => {
std::thread::sleep(std::time::Duration::from_millis(50));
}
}
}
Err(OlError::new(
ERR_MODEL_RELAY_SETTINGS,
"the file kept changing under the restore".to_string(),
))
}
const INDENT_UNIT: &str = " ";
fn splice_set(
raw: &str,
key: &str,
url: &str,
expect: &SlotValue,
) -> Result<Option<String>, OlError> {
let parsed = parse_to_ast(raw, &collect_options(), &parse_options())
.map_err(|e| OlError::new(ERR_MODEL_RELAY_IO, e.to_string()))?;
let json_url =
serde_json::to_string(url).map_err(|e| OlError::new(ERR_MODEL_RELAY_IO, e.to_string()))?;
let Some(value) = parsed.value else {
if !matches!(expect, SlotValue::Absent) {
return Ok(None);
}
let mut out = raw.to_string();
if !out.is_empty() && !out.ends_with('\n') {
out.push('\n');
}
out.push_str(&format!("{{\n{INDENT_UNIT}\"{key}\": {json_url}\n}}\n"));
return Ok(Some(out));
};
let Value::Object(obj) = value else {
return Ok(None);
};
let current_prop = obj.properties.iter().rev().find(|p| p.name.as_str() == key);
let current = match current_prop {
None => Some(SlotValue::Absent),
Some(p) => match &p.value {
Value::NullKeyword(_) => Some(SlotValue::Null),
Value::StringLit(s) => Some(SlotValue::Text(s.value.to_string())),
_ => None,
},
};
if current.as_ref() != Some(expect) {
return Ok(None);
}
let tokens = parsed.tokens.unwrap_or_default();
match current_prop {
Some(p) => {
let vr = p.value.range();
Ok(Some(format!(
"{}{json_url}{}",
&raw[..vr.start],
&raw[vr.end..]
)))
}
None => Ok(Some(insert_property(raw, &obj, &tokens, key, &json_url))),
}
}
fn comma_between(
tokens: &[TokenAndRange],
from: usize,
to: usize,
) -> Option<jsonc_parser::common::Range> {
tokens
.iter()
.find(|t| matches!(t.token, Token::Comma) && t.range.start >= from && t.range.end <= to)
.map(|t| t.range)
}
fn comma_before(
tokens: &[TokenAndRange],
before: usize,
object_start: usize,
) -> Option<jsonc_parser::common::Range> {
tokens
.iter()
.rfind(|t| {
matches!(t.token, Token::Comma)
&& t.range.end <= before
&& t.range.start >= object_start
})
.map(|t| t.range)
}
fn line_indent(raw: &str, pos: usize) -> String {
let line_start = raw[..pos].rfind('\n').map(|i| i + 1).unwrap_or(0);
raw[line_start..pos]
.chars()
.take_while(|c| *c == ' ' || *c == '\t')
.collect()
}
fn insert_property(
raw: &str,
obj: &jsonc_parser::ast::Object,
tokens: &[TokenAndRange],
key: &str,
json_url: &str,
) -> String {
let c = obj.range.end - 1; let insertion = format!("\"{key}\": {json_url}");
if obj.properties.is_empty() {
let inner_has_newline = raw[obj.range.start + 1..c].contains('\n');
if !inner_has_newline {
return format!("{}{insertion}{}", &raw[..c], &raw[c..]);
}
let after_brace = obj.range.start + 1;
let nl_at = raw[after_brace..]
.find('\n')
.map(|off| after_brace + off + 1)
.unwrap_or(after_brace);
let close_indent = line_indent(raw, c);
let indent = format!("{close_indent}{INDENT_UNIT}");
return format!("{}{indent}{insertion}\n{}", &raw[..nl_at], &raw[nl_at..]);
}
let p = obj.properties.last().expect("non-empty checked above");
let p_end = p.value.range().end;
let single_line = !raw[obj.range.start..c].contains('\n');
if single_line {
return match comma_between(tokens, p_end, c) {
None => format!("{}, {insertion}{}", &raw[..p_end], &raw[p_end..]),
Some(cr) => format!("{} {insertion},{}", &raw[..cr.end], &raw[cr.end..]),
};
}
let indent = line_indent(raw, p.range.start);
let comma = comma_between(tokens, p_end, c);
let m = comma.map(|cr| cr.end).unwrap_or(p_end);
let brace_shared = !raw[m..c].contains('\n');
if brace_shared {
return match comma {
None => format!("{},\n{indent}{insertion}{}", &raw[..p_end], &raw[p_end..]),
Some(cr) => format!("{}\n{indent}{insertion},{}", &raw[..cr.end], &raw[cr.end..]),
};
}
let e = raw[m..c].find('\n').map(|off| m + off).unwrap_or(m);
match comma {
None => format!(
"{},{}\n{indent}{insertion}{}",
&raw[..p_end],
&raw[p_end..e],
&raw[e..]
),
Some(_) => format!("{}\n{indent}{insertion},{}", &raw[..e], &raw[e..]),
}
}
fn splice_restore(
raw: &str,
key: &str,
ours: Option<&str>,
prior: &SlotValue,
) -> Result<Option<String>, OlError> {
let parsed = parse_to_ast(raw, &collect_options(), &parse_options())
.map_err(|e| OlError::new(ERR_MODEL_RELAY_IO, e.to_string()))?;
let Some(Value::Object(obj)) = parsed.value else {
return Ok(None);
};
let Some(prop) = obj.properties.iter().rev().find(|p| p.name.as_str() == key) else {
return Ok(None);
};
let holds_ours = matches!(&prop.value, Value::StringLit(s) if Some(s.value.as_ref()) == ours);
if !holds_ours {
return Ok(None);
}
match prior {
SlotValue::Text(t) => {
let json = serde_json::to_string(t)
.map_err(|e| OlError::new(ERR_MODEL_RELAY_IO, e.to_string()))?;
let vr = prop.value.range();
Ok(Some(format!(
"{}{json}{}",
&raw[..vr.start],
&raw[vr.end..]
)))
}
SlotValue::Null => {
let vr = prop.value.range();
Ok(Some(format!("{}null{}", &raw[..vr.start], &raw[vr.end..])))
}
SlotValue::Absent => {
let tokens = parsed.tokens.unwrap_or_default();
Ok(Some(delete_property(raw, &obj, &tokens, prop)))
}
}
}
fn delete_property(
raw: &str,
obj: &jsonc_parser::ast::Object,
tokens: &[TokenAndRange],
p: &jsonc_parser::ast::ObjectProp,
) -> String {
let c = obj.range.end - 1;
let following = comma_between(tokens, p.range.end, c);
let preceding = comma_before(tokens, p.range.start, obj.range.start);
let single_line = !raw[obj.range.start..c].contains('\n');
if single_line {
return match following {
Some(cr) => {
let mut start = p.range.start;
if start > 0 && raw.as_bytes()[start - 1] == b' ' {
start -= 1;
}
format!("{}{}", &raw[..start], &raw[cr.end..])
}
None => match preceding {
Some(pcr) => format!("{}{}", &raw[..pcr.start], &raw[p.range.end..]),
None => format!("{}{}", &raw[..p.range.start], &raw[p.range.end..]),
},
};
}
let m = following.map(|cr| cr.end).unwrap_or(p.range.end);
let brace_shared = !raw[m..c].contains('\n');
if brace_shared {
return match (following, preceding) {
(Some(fcr), Some(pcr)) => format!("{}{}", &raw[..pcr.end], &raw[fcr.end..]),
(Some(fcr), None) => format!("{}{}", &raw[..p.range.start], &raw[fcr.end..]),
(None, Some(pcr)) => format!("{}{}", &raw[..pcr.start], &raw[p.range.end..]),
(None, None) => format!("{}{}", &raw[..p.range.start], &raw[p.range.end..]),
};
}
let line_start = raw[..p.range.start].rfind('\n').map(|i| i + 1).unwrap_or(0);
let after = following.map(|cr| cr.end).unwrap_or(p.range.end);
let line_end = raw[after..]
.find('\n')
.map(|off| after + off + 1)
.unwrap_or(raw.len());
if following.is_some() {
format!("{}{}", &raw[..line_start], &raw[line_end..])
} else if let Some(pcr) = preceding {
format!(
"{}{}{}",
&raw[..pcr.start],
&raw[pcr.end..line_start],
&raw[line_end..]
)
} else {
format!("{}{}", &raw[..line_start], &raw[line_end..])
}
}
#[cfg(test)]
mod tests {
use super::*;
fn with_openlatch_dir<T>(f: impl FnOnce() -> 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 prev = std::env::var_os("OPENLATCH_DIR");
std::env::set_var("OPENLATCH_DIR", tmp.path());
let out = f();
match prev {
Some(v) => std::env::set_var("OPENLATCH_DIR", v),
None => std::env::remove_var("OPENLATCH_DIR"),
}
out
}
const URL: &str = "http://127.0.0.1:7600";
const AGENT: &str = "fake-a";
const KEY: &str = "fake.proxy";
#[test]
fn insert_then_release_is_byte_identical_for_a_commented_file() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
let original = "{\n // customer\n \"a\": 1 // note\n}\n";
std::fs::write(&file, original).expect("seed");
let write = write_proxy_key(&file, KEY, URL, AGENT).expect("write");
assert_eq!(write, KeyWrite::Written);
let after_write = std::fs::read_to_string(&file).expect("read after write");
assert_eq!(
after_write,
format!(
"{{\n // customer\n \"a\": 1, // note\n \"{KEY}\": {}\n}}\n",
serde_json::to_string(URL).unwrap()
)
);
assert!(reads_back(&after_write, KEY, URL));
assert_eq!(
write_proxy_key(&file, KEY, URL, AGENT).expect("second write"),
KeyWrite::AlreadyOurs
);
assert_eq!(std::fs::read_to_string(&file).unwrap(), after_write);
assert!(release_proxy_key(&file, KEY, AGENT, ReleasedBy::Teardown).expect("release"));
assert_eq!(std::fs::read_to_string(&file).unwrap(), original);
});
}
#[test]
fn a_foreign_value_is_never_overwritten_and_never_recorded() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
let original = "{\n \"fake.proxy\": \"http://corp:3128\"\n}\n";
std::fs::write(&file, original).expect("seed");
let write = write_proxy_key(&file, KEY, URL, AGENT).expect("write");
assert_eq!(write, KeyWrite::Foreign("http://corp:3128".to_string()));
assert_eq!(std::fs::read_to_string(&file).unwrap(), original);
assert!(model_relay_endpoints::endpoint_records(RECORD_PREFIX)
.expect("read ledger")
.is_empty());
});
}
#[test]
fn our_loopback_without_a_record_is_foreign() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
std::fs::write(&file, format!("{{\n \"fake.proxy\": \"{URL}\"\n}}\n")).unwrap();
assert_eq!(
write_proxy_key(&file, KEY, URL, AGENT).expect("write over loopback"),
KeyWrite::Foreign(URL.to_string())
);
});
}
#[test]
fn an_absent_file_is_created_only_when_its_directory_exists() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
assert_eq!(
write_proxy_key(&file, KEY, URL, AGENT).expect("write"),
KeyWrite::Written
);
assert_eq!(
std::fs::read_to_string(&file).unwrap(),
format!(
"{{\n \"{KEY}\": {}\n}}\n",
serde_json::to_string(URL).unwrap()
)
);
assert!(release_proxy_key(&file, KEY, AGENT, ReleasedBy::Teardown).expect("release"));
assert_eq!(std::fs::read_to_string(&file).unwrap(), "{\n}\n");
let missing_dir = dir.path().join("nope").join("settings.json");
let err = write_proxy_key(&missing_dir, KEY, URL, AGENT).expect_err("no directory");
assert_eq!(err.code, ERR_MODEL_RELAY_SETTINGS);
assert!(!missing_dir.exists());
});
}
#[test]
fn a_present_unset_key_changes_only_its_value_span() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
let original = "{\n // customer\n \"fake.proxy\": null,\n \"b\": 2,\n}\n";
std::fs::write(&file, original).expect("seed");
assert_eq!(
write_proxy_key(&file, KEY, URL, AGENT).expect("write"),
KeyWrite::Written
);
assert_eq!(
std::fs::read_to_string(&file).unwrap(),
format!(
"{{\n // customer\n \"fake.proxy\": {},\n \"b\": 2,\n}}\n",
serde_json::to_string(URL).unwrap()
)
);
});
}
#[test]
fn an_absent_key_is_inserted_in_the_files_indentation() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let json_url = serde_json::to_string(URL).unwrap();
let tab_file = dir.path().join("tab.json");
std::fs::write(&tab_file, "{\n\t\"a\": 1\n}\n").unwrap();
assert_eq!(
write_proxy_key(&tab_file, KEY, URL, "tab-agent").expect("tab write"),
KeyWrite::Written
);
assert_eq!(
std::fs::read_to_string(&tab_file).unwrap(),
format!("{{\n\t\"a\": 1,\n\t\"{KEY}\": {json_url}\n}}\n")
);
let four_file = dir.path().join("four.json");
std::fs::write(&four_file, "{\n \"a\": 1,\n}\n").unwrap();
assert_eq!(
write_proxy_key(&four_file, KEY, URL, "four-agent").expect("four write"),
KeyWrite::Written
);
assert_eq!(
std::fs::read_to_string(&four_file).unwrap(),
format!("{{\n \"a\": 1,\n \"{KEY}\": {json_url},\n}}\n")
);
});
}
#[test]
fn insert_then_release_is_byte_identical_across_every_shape() {
with_openlatch_dir(|| {
let u = format!("\"{KEY}\": {}", serde_json::to_string(URL).unwrap());
let shapes: [(&str, String); 7] = [
(
"{\n // customer\n \"a\": 1 // note\n}\n",
format!("{{\n // customer\n \"a\": 1, // note\n {u}\n}}\n"),
),
(
"{\n \"a\": 1,\n}\n",
format!("{{\n \"a\": 1,\n {u},\n}}\n"),
),
("{ \"a\": 1 }\n", format!("{{ \"a\": 1, {u} }}\n")),
("{}\n", format!("{{{u}}}\n")),
(
"{ \"a\": 1,\n \"b\": 2 }\n",
format!("{{ \"a\": 1,\n \"b\": 2,\n {u} }}\n"),
),
(
"{ \"a\": 1,\n \"b\": 2 }",
format!("{{ \"a\": 1,\n \"b\": 2,\n {u} }}"),
),
(
"{ \"a\": 1,\n \"b\": 2, }\n",
format!("{{ \"a\": 1,\n \"b\": 2,\n {u}, }}\n"),
),
];
for (i, (original, expected_after_write)) in shapes.iter().enumerate() {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
std::fs::write(&file, original).expect("seed");
let agent = format!("shape-{i}");
assert_eq!(
write_proxy_key(&file, KEY, URL, &agent).expect("write"),
KeyWrite::Written,
"shape {i}"
);
let after_write = std::fs::read_to_string(&file).unwrap();
assert_eq!(
after_write, *expected_after_write,
"shape {i} after the write"
);
assert!(
reads_back(&after_write, KEY, URL),
"shape {i} must read back"
);
assert!(
release_proxy_key(&file, KEY, &agent, ReleasedBy::Teardown).expect("release"),
"shape {i}"
);
assert_eq!(
std::fs::read_to_string(&file).unwrap(),
*original,
"shape {i} restore"
);
}
let refused_dir = tempfile::tempdir().expect("tempdir");
let refused_file = refused_dir.path().join("settings.json");
let refused = "{ \"a\": 1 /* see\n ticket */\n}\n";
std::fs::write(&refused_file, refused).unwrap();
assert_eq!(
write_proxy_key(&refused_file, KEY, URL, "refused-agent").expect("write refused"),
KeyWrite::Unparseable(preflight::SETTINGS_LAYOUT.to_string())
);
assert_eq!(std::fs::read_to_string(&refused_file).unwrap(), refused);
let refused_rk = record_key("refused-agent", KEY, &refused_file);
assert!(
find_record(&refused_rk).expect("ledger").is_none(),
"step 1 refuses the layout before any record is put"
);
});
}
#[test]
fn observe_key_sees_the_refused_layout() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
let refused = "{ \"a\": 1 /* see\n ticket */\n}";
std::fs::write(&file, refused).unwrap();
assert_eq!(
observe_key(&file, KEY, AGENT).expect("observe"),
Observed::Unparseable(preflight::SETTINGS_LAYOUT.to_string())
);
assert_eq!(std::fs::read_to_string(&file).unwrap(), refused);
let seven_shapes = [
"{\n // customer\n \"a\": 1 // note\n}\n",
"{\n \"a\": 1,\n}\n",
"{ \"a\": 1 }\n",
"{}\n",
"{ \"a\": 1,\n \"b\": 2 }\n",
"{ \"a\": 1,\n \"b\": 2 }",
"{ \"a\": 1,\n \"b\": 2, }\n",
];
for shape in seven_shapes {
std::fs::write(&file, shape).unwrap();
assert_eq!(
observe_key(&file, KEY, AGENT).expect("observe shape"),
Observed::Unset(SlotValue::Absent),
"{shape:?}"
);
assert_eq!(std::fs::read_to_string(&file).unwrap(), shape);
}
});
}
#[test]
fn an_unparseable_or_non_object_file_is_skipped() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
std::fs::write(&file, "{ \"a\": ").unwrap();
assert!(matches!(
write_proxy_key(&file, KEY, URL, AGENT).expect("write on garbage"),
KeyWrite::Unparseable(_)
));
assert_eq!(std::fs::read_to_string(&file).unwrap(), "{ \"a\": ");
std::fs::write(&file, "[]").unwrap();
assert!(matches!(
write_proxy_key(&file, KEY, URL, AGENT).expect("write on array"),
KeyWrite::Unparseable(_)
));
assert_eq!(std::fs::read_to_string(&file).unwrap(), "[]");
assert!(model_relay_endpoints::endpoint_records(RECORD_PREFIX)
.expect("ledger")
.is_empty());
});
}
#[test]
fn a_second_write_is_already_ours_and_writes_nothing() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
std::fs::write(&file, "{\n \"a\": 1\n}\n").unwrap();
assert_eq!(
write_proxy_key(&file, KEY, URL, AGENT).expect("first"),
KeyWrite::Written
);
let fp_before = crate::fs_secure::fingerprint(&file).expect("fingerprint");
assert_eq!(
write_proxy_key(&file, KEY, URL, AGENT).expect("second"),
KeyWrite::AlreadyOurs
);
let fp_after = crate::fs_secure::fingerprint(&file).expect("fingerprint");
assert_eq!(fp_before, fp_after, "the file must not be touched");
});
}
#[test]
fn a_port_change_rewrites_our_value_and_keeps_the_prior() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
std::fs::write(&file, "{\n \"fake.proxy\": null\n}\n").unwrap();
assert_eq!(
write_proxy_key(&file, KEY, "http://127.0.0.1:7600", AGENT).expect("first"),
KeyWrite::Written
);
assert_eq!(
write_proxy_key(&file, KEY, "http://127.0.0.1:7700", AGENT).expect("second"),
KeyWrite::Written
);
let content = std::fs::read_to_string(&file).unwrap();
assert!(content.contains("7700"));
assert!(!content.contains("7600"));
let rk = record_key(AGENT, KEY, &file);
let rec = find_record(&rk)
.expect("read ledger")
.expect("record exists");
assert_eq!(
rec.prior,
SlotValue::Null,
"the customer's original null must survive a port change"
);
});
}
#[test]
fn release_restores_the_prior_representation() {
with_openlatch_dir(|| {
for (prior, expected) in [
(SlotValue::Null, "null".to_string()),
(SlotValue::Text(String::new()), "\"\"".to_string()),
(SlotValue::Text("x".to_string()), "\"x\"".to_string()),
] {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
std::fs::write(
&file,
format!(
"{{\n \"fake.proxy\": {}\n}}\n",
serde_json::to_string(URL).unwrap()
),
)
.unwrap();
let rk = record_key(AGENT, KEY, &file);
model_relay_endpoints::put_endpoint(
&rk,
EndpointRecord {
port: 7600,
origin: String::new(),
prior: prior.clone(),
last_written: Some(URL.to_string()),
file: file.clone(),
state: SlotState::Wired,
released_by: None,
changed_at: 0,
proven_at: None,
proven_by: None,
family: None,
pending_event: None,
misconfigured_event: None,
},
)
.expect("seed record");
assert!(
release_proxy_key(&file, KEY, AGENT, ReleasedBy::Teardown).expect("release"),
"{prior:?}"
);
assert_eq!(
std::fs::read_to_string(&file).unwrap(),
format!("{{\n \"fake.proxy\": {expected}\n}}\n"),
"{prior:?}"
);
}
});
}
#[test]
fn release_leaves_a_value_the_developer_changed() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
std::fs::write(&file, "{\n \"a\": 1\n}\n").unwrap();
write_proxy_key(&file, KEY, URL, AGENT).expect("write");
let developer_changed = std::fs::read_to_string(&file)
.unwrap()
.replace(URL, "http://mine");
std::fs::write(&file, &developer_changed).unwrap();
assert!(
!release_proxy_key(&file, KEY, AGENT, ReleasedBy::Teardown).expect("release"),
"the developer's own value must never be touched"
);
assert_eq!(std::fs::read_to_string(&file).unwrap(), developer_changed);
let rk = record_key(AGENT, KEY, &file);
let rec = find_record(&rk)
.expect("read ledger")
.expect("record exists");
assert_eq!(rec.state, SlotState::Released);
});
}
#[test]
fn release_tombstones_even_when_the_file_is_gone() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
std::fs::write(&file, "{\n \"a\": 1\n}\n").unwrap();
write_proxy_key(&file, KEY, URL, AGENT).expect("write");
std::fs::remove_file(&file).unwrap();
assert!(!release_proxy_key(&file, KEY, AGENT, ReleasedBy::Teardown).expect("release"));
let rk = record_key(AGENT, KEY, &file);
let rec = find_record(&rk)
.expect("read ledger")
.expect("record exists");
assert_eq!(rec.state, SlotState::Released);
assert_eq!(rec.released_by, Some(ReleasedBy::Teardown));
});
}
#[test]
fn records_live_under_the_proxy_prefix_only() {
with_openlatch_dir(|| {
let dir = tempfile::tempdir().expect("tempdir");
let file = dir.path().join("settings.json");
std::fs::write(&file, "{\n \"a\": 1\n}\n").unwrap();
write_proxy_key(&file, KEY, URL, AGENT).expect("write");
assert!(model_relay_endpoints::endpoint_records("cline:")
.expect("cline records")
.is_empty());
let proxy_records =
model_relay_endpoints::endpoint_records(RECORD_PREFIX).expect("proxy records");
assert_eq!(proxy_records.len(), 1);
assert_eq!(proxy_records[0].1.origin, "");
});
}
}