use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::OnceLock;
use serde::{Deserialize, Serialize};
use crate::schema::ModelSchema;
pub const DEFAULT_CATALOG_URL: &str = "https://car.parslee.ai/catalog.json";
pub const DEFAULT_CATALOG_PUBKEY: &str = "FpylXbP3mL67k6B+sRAko0MY5kYPKbdFDT3nI1WHyCI=";
pub fn catalog_url() -> String {
std::env::var("CAR_CATALOG_URL")
.ok()
.filter(|url| !url.trim().is_empty())
.unwrap_or_else(|| DEFAULT_CATALOG_URL.to_string())
}
pub fn catalog_public_key() -> String {
std::env::var("CAR_CATALOG_PUBKEY")
.ok()
.filter(|key| !key.trim().is_empty())
.unwrap_or_else(|| DEFAULT_CATALOG_PUBKEY.to_string())
}
pub const CATALOG_CACHE_FILE: &str = "catalog-cache.json";
const MAX_CATALOG_BYTES: u64 = 8 * 1024 * 1024;
const MAX_SIGNATURE_BYTES: u64 = 4 * 1024;
async fn read_capped(resp: reqwest::Response, limit: u64, what: &str) -> Result<Vec<u8>, String> {
use futures::StreamExt;
if resp.content_length().is_some_and(|len| len > limit) {
return Err(format!("{what} too large (over {limit} bytes)"));
}
let mut body = Vec::new();
let mut stream = resp.bytes_stream();
while let Some(chunk) = stream.next().await {
let chunk = chunk.map_err(|e| format!("read {what}: {e}"))?;
if (body.len() + chunk.len()) as u64 > limit {
return Err(format!("{what} too large (over {limit} bytes)"));
}
body.extend_from_slice(&chunk);
}
Ok(body)
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CatalogDoc {
pub version: u64,
#[serde(default)]
pub models: Vec<ModelSchema>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub revoked: Vec<String>,
}
#[derive(Deserialize)]
struct RawCatalogDoc {
version: u64,
#[serde(default)]
models: Vec<serde_json::Value>,
#[serde(default)]
revoked: Vec<String>,
}
fn parse_catalog_row(row: serde_json::Value) -> Option<ModelSchema> {
let id = row
.get("id")
.and_then(|id| id.as_str())
.unwrap_or("<no id>")
.to_string();
if let Some(min) = row.get("min_car_version").and_then(|v| v.as_str()) {
if !version_at_least(env!("CARGO_PKG_VERSION"), min) {
tracing::debug!(%id, min, "skipping catalog row that needs a newer CAR");
return None;
}
}
match serde_json::from_value(row) {
Ok(schema) => Some(schema),
Err(error) => {
tracing::warn!(%id, %error, "skipping catalog row this CAR cannot read");
None
}
}
}
fn version_at_least(have: &str, need: &str) -> bool {
let parse = |v: &str| -> Option<Vec<u64>> {
v.split(['-', '+'])
.next()?
.split('.')
.map(|part| part.parse().ok())
.collect()
};
match (parse(have), parse(need)) {
(Some(have), Some(need)) => {
let width = have.len().max(need.len());
let pad = |mut v: Vec<u64>| {
v.resize(width, 0);
v
};
pad(have) >= pad(need)
}
_ => false,
}
}
#[derive(Debug, Clone)]
pub struct VerifiedCatalog {
doc: CatalogDoc,
signed_body: String,
signature: String,
}
impl VerifiedCatalog {
pub fn version(&self) -> u64 {
self.doc.version
}
pub fn model_count(&self) -> usize {
self.doc.models.len()
}
pub fn into_models(self) -> Vec<ModelSchema> {
self.doc.models
}
pub fn revoked(&self) -> &[String] {
&self.doc.revoked
}
pub fn envelope(&self) -> CatalogCacheEnvelope {
CatalogCacheEnvelope {
signed_body: self.signed_body.clone(),
signature: self.signature.clone(),
}
}
}
pub const RETAINED_FILE: &str = "catalog-retained.json";
pub const REVOKED_NAMES_FILE: &str = "catalog-revoked.json";
pub fn revoked_names_path(state_root: &Path) -> PathBuf {
state_root.join(REVOKED_NAMES_FILE)
}
pub const REVOKED_MODELS_FILE: &str = "catalog-revoked-models.json";
pub fn revoked_models_path(state_root: &Path) -> PathBuf {
state_root.join(REVOKED_MODELS_FILE)
}
pub fn load_revoked_names(path: &Path) -> std::collections::BTreeMap<String, String> {
std::fs::read_to_string(path)
.ok()
.and_then(|json| serde_json::from_str(&json).ok())
.unwrap_or_default()
}
pub fn save_revoked_names(
path: &Path,
names: &std::collections::BTreeMap<String, String>,
) -> std::io::Result<()> {
if names.is_empty() {
return match std::fs::remove_file(path) {
Err(error) if error.kind() != std::io::ErrorKind::NotFound => Err(error),
_ => Ok(()),
};
}
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let json = serde_json::to_string_pretty(names).map_err(std::io::Error::other)?;
let tmp = unique_temp_path(path);
std::fs::write(&tmp, json)?;
if let Err(error) = std::fs::rename(&tmp, path) {
let _ = std::fs::remove_file(&tmp);
return Err(error);
}
Ok(())
}
pub const STATUS_FILE: &str = "catalog-status.json";
pub const STALE_AFTER_SECS: u64 = 3 * 24 * 60 * 60;
pub const LOOP_STALL_SECS: u64 = 60 * 60;
pub const DAEMON_ALIVE_SECS: u64 = 60 * 60;
pub const DAEMON_RUN_GAP_SECS: u64 = 30 * 60;
pub fn status_path(state_root: &Path) -> PathBuf {
state_root.join(STATUS_FILE)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CheckFailure {
Unreachable,
Rejected,
Older,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct CatalogStatus {
#[serde(default)]
pub source: Option<String>,
#[serde(default)]
pub last_attempt_at: Option<u64>,
#[serde(default)]
pub last_verified_at: Option<u64>,
#[serde(default)]
pub verified_version: Option<u64>,
#[serde(default)]
pub failing_since: Option<u64>,
#[serde(default)]
pub last_error: Option<String>,
#[serde(default)]
pub last_failure: Option<CheckFailure>,
#[serde(default)]
pub off_since: Option<u64>,
#[serde(default)]
pub loop_alive_at: Option<u64>,
#[serde(default)]
pub daemon_alive_at: Option<u64>,
#[serde(default)]
pub daemon_alive_since: Option<u64>,
}
pub fn load_status(path: &Path) -> CatalogStatus {
std::fs::read_to_string(path)
.ok()
.and_then(|json| serde_json::from_str(&json).ok())
.unwrap_or_default()
}
fn save_status(path: &Path, status: &CatalogStatus) -> std::io::Result<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let json = serde_json::to_string_pretty(status).map_err(std::io::Error::other)?;
let tmp = unique_temp_path(path);
std::fs::write(&tmp, json)?;
if let Err(error) = std::fs::rename(&tmp, path) {
let _ = std::fs::remove_file(&tmp);
return Err(error);
}
Ok(())
}
static STATUS_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
fn update_status(path: &Path, what: &str, change: impl FnOnce(&mut CatalogStatus) -> bool) {
let _guard = STATUS_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let mut status = load_status(path);
if change(&mut status) {
if let Err(error) = save_status(path, &status) {
tracing::warn!(path = %path.display(), %error, "could not record {what}");
}
}
}
pub fn record_check(
path: &Path,
source: &str,
outcome: Result<u64, (CheckFailure, &str)>,
now: u64,
) {
update_status(path, "the catalog check", |status| {
status.source = Some(
source
.split(['?', '#'])
.next()
.unwrap_or(source)
.to_string(),
);
status.last_attempt_at = Some(now);
status.off_since = None;
match outcome {
Ok(version) => {
status.last_verified_at = Some(now);
status.verified_version = Some(version);
status.failing_since = None;
status.last_error = None;
status.last_failure = None;
}
Err((failure, error)) => {
status.failing_since.get_or_insert(now);
status.last_error = Some(error.chars().take(300).collect());
status.last_failure = Some(failure);
}
}
true
});
}
pub fn record_loop_alive(path: &Path, now: u64) {
update_status(path, "the catalog loop tick", |status| {
status.loop_alive_at = Some(now);
true
});
}
pub fn record_daemon_alive(path: &Path, now: u64) {
update_status(path, "the daemon heartbeat", |status| {
let continuing = status
.daemon_alive_at
.is_some_and(|at| now.saturating_sub(at) <= DAEMON_RUN_GAP_SECS);
if !continuing || status.daemon_alive_since.is_none() {
status.daemon_alive_since = Some(now);
}
status.daemon_alive_at = Some(now);
true
});
}
pub fn record_off(path: &Path, now: u64) {
update_status(path, "the catalog check", |status| {
if status.off_since.is_some() {
return false;
}
status.off_since = Some(now);
true
});
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "state", rename_all = "snake_case")]
pub enum CatalogFreshness {
#[default]
NeverChecked,
Off { since: u64 },
Current {
last_verified_at: Option<u64>,
verified_version: Option<u64>,
},
Stalled {
loop_alive_at: u64,
daemon_alive_at: u64,
},
Stale {
source: Option<String>,
last_verified_at: Option<u64>,
failing_since: u64,
failure: Option<CheckFailure>,
last_error: Option<String>,
},
}
pub fn catalog_freshness(status: &CatalogStatus) -> CatalogFreshness {
if let Some(since) = status.off_since {
return CatalogFreshness::Off { since };
}
let Some(last_attempt) = status.last_attempt_at else {
return CatalogFreshness::NeverChecked;
};
match status.failing_since {
Some(since) if last_attempt.saturating_sub(since) >= STALE_AFTER_SECS => {
CatalogFreshness::Stale {
source: status.source.clone(),
last_verified_at: status.last_verified_at,
failing_since: since,
failure: status.last_failure,
last_error: status.last_error.clone(),
}
}
_ => CatalogFreshness::Current {
last_verified_at: status.last_verified_at,
verified_version: status.verified_version,
},
}
}
pub fn catalog_freshness_with_liveness(status: &CatalogStatus, now: u64) -> CatalogFreshness {
let freshness = catalog_freshness(status);
if matches!(
freshness,
CatalogFreshness::Off { .. } | CatalogFreshness::Stale { .. }
) {
return freshness;
}
if let (Some(loop_alive_at), Some(daemon_alive_at)) =
(status.loop_alive_at, status.daemon_alive_at)
{
let run_start = status.daemon_alive_since.unwrap_or(daemon_alive_at);
if now.saturating_sub(daemon_alive_at) <= DAEMON_ALIVE_SECS
&& daemon_alive_at.saturating_sub(loop_alive_at.max(run_start)) >= LOOP_STALL_SECS
{
return CatalogFreshness::Stalled {
loop_alive_at,
daemon_alive_at,
};
}
}
freshness
}
pub fn retained_path(state_root: &Path) -> PathBuf {
state_root.join(RETAINED_FILE)
}
#[derive(Debug, Default, Serialize, Deserialize)]
struct RetainedFile {
envelopes: Vec<CatalogCacheEnvelope>,
rows: Vec<RetainedRow>,
}
#[derive(Debug, Serialize, Deserialize)]
struct RetainedRow {
id: String,
envelope: usize,
}
pub fn save_retained(
path: &Path,
rows: &[(String, std::sync::Arc<CatalogCacheEnvelope>)],
) -> std::io::Result<()> {
if rows.is_empty() {
return match std::fs::remove_file(path) {
Err(error) if error.kind() != std::io::ErrorKind::NotFound => Err(error),
_ => Ok(()),
};
}
let mut file = RetainedFile::default();
for (id, envelope) in rows {
let index = match file.envelopes.iter().position(|e| e == envelope.as_ref()) {
Some(index) => index,
None => {
file.envelopes.push(envelope.as_ref().clone());
file.envelopes.len() - 1
}
};
file.rows.push(RetainedRow {
id: id.clone(),
envelope: index,
});
}
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let json = serde_json::to_string_pretty(&file).map_err(std::io::Error::other)?;
let tmp = unique_temp_path(path);
std::fs::write(&tmp, json)?;
if let Err(error) = std::fs::rename(&tmp, path) {
let _ = std::fs::remove_file(&tmp);
return Err(error);
}
Ok(())
}
pub fn load_retained(
path: &Path,
public_key_b64: Option<&str>,
) -> Vec<(ModelSchema, std::sync::Arc<CatalogCacheEnvelope>)> {
let Some(public_key_b64) = public_key_b64.map(str::trim).filter(|k| !k.is_empty()) else {
return Vec::new();
};
let Some(file) = std::fs::read_to_string(path)
.ok()
.and_then(|json| serde_json::from_str::<RetainedFile>(&json).ok())
else {
return Vec::new();
};
let verified: Vec<Option<(VerifiedCatalog, std::sync::Arc<CatalogCacheEnvelope>)>> = file
.envelopes
.into_iter()
.map(|envelope| {
verify_signed_catalog(
envelope.signed_body.clone(),
envelope.signature.clone(),
public_key_b64,
)
.ok()
.map(|catalog| (catalog, std::sync::Arc::new(envelope)))
})
.collect();
file.rows
.into_iter()
.filter_map(|row| {
let (catalog, envelope) = verified.get(row.envelope)?.as_ref()?;
let schema = catalog.doc.models.iter().find(|m| m.id == row.id)?.clone();
Some((schema, envelope.clone()))
})
.collect()
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CatalogCacheEnvelope {
pub(crate) signed_body: String,
pub(crate) signature: String,
}
pub fn cache_path(state_root: &Path) -> PathBuf {
state_root.join(CATALOG_CACHE_FILE)
}
pub fn load_cache(path: &Path, public_key_b64: Option<&str>) -> Vec<ModelSchema> {
load_verified(path, public_key_b64)
.map(VerifiedCatalog::into_models)
.unwrap_or_default()
}
pub fn load_verified(path: &Path, public_key_b64: Option<&str>) -> Option<VerifiedCatalog> {
let public_key_b64 = public_key_b64
.map(str::trim)
.filter(|key| !key.is_empty())?;
let json = std::fs::read_to_string(path).ok()?;
let envelope: CatalogCacheEnvelope = serde_json::from_str(&json).ok()?;
verify_signed_catalog(envelope.signed_body, envelope.signature, public_key_b64).ok()
}
pub fn load_doc(path: &Path, public_key_b64: Option<&str>) -> Option<CatalogDoc> {
load_verified(path, public_key_b64).map(|verified| verified.doc)
}
fn verify_signed_catalog(
signed_body: String,
signature: String,
public_key_b64: &str,
) -> Result<VerifiedCatalog, String> {
let mut last_error = String::from("no catalog public key configured");
let verified = public_key_b64
.split(',')
.map(str::trim)
.filter(|key| !key.is_empty())
.any(|key| {
match car_bundle::verify_detached(signed_body.as_bytes(), signature.trim(), key) {
Ok(()) => true,
Err(e) => {
last_error = e.to_string();
false
}
}
});
if !verified {
return Err(format!(
"catalog signature verification failed: {last_error}"
));
}
let raw: RawCatalogDoc =
serde_json::from_str(&signed_body).map_err(|e| format!("parse catalog: {e}"))?;
let doc = CatalogDoc {
version: raw.version,
models: raw
.models
.into_iter()
.filter_map(parse_catalog_row)
.collect(),
revoked: raw.revoked,
};
Ok(VerifiedCatalog {
doc,
signed_body,
signature: signature.trim().to_string(),
})
}
pub async fn fetch_and_verify(
http: &reqwest::Client,
url: &str,
public_key_b64: &str,
) -> Result<VerifiedCatalog, String> {
let (bytes, sig) = fetch_signed(http, url).await?;
verify_fetched(bytes, sig, public_key_b64)
}
pub async fn fetch_signed(http: &reqwest::Client, url: &str) -> Result<(Vec<u8>, String), String> {
let resp = http
.get(url)
.send()
.await
.map_err(|e| format!("fetch catalog: {e}"))?
.error_for_status()
.map_err(|e| format!("fetch catalog: {e}"))?;
let bytes = read_capped(resp, MAX_CATALOG_BYTES, "catalog").await?;
let sig_resp = http
.get(format!("{url}.sig"))
.send()
.await
.map_err(|e| format!("fetch signature: {e}"))?
.error_for_status()
.map_err(|e| format!("fetch signature: {e}"))?;
let sig = String::from_utf8(read_capped(sig_resp, MAX_SIGNATURE_BYTES, "signature").await?)
.map_err(|e| format!("read signature: {e}"))?;
Ok((bytes, sig))
}
pub fn verify_fetched(
bytes: Vec<u8>,
sig: String,
public_key_b64: &str,
) -> Result<VerifiedCatalog, String> {
let signed_body = String::from_utf8(bytes).map_err(|e| format!("parse catalog: {e}"))?;
verify_signed_catalog(signed_body, sig, public_key_b64)
}
pub fn save_verified(path: &Path, verified: &VerifiedCatalog) -> std::io::Result<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let envelope = CatalogCacheEnvelope {
signed_body: verified.signed_body.clone(),
signature: verified.signature.clone(),
};
let json = serde_json::to_string_pretty(&envelope).map_err(std::io::Error::other)?;
let tmp = unique_temp_path(path);
std::fs::write(&tmp, json)?;
if let Err(error) = std::fs::rename(&tmp, path) {
let _ = std::fs::remove_file(&tmp);
return Err(error);
}
Ok(())
}
fn unique_temp_path(path: &Path) -> PathBuf {
static TEMP_SEQUENCE: AtomicU64 = AtomicU64::new(0);
let sequence = TEMP_SEQUENCE.fetch_add(1, Ordering::Relaxed);
let file_name = path
.file_name()
.and_then(|name| name.to_str())
.unwrap_or(CATALOG_CACHE_FILE);
path.with_file_name(format!(
".{file_name}.{}.{}.tmp",
std::process::id(),
sequence
))
}
fn cache_update_lock() -> &'static tokio::sync::Mutex<()> {
static LOCK: OnceLock<tokio::sync::Mutex<()>> = OnceLock::new();
LOCK.get_or_init(|| tokio::sync::Mutex::new(()))
}
fn catalog_lock_path(path: &Path) -> PathBuf {
let mut lock_path = path.as_os_str().to_owned();
lock_path.push(".lock");
PathBuf::from(lock_path)
}
fn acquire_catalog_lock(path: &Path) -> std::io::Result<std::fs::File> {
let lock_path = catalog_lock_path(path);
if let Some(parent) = lock_path.parent() {
std::fs::create_dir_all(parent)?;
}
let lock = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&lock_path)?;
#[cfg(test)]
{
use std::fs::TryLockError;
match lock.try_lock() {
Ok(()) => Ok(lock),
Err(TryLockError::WouldBlock) => {
if let Some(contended_path) = std::env::var_os("CAR_CATALOG_TEST_LOCK_CONTENDED") {
std::fs::write(contended_path, b"contended")?;
}
lock.lock()?;
Ok(lock)
}
Err(TryLockError::Error(error)) => Err(error),
}
}
#[cfg(not(test))]
{
lock.lock()?;
Ok(lock)
}
}
#[cfg(test)]
fn pause_install_after_authenticated_read() -> Result<(), String> {
let Some(ready_path) = std::env::var_os("CAR_CATALOG_TEST_READ_READY") else {
return Ok(());
};
let release_path = std::env::var_os("CAR_CATALOG_TEST_READ_RELEASE")
.ok_or_else(|| "CAR_CATALOG_TEST_READ_RELEASE is required with READ_READY".to_string())?;
std::fs::write(&ready_path, b"ready")
.map_err(|error| format!("signal authenticated catalog read: {error}"))?;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
while !Path::new(&release_path).exists() {
if std::time::Instant::now() >= deadline {
return Err("timed out waiting to release authenticated catalog read".to_string());
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
Ok(())
}
fn install_if_newer_locked(
path: &Path,
verified: &VerifiedCatalog,
public_key_b64: &str,
) -> Result<usize, String> {
let _lock = acquire_catalog_lock(path).map_err(|error| {
format!(
"acquire catalog cache lock {}: {error}",
catalog_lock_path(path).display()
)
})?;
let cached_version = load_doc(path, Some(public_key_b64))
.map(|doc| doc.version)
.unwrap_or(0);
#[cfg(test)]
pause_install_after_authenticated_read()?;
if verified.version() <= cached_version {
return Err(format!(
"catalog version {} is not newer than the cached version {} (rejected)",
verified.version(),
cached_version
));
}
let count = verified.model_count();
save_verified(path, verified).map_err(|e| e.to_string())?;
Ok(count)
}
pub async fn install_if_newer(
path: &Path,
verified: &VerifiedCatalog,
public_key_b64: &str,
) -> Result<usize, String> {
let _guard = cache_update_lock().lock().await;
let path = path.to_path_buf();
let verified = verified.clone();
let public_key_b64 = public_key_b64.to_string();
tokio::task::spawn_blocking(move || install_if_newer_locked(&path, &verified, &public_key_b64))
.await
.map_err(|error| format!("catalog cache install task failed: {error}"))?
}
#[cfg(test)]
pub(crate) fn signed_test_catalog(doc: CatalogDoc, seed: u8) -> (VerifiedCatalog, String) {
use base64::Engine;
use ed25519_dalek::{Signer, SigningKey};
let signed_body = serde_json::to_string(&doc).unwrap();
let signing_key = SigningKey::from_bytes(&[seed; 32]);
let signature = base64::engine::general_purpose::STANDARD
.encode(signing_key.sign(signed_body.as_bytes()).to_bytes());
let public_key =
base64::engine::general_purpose::STANDARD.encode(signing_key.verifying_key().as_bytes());
let verified = verify_signed_catalog(signed_body, signature, &public_key).unwrap();
(verified, public_key)
}
#[cfg(test)]
mod tests {
#[test]
fn catalog_freshness_counts_failed_checks_not_time() {
let day = 24 * 60 * 60;
let now = 100 * day;
assert_eq!(
catalog_freshness(&CatalogStatus::default()),
CatalogFreshness::NeverChecked
);
let verified = CatalogStatus {
source: Some("https://x/catalog.json".into()),
last_attempt_at: Some(now - 30 * day),
last_verified_at: Some(now - 30 * day),
verified_version: Some(7),
..CatalogStatus::default()
};
assert_eq!(
catalog_freshness(&verified),
CatalogFreshness::Current {
last_verified_at: Some(now - 30 * day),
verified_version: Some(7),
},
);
let failing = |since: u64, last_attempt: u64| CatalogStatus {
last_attempt_at: Some(last_attempt),
failing_since: Some(since),
last_failure: Some(CheckFailure::Unreachable),
last_error: Some("fetch catalog: timed out".into()),
..verified.clone()
};
assert!(matches!(
catalog_freshness(&failing(now - 30 * day, now - 30 * day)),
CatalogFreshness::Current { .. }
));
assert!(matches!(
catalog_freshness(&failing(now - day, now)),
CatalogFreshness::Current { .. }
));
assert_eq!(
catalog_freshness(&failing(now - STALE_AFTER_SECS, now)),
CatalogFreshness::Stale {
source: Some("https://x/catalog.json".into()),
last_verified_at: Some(now - 30 * day),
failing_since: now - STALE_AFTER_SECS,
failure: Some(CheckFailure::Unreachable),
last_error: Some("fetch catalog: timed out".into()),
}
);
let off = CatalogStatus {
off_since: Some(now - 5 * day),
..failing(now - 10 * day, now)
};
assert_eq!(
catalog_freshness(&off),
CatalogFreshness::Off {
since: now - 5 * day
}
);
}
#[test]
fn a_loop_trailing_its_daemon_heartbeat_has_stalled() {
let hour = 60 * 60;
let now = 1_000 * hour;
let alive = |loop_at: Option<u64>, daemon_at: Option<u64>| CatalogStatus {
last_attempt_at: Some(now - 10 * hour),
last_verified_at: Some(now - 10 * hour),
loop_alive_at: loop_at,
daemon_alive_at: daemon_at,
daemon_alive_since: Some(now - 24 * hour),
..CatalogStatus::default()
};
assert_eq!(
catalog_freshness_with_liveness(&alive(Some(now - 3 * hour), Some(now - 60)), now),
CatalogFreshness::Stalled {
loop_alive_at: now - 3 * hour,
daemon_alive_at: now - 60,
}
);
assert!(matches!(
catalog_freshness_with_liveness(&alive(Some(now - 900), Some(now - 60)), now),
CatalogFreshness::Current { .. }
));
let woke = CatalogStatus {
daemon_alive_since: Some(now - 60),
..alive(Some(now - 72 * hour), Some(now - 60))
};
assert!(matches!(
catalog_freshness_with_liveness(&woke, now),
CatalogFreshness::Current { .. }
));
assert!(matches!(
catalog_freshness_with_liveness(
&alive(Some(now - 9 * hour), Some(now - 2 * hour)),
now
),
CatalogFreshness::Current { .. }
));
assert!(matches!(
catalog_freshness_with_liveness(&alive(None, Some(now - 60)), now),
CatalogFreshness::Current { .. }
));
let off = CatalogStatus {
off_since: Some(now - 9 * hour),
..alive(Some(now - 3 * hour), Some(now - 60))
};
assert!(matches!(
catalog_freshness_with_liveness(&off, now),
CatalogFreshness::Off { .. }
));
}
#[test]
fn concurrent_status_writers_keep_every_field() {
let dir = tempfile::tempdir().unwrap();
let path = std::sync::Arc::new(status_path(dir.path()));
record_check(&path, "u", Err((CheckFailure::Unreachable, "refused")), 10);
let writers: Vec<_> = (0..4)
.map(|w| {
let path = path.clone();
std::thread::spawn(move || {
for i in 0..200u64 {
match w {
0 => record_daemon_alive(&path, 1_000 + i),
1 => record_loop_alive(&path, 1_000 + i),
_ => record_check(
&path,
"u",
Err((CheckFailure::Unreachable, "refused")),
20 + i,
),
}
}
})
})
.collect();
for writer in writers {
writer.join().unwrap();
}
let status = load_status(&path);
assert_eq!(status.failing_since, Some(10), "{status:?}");
assert_eq!(status.daemon_alive_at, Some(1_199));
assert_eq!(status.loop_alive_at, Some(1_199));
}
#[test]
fn a_gap_in_the_daemon_heartbeat_starts_a_new_run() {
let dir = tempfile::tempdir().unwrap();
let path = status_path(dir.path());
record_daemon_alive(&path, 1_000);
record_daemon_alive(&path, 1_000 + 900);
assert_eq!(load_status(&path).daemon_alive_since, Some(1_000));
record_daemon_alive(&path, 1_000 + 900 + DAEMON_RUN_GAP_SECS + 1);
let status = load_status(&path);
assert_eq!(status.daemon_alive_since, status.daemon_alive_at);
}
#[test]
fn recording_checks_tracks_the_first_failure_and_clears_on_success() {
let dir = tempfile::tempdir().unwrap();
let path = status_path(dir.path());
record_off(&path, 5);
assert_eq!(load_status(&path).off_since, Some(5));
record_check(
&path,
"u",
Err((CheckFailure::Unreachable, "fetch catalog: refused")),
10,
);
record_check(
&path,
"u",
Err((CheckFailure::Rejected, "does not verify")),
20,
);
let status = load_status(&path);
assert_eq!(status.off_since, None);
assert_eq!(status.failing_since, Some(10));
assert_eq!(status.last_failure, Some(CheckFailure::Rejected));
assert_eq!(status.last_error.as_deref(), Some("does not verify"));
record_check(&path, "https://x/c.json?token=secret", Ok(9), 30);
let status = load_status(&path);
assert_eq!(status.failing_since, None);
assert_eq!(status.last_failure, None);
assert_eq!(status.last_error, None);
assert_eq!(status.last_verified_at, Some(30));
assert_eq!(status.verified_version, Some(9));
assert_eq!(status.source.as_deref(), Some("https://x/c.json"));
}
#[tokio::test]
async fn an_oversized_streamed_catalog_is_refused_without_reading_it_all() {
use std::io::Write;
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let url = format!("http://{}/catalog.json", listener.local_addr().unwrap());
std::thread::spawn(move || {
if let Ok((mut stream, _)) = listener.accept() {
let mut request = [0u8; 1024];
let _ = std::io::Read::read(&mut stream, &mut request);
let _ = stream.write_all(b"HTTP/1.1 200 OK\r\nConnection: close\r\n\r\n");
let chunk = vec![b'x'; 64 * 1024];
for _ in 0..(MAX_CATALOG_BYTES / chunk.len() as u64 + 8) {
if stream.write_all(&chunk).is_err() {
return;
}
}
}
});
let http = crate::tls_client::catalog_refresh_client();
let error = fetch_and_verify(&http, &url, DEFAULT_CATALOG_PUBKEY)
.await
.unwrap_err();
assert!(error.contains("too large"), "{error}");
}
#[test]
fn an_unreadable_or_too_new_row_is_skipped_not_the_whole_catalog() {
let good = crate::registry::builtin_catalog().remove(0);
let mut unknown = serde_json::to_value(&good).unwrap();
unknown["id"] = "future/unknown-capability".into();
unknown["capabilities"] = serde_json::json!(["telepathy"]);
let mut too_new = serde_json::to_value(&good).unwrap();
too_new["id"] = "future/needs-newer-car".into();
too_new["min_car_version"] = "999.0.0".into();
let mut old_enough = serde_json::to_value(&good).unwrap();
old_enough["id"] = "future/old-enough".into();
old_enough["min_car_version"] = "0.0.1".into();
let body = serde_json::json!({
"version": 4,
"models": [serde_json::to_value(&good).unwrap(), unknown, too_new, old_enough],
})
.to_string();
use base64::Engine;
use ed25519_dalek::{Signer, SigningKey};
let key = SigningKey::from_bytes(&[41; 32]);
let signature =
base64::engine::general_purpose::STANDARD.encode(key.sign(body.as_bytes()).to_bytes());
let public =
base64::engine::general_purpose::STANDARD.encode(key.verifying_key().as_bytes());
let verified = verify_signed_catalog(body, signature, &public).unwrap();
let ids: Vec<String> = verified.into_models().into_iter().map(|m| m.id).collect();
assert_eq!(ids, vec![good.id, "future/old-enough".to_string()]);
}
#[test]
fn versions_compare_numerically() {
assert!(version_at_least("0.57.0", "0.57.0"));
assert!(version_at_least("0.57.1", "0.57"));
assert!(version_at_least("0.100.0", "0.99.9"));
assert!(!version_at_least("0.56.9", "0.57.0"));
assert!(version_at_least("0.57.0-nightly.20261002", "0.57.0"));
assert!(!version_at_least("0.57.0", "garbage"));
}
#[test]
fn a_catalog_verifies_against_any_key_in_the_list_and_carries_its_revocations() {
let (signed_by_new, new_key) = signed_test_catalog(
CatalogDoc {
version: 3,
models: vec![],
revoked: vec!["mlx/bad:4bit".into()],
},
51,
);
let (_, old_key) = signed_test_catalog(
CatalogDoc {
version: 1,
models: vec![],
revoked: Vec::new(),
},
52,
);
let envelope = signed_by_new.envelope();
let verify = |keys: &str| {
verify_signed_catalog(
envelope.signed_body.clone(),
envelope.signature.clone(),
keys,
)
};
let both = verify(&format!("{old_key}, {new_key}")).expect("either key verifies");
assert_eq!(both.revoked(), ["mlx/bad:4bit".to_string()]);
assert!(
verify(&old_key).is_err(),
"the old key alone did not sign it"
);
assert!(verify(" , ").is_err(), "no key, no catalog");
}
#[tokio::test]
async fn a_version_pinned_by_a_delisted_key_does_not_block_a_listed_one() {
let tmp = tempfile::tempdir().unwrap();
let path = cache_path(tmp.path());
let doc = |version| CatalogDoc {
version,
models: vec![],
revoked: Vec::new(),
};
let (hostile, _compromised_key) = signed_test_catalog(doc(u64::MAX / 2), 71);
save_verified(&path, &hostile).unwrap();
let (legitimate, current_key) = signed_test_catalog(doc(5), 72);
assert_eq!(
install_if_newer(&path, &legitimate, ¤t_key).await,
Ok(0)
);
assert_eq!(
load_doc(&path, Some(¤t_key)).map(|d| d.version),
Some(5)
);
}
#[test]
fn the_default_catalog_key_is_a_well_formed_ed25519_key() {
use base64::Engine;
let bogus = base64::engine::general_purpose::STANDARD.encode([0u8; 64]);
let error = car_bundle::verify_detached(b"{}", &bogus, DEFAULT_CATALOG_PUBKEY).unwrap_err();
assert!(
matches!(error, car_bundle::BundleError::SignatureInvalid(_)),
"{error}"
);
let truncated = &DEFAULT_CATALOG_PUBKEY[..20];
assert!(matches!(
car_bundle::verify_detached(b"{}", &bogus, truncated).unwrap_err(),
car_bundle::BundleError::KeyMalformed(_)
));
assert!(DEFAULT_CATALOG_URL.starts_with("https://"));
}
use super::*;
use std::process::{Child, Command, ExitStatus};
use std::time::{Duration, Instant};
const CROSS_PROCESS_TEST: &str =
"catalog::tests::cross_process_installs_preserve_highest_authenticated_version";
fn wait_for_path(path: &Path, timeout: Duration) -> Result<(), String> {
let deadline = Instant::now() + timeout;
while !path.exists() {
if Instant::now() >= deadline {
return Err(format!("timed out waiting for {}", path.display()));
}
std::thread::sleep(Duration::from_millis(10));
}
Ok(())
}
fn wait_for_child(child: &mut Child, timeout: Duration) -> Result<ExitStatus, String> {
let deadline = Instant::now() + timeout;
loop {
match child.try_wait() {
Ok(Some(status)) => return Ok(status),
Ok(None) if Instant::now() < deadline => {
std::thread::sleep(Duration::from_millis(10));
}
Ok(None) => {
let _ = child.kill();
let _ = child.wait();
return Err("timed out waiting for catalog child process".to_string());
}
Err(error) => {
let _ = child.kill();
let _ = child.wait();
return Err(format!("wait for catalog child process: {error}"));
}
}
}
}
fn spawn_catalog_child(
cache_path: &Path,
version: u64,
ready_path: Option<&Path>,
release_path: Option<&Path>,
contended_path: Option<&Path>,
) -> Child {
let mut command = Command::new(std::env::current_exe().unwrap());
command
.arg(CROSS_PROCESS_TEST)
.arg("--exact")
.arg("--nocapture")
.arg("--test-threads=1")
.env("CAR_CATALOG_TEST_CHILD_VERSION", version.to_string())
.env("CAR_CATALOG_TEST_CHILD_CACHE", cache_path);
if let Some(path) = ready_path {
command.env("CAR_CATALOG_TEST_READ_READY", path);
}
if let Some(path) = release_path {
command.env("CAR_CATALOG_TEST_READ_RELEASE", path);
}
if let Some(path) = contended_path {
command.env("CAR_CATALOG_TEST_LOCK_CONTENDED", path);
}
command.spawn().expect("spawn catalog test child")
}
#[test]
fn missing_cache_is_empty() {
let p = std::env::temp_dir().join("car-catalog-none-xyz.json");
let _ = std::fs::remove_file(&p);
assert!(load_cache(&p, None).is_empty());
}
#[test]
fn cache_path_sits_directly_under_the_state_root() {
let p = cache_path(Path::new("/home/u/.car"));
assert_eq!(p, Path::new("/home/u/.car/catalog-cache.json"));
let relocated = cache_path(Path::new("/tmp/car-alt"));
assert_eq!(relocated, Path::new("/tmp/car-alt/catalog-cache.json"));
}
#[test]
fn atomic_temp_paths_are_unique_siblings() {
let path = Path::new("/home/u/.car/catalog-cache.json");
let first = unique_temp_path(path);
let second = unique_temp_path(path);
assert_eq!(first.parent(), path.parent());
assert_eq!(second.parent(), path.parent());
assert_ne!(first, second);
assert!(first
.file_name()
.unwrap()
.to_string_lossy()
.ends_with(".tmp"));
assert!(second
.file_name()
.unwrap()
.to_string_lossy()
.ends_with(".tmp"));
}
#[test]
fn valid_signed_envelope_round_trips_curated_models() {
let tmp = tempfile::tempdir().unwrap();
let path = cache_path(tmp.path());
let schema = crate::openrouter::builtin_schemas()
.into_iter()
.next()
.expect("managed OpenRouter schema");
let (verified, public_key) = signed_test_catalog(
CatalogDoc {
revoked: Vec::new(),
version: 7,
models: vec![schema.clone()],
},
7,
);
save_verified(&path, &verified).unwrap();
let loaded = load_cache(&path, Some(&public_key));
assert_eq!(loaded.len(), 1);
assert_eq!(loaded[0].id, schema.id);
assert_eq!(loaded[0].trust_tier, crate::schema::TrustTier::Curated);
assert_eq!(
load_doc(&path, Some(&public_key))
.map(|doc| doc.version)
.unwrap(),
7
);
}
#[tokio::test]
async fn tampered_body_or_version_and_wrong_or_missing_key_fail_closed() {
let tmp = tempfile::tempdir().unwrap();
let path = cache_path(tmp.path());
let schema = crate::openrouter::builtin_schemas()
.into_iter()
.next()
.expect("managed OpenRouter schema");
let (verified, public_key) = signed_test_catalog(
CatalogDoc {
revoked: Vec::new(),
version: 7,
models: vec![schema],
},
11,
);
let (_, wrong_key) = signed_test_catalog(
CatalogDoc {
revoked: Vec::new(),
version: 1,
models: vec![],
},
12,
);
save_verified(&path, &verified).unwrap();
assert!(load_cache(&path, None).is_empty());
assert!(load_cache(&path, Some("")).is_empty());
assert!(load_cache(&path, Some(&wrong_key)).is_empty());
let mut envelope: CatalogCacheEnvelope =
serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap();
envelope.signed_body = envelope
.signed_body
.replace("\"version\":7", "\"version\":18446744073709551615");
std::fs::write(&path, serde_json::to_vec_pretty(&envelope).unwrap()).unwrap();
assert!(load_cache(&path, Some(&public_key)).is_empty());
assert!(
load_doc(&path, Some(&public_key)).is_none(),
"a tampered cached version must not participate in anti-rollback comparison"
);
let (newer, same_public_key) = signed_test_catalog(
CatalogDoc {
revoked: Vec::new(),
version: 8,
models: vec![],
},
11,
);
assert_eq!(public_key, same_public_key);
assert_eq!(
install_if_newer(&path, &newer, &public_key).await,
Ok(0),
"a forged high cached version must not block a newer authenticated refresh"
);
assert_eq!(
load_doc(&path, Some(&public_key)).map(|doc| doc.version),
Some(8)
);
}
#[test]
fn invalid_or_missing_signature_and_legacy_plain_json_fail_closed() {
let tmp = tempfile::tempdir().unwrap();
let path = cache_path(tmp.path());
let doc = CatalogDoc {
revoked: Vec::new(),
version: u64::MAX,
models: crate::openrouter::builtin_schemas()
.into_iter()
.take(1)
.collect(),
};
let (verified, public_key) = signed_test_catalog(doc.clone(), 21);
save_verified(&path, &verified).unwrap();
let mut value: serde_json::Value =
serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap();
value["signature"] = serde_json::Value::String("AAAA".to_string());
std::fs::write(&path, serde_json::to_vec_pretty(&value).unwrap()).unwrap();
assert!(load_cache(&path, Some(&public_key)).is_empty());
save_verified(&path, &verified).unwrap();
let mut value: serde_json::Value =
serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap();
value.as_object_mut().unwrap().remove("signature");
std::fs::write(&path, serde_json::to_vec_pretty(&value).unwrap()).unwrap();
assert!(load_cache(&path, Some(&public_key)).is_empty());
std::fs::write(&path, serde_json::to_vec_pretty(&doc).unwrap()).unwrap();
assert!(load_cache(&path, Some(&public_key)).is_empty());
assert!(load_doc(&path, Some(&public_key)).is_none());
}
#[tokio::test]
async fn concurrent_installs_leave_the_highest_authenticated_version() {
let tmp = tempfile::tempdir().unwrap();
let path = cache_path(tmp.path());
let (version_n, public_key) = signed_test_catalog(
CatalogDoc {
revoked: Vec::new(),
version: 41,
models: vec![],
},
31,
);
let (version_n_plus_one, same_public_key) = signed_test_catalog(
CatalogDoc {
revoked: Vec::new(),
version: 42,
models: vec![],
},
31,
);
assert_eq!(public_key, same_public_key);
let (_lower, higher) = tokio::join!(
install_if_newer(&path, &version_n, &public_key),
install_if_newer(&path, &version_n_plus_one, &public_key)
);
assert!(higher.is_ok());
assert_eq!(
load_doc(&path, Some(&public_key)).map(|doc| doc.version),
Some(42)
);
let replay = install_if_newer(&path, &version_n, &public_key).await;
assert!(replay.is_err());
assert_eq!(
load_doc(&path, Some(&public_key)).map(|doc| doc.version),
Some(42)
);
}
#[test]
fn cross_process_installs_preserve_highest_authenticated_version() {
if let Some(version) = std::env::var_os("CAR_CATALOG_TEST_CHILD_VERSION") {
let version = version
.to_string_lossy()
.parse::<u64>()
.expect("child catalog version");
let cache_path = PathBuf::from(
std::env::var_os("CAR_CATALOG_TEST_CHILD_CACHE").expect("child catalog cache path"),
);
let (verified, public_key) = signed_test_catalog(
CatalogDoc {
revoked: Vec::new(),
version,
models: vec![],
},
31,
);
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
runtime
.block_on(install_if_newer(&cache_path, &verified, &public_key))
.expect("child catalog install");
return;
}
let tmp = tempfile::tempdir().unwrap();
let path = cache_path(tmp.path());
let ready = tmp.path().join("version-41-read");
let release = tmp.path().join("release-version-41");
let (initial, public_key) = signed_test_catalog(
CatalogDoc {
revoked: Vec::new(),
version: 40,
models: vec![],
},
31,
);
save_verified(&path, &initial).unwrap();
let contended = tmp.path().join("version-42-contended");
let mut version_41 = spawn_catalog_child(&path, 41, Some(&ready), Some(&release), None);
if let Err(error) = wait_for_path(&ready, Duration::from_secs(10)) {
let _ = version_41.kill();
let _ = version_41.wait();
panic!("{error}");
}
let mut version_42 = spawn_catalog_child(&path, 42, None, None, Some(&contended));
let contention_status = wait_for_path(&contended, Duration::from_secs(10));
std::fs::write(&release, b"release").unwrap();
let version_42_status = wait_for_child(&mut version_42, Duration::from_secs(10));
let version_41_status = wait_for_child(&mut version_41, Duration::from_secs(10));
assert!(
contention_status.is_ok(),
"version 42 never contended on the cross-process cache lock: {contention_status:?}"
);
assert!(
version_42_status.unwrap().success(),
"version 42 child failed"
);
assert!(
version_41_status.unwrap().success(),
"version 41 child failed"
);
assert_eq!(
load_doc(&path, Some(&public_key)).map(|doc| doc.version),
Some(42),
"a stale cross-process writer must not replace a newer authenticated catalog"
);
let (stale, same_public_key) = signed_test_catalog(
CatalogDoc {
revoked: Vec::new(),
version: 41,
models: vec![],
},
31,
);
assert_eq!(public_key, same_public_key);
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
let stale_result = runtime.block_on(install_if_newer(&path, &stale, &public_key));
assert!(
stale_result.is_err(),
"an authenticated but stale replay must be rejected after the process race"
);
assert_eq!(
load_doc(&path, Some(&public_key)).map(|doc| doc.version),
Some(42)
);
}
}