use serde::{Deserialize, Serialize};
use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
#[derive(Debug, thiserror::Error)]
pub enum RegistryError {
#[error("invalid agent name (must be non-empty, alphanumeric + `-_.`): {0:?}")]
InvalidName(String),
#[error("could not resolve home directory")]
NoHomeDir,
#[error("registry I/O error: {0}")]
Io(#[from] std::io::Error),
#[error("registry JSON error: {0}")]
Json(#[from] serde_json::Error),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum AgentStatus {
Running,
#[default]
Idle,
Errored,
Stopping,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentEntry {
pub name: String,
pub dashboard_url: String,
#[serde(default)]
pub status: AgentStatus,
#[serde(default)]
pub display_name: Option<String>,
#[serde(default)]
pub port: Option<u16>,
#[serde(default)]
pub pid: Option<u32>,
#[serde(default)]
pub registered_at: u64,
#[serde(default)]
pub last_heartbeat_at: u64,
}
impl AgentEntry {
pub fn new(name: impl Into<String>, dashboard_url: impl Into<String>) -> Self {
let now = now_secs();
Self {
name: name.into(),
dashboard_url: dashboard_url.into(),
status: AgentStatus::default(),
display_name: None,
port: None,
pid: None,
registered_at: now,
last_heartbeat_at: now,
}
}
pub fn with_display_name(mut self, label: impl Into<String>) -> Self {
self.display_name = Some(label.into());
self
}
pub fn with_port(mut self, port: u16) -> Self {
self.port = Some(port);
self
}
pub fn with_pid(mut self, pid: u32) -> Self {
self.pid = Some(pid);
self
}
pub fn with_status(mut self, status: AgentStatus) -> Self {
self.status = status;
self
}
}
#[derive(Debug, Clone)]
pub struct AgentRegistry {
dir: PathBuf,
}
impl AgentRegistry {
pub fn user_default() -> Result<Self, RegistryError> {
let home = std::env::var_os("HOME")
.or_else(|| std::env::var_os("USERPROFILE"))
.ok_or(RegistryError::NoHomeDir)?;
let dir = PathBuf::from(home).join(".car").join("registry");
Self::open(dir)
}
pub fn open(dir: impl Into<PathBuf>) -> Result<Self, RegistryError> {
let dir = dir.into();
std::fs::create_dir_all(&dir)?;
Ok(Self { dir })
}
pub fn dir(&self) -> &Path {
&self.dir
}
pub fn register(&self, entry: &AgentEntry) -> Result<(), RegistryError> {
let path = self.entry_path(&entry.name)?;
let mut entry = entry.clone();
if entry.registered_at == 0 {
entry.registered_at = now_secs();
}
if entry.last_heartbeat_at == 0 {
entry.last_heartbeat_at = entry.registered_at;
}
write_json_atomic(&path, &entry)?;
Ok(())
}
pub fn heartbeat(&self, name: &str) -> Result<bool, RegistryError> {
let path = self.entry_path(name)?;
if !path.exists() {
return Ok(false);
}
let bytes = std::fs::read(&path)?;
let mut entry: AgentEntry = serde_json::from_slice(&bytes)?;
entry.last_heartbeat_at = now_secs();
write_json_atomic(&path, &entry)?;
Ok(true)
}
pub fn unregister(&self, name: &str) -> Result<(), RegistryError> {
let path = self.entry_path(name)?;
if path.exists() {
std::fs::remove_file(&path)?;
}
Ok(())
}
pub fn list(&self) -> Result<Vec<AgentEntry>, RegistryError> {
let mut out = Vec::new();
for entry in std::fs::read_dir(&self.dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) != Some("json") {
continue;
}
if let Ok(bytes) = std::fs::read(&path) {
if let Ok(parsed) = serde_json::from_slice::<AgentEntry>(&bytes) {
out.push(parsed);
}
}
}
out.sort_by(|a, b| a.name.cmp(&b.name));
Ok(out)
}
pub fn reap_stale(&self, max_age_secs: u64) -> Result<Vec<String>, RegistryError> {
let cutoff = now_secs().saturating_sub(max_age_secs);
let entries = self.list()?;
let mut reaped = Vec::new();
for entry in entries {
if entry.last_heartbeat_at < cutoff {
if let Ok(path) = self.entry_path(&entry.name) {
if path.exists() {
let _ = std::fs::remove_file(&path);
reaped.push(entry.name);
}
}
}
}
Ok(reaped)
}
fn entry_path(&self, name: &str) -> Result<PathBuf, RegistryError> {
validate_name(name)?;
Ok(self.dir.join(format!("{}.json", name)))
}
}
fn validate_name(name: &str) -> Result<(), RegistryError> {
if name.is_empty() {
return Err(RegistryError::InvalidName(name.to_string()));
}
if !name
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '.')
{
return Err(RegistryError::InvalidName(name.to_string()));
}
if name == "." || name == ".." {
return Err(RegistryError::InvalidName(name.to_string()));
}
Ok(())
}
fn now_secs() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
fn write_json_atomic<T: Serialize>(path: &Path, value: &T) -> Result<(), RegistryError> {
let parent = path.parent().ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"registry path has no parent",
)
})?;
let tmp = parent.join(format!(
".{}.tmp",
path.file_name()
.and_then(|s| s.to_str())
.unwrap_or("registry-write")
));
let json = serde_json::to_vec_pretty(value)?;
std::fs::write(&tmp, json)?;
std::fs::rename(&tmp, path)?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn temp_registry() -> (tempfile::TempDir, AgentRegistry) {
let tmp = tempfile::TempDir::new().unwrap();
let reg = AgentRegistry::open(tmp.path()).unwrap();
(tmp, reg)
}
#[test]
fn register_then_list_round_trips() {
let (_tmp, reg) = temp_registry();
reg.register(
&AgentEntry::new("trader-paper", "http://127.0.0.1:8731")
.with_display_name("Trader (paper)")
.with_port(8731)
.with_status(AgentStatus::Running),
)
.unwrap();
let listed = reg.list().unwrap();
assert_eq!(listed.len(), 1);
assert_eq!(listed[0].name, "trader-paper");
assert_eq!(listed[0].port, Some(8731));
assert_eq!(listed[0].status, AgentStatus::Running);
assert!(listed[0].registered_at > 0);
}
#[test]
fn heartbeat_bumps_timestamp() {
let (_tmp, reg) = temp_registry();
reg.register(&AgentEntry::new("a", "http://x")).unwrap();
let before = reg.list().unwrap()[0].last_heartbeat_at;
std::thread::sleep(std::time::Duration::from_secs(1));
let touched = reg.heartbeat("a").unwrap();
assert!(touched);
let after = reg.list().unwrap()[0].last_heartbeat_at;
assert!(after >= before);
}
#[test]
fn heartbeat_unknown_returns_false() {
let (_tmp, reg) = temp_registry();
assert!(!reg.heartbeat("nobody-home").unwrap());
}
#[test]
fn unregister_idempotent() {
let (_tmp, reg) = temp_registry();
reg.register(&AgentEntry::new("x", "http://x")).unwrap();
reg.unregister("x").unwrap();
assert!(reg.list().unwrap().is_empty());
reg.unregister("x").unwrap();
}
#[test]
fn reap_stale_removes_old_entries() {
let (_tmp, reg) = temp_registry();
let mut old = AgentEntry::new("zombie", "http://z");
old.last_heartbeat_at = 1; reg.register(&old).unwrap();
let fresh = AgentEntry::new("alive", "http://a");
reg.register(&fresh).unwrap();
let reaped = reg.reap_stale(60).unwrap();
assert_eq!(reaped, vec!["zombie"]);
let remaining = reg.list().unwrap();
assert_eq!(remaining.len(), 1);
assert_eq!(remaining[0].name, "alive");
}
#[test]
fn invalid_names_rejected() {
let (_tmp, reg) = temp_registry();
assert!(reg
.register(&AgentEntry::new("..", "http://x"))
.is_err());
assert!(reg
.register(&AgentEntry::new("evil/name", "http://x"))
.is_err());
assert!(reg
.register(&AgentEntry::new("space name", "http://x"))
.is_err());
assert!(reg
.register(&AgentEntry::new("", "http://x"))
.is_err());
}
#[test]
fn corrupt_entries_skipped_in_list() {
let (tmp, reg) = temp_registry();
reg.register(&AgentEntry::new("ok", "http://o")).unwrap();
std::fs::write(tmp.path().join("garbage.json"), b"not json").unwrap();
let entries = reg.list().unwrap();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].name, "ok");
}
}