use serde::{Deserialize, Serialize};
use std::fs;
use std::io::Write;
use std::path::{Path, PathBuf};
use std::time::{Duration, SystemTime};
pub fn status_dir() -> PathBuf {
#[cfg(test)]
{
if let Some(p) = test_override::get() {
return p;
}
}
crate::config::opencrabs_home().join("tmp").join("detached")
}
pub fn legacy_dir() -> PathBuf {
status_dir().with_file_name("subagents")
}
#[cfg(test)]
pub(crate) mod test_override {
use std::cell::RefCell;
use std::path::PathBuf;
thread_local! {
static DIR: RefCell<Option<PathBuf>> = const { RefCell::new(None) };
}
pub fn set(p: PathBuf) {
DIR.with(|d| *d.borrow_mut() = Some(p));
}
pub fn get() -> Option<PathBuf> {
DIR.with(|d| d.borrow().clone())
}
pub fn clear() {
DIR.with(|d| *d.borrow_mut() = None);
}
}
pub fn ensure_dir() -> std::io::Result<()> {
let dir = status_dir();
if !dir.exists() {
fs::create_dir_all(&dir)?;
}
Ok(())
}
pub fn status_path(id: &str) -> PathBuf {
status_dir().join(format!("{}.json", id))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum WorkKind {
Command,
Agent,
}
impl Default for WorkKind {
fn default() -> Self {
WorkKind::Agent
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub enum WorkState {
Pending,
Running,
AwaitingInput,
Completed,
Failed,
Interrupted,
}
impl WorkState {
pub fn is_terminal(&self) -> bool {
matches!(
self,
WorkState::Completed | WorkState::Failed | WorkState::Interrupted
)
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct ProgressSnapshot {
#[serde(default = "usize::default")]
pub iteration: usize,
#[serde(default)]
pub last_tool: Option<String>,
#[serde(default)]
pub last_event: Option<String>,
#[serde(default)]
pub updated_at: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkFinish {
pub completed_at: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub success: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub code: Option<i32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub elapsed_secs: Option<f32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_bytes: Option<usize>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_full: Option<String>,
}
#[derive(Debug, Clone, Copy)]
pub struct CommandExit {
pub success: bool,
pub code: i32,
pub elapsed_secs: f32,
pub output_bytes: usize,
}
fn finish_now() -> WorkFinish {
WorkFinish {
completed_at: now_rfc3339(),
success: None,
code: None,
elapsed_secs: None,
output_bytes: None,
error: None,
output_full: None,
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkStatus {
pub id: String,
#[serde(default)]
pub kind: WorkKind,
pub session_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent_session_id: Option<String>,
pub label: String,
pub task: String,
pub spawned_at: String,
pub state: WorkState,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub progress: Option<ProgressSnapshot>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finish: Option<WorkFinish>,
}
impl WorkStatus {
pub fn new_agent(
id: &str,
label: &str,
session_id: &str,
prompt: &str,
parent_session_id: Option<&str>,
) -> std::io::Result<Self> {
ensure_dir()?;
let status = Self {
id: id.to_string(),
kind: WorkKind::Agent,
session_id: session_id.to_string(),
parent_session_id: parent_session_id.map(str::to_string),
label: label.to_string(),
task: prompt.to_string(),
spawned_at: now_rfc3339(),
state: WorkState::Pending,
progress: None,
finish: None,
};
status.write()?;
Ok(status)
}
pub fn new_command(
id: &str,
session_id: &str,
label: &str,
command: &str,
) -> std::io::Result<Self> {
ensure_dir()?;
let status = Self {
id: id.to_string(),
kind: WorkKind::Command,
session_id: session_id.to_string(),
label: label.to_string(),
task: command.to_string(),
spawned_at: now_rfc3339(),
state: WorkState::Running,
progress: None,
finish: None,
parent_session_id: None,
};
status.write()?;
Ok(status)
}
pub fn finish_command(
id: &str,
session_id: &str,
label: &str,
command: &str,
exit: CommandExit,
) -> std::io::Result<()> {
let mut status = Self::read(id).unwrap_or_else(|| Self {
id: id.to_string(),
kind: WorkKind::Command,
session_id: session_id.to_string(),
label: label.to_string(),
task: command.to_string(),
spawned_at: now_rfc3339(),
state: WorkState::Running,
progress: None,
finish: None,
parent_session_id: None,
});
status.state = if exit.success {
WorkState::Completed
} else {
WorkState::Failed
};
let mut finish = finish_now();
finish.success = Some(exit.success);
finish.code = Some(exit.code);
finish.elapsed_secs = Some(exit.elapsed_secs);
finish.output_bytes = Some(exit.output_bytes);
status.finish = Some(finish);
status.write()
}
pub fn mark_running(&mut self) -> std::io::Result<()> {
self.state = WorkState::Running;
self.write()
}
pub fn mark_awaiting_input(&mut self) -> std::io::Result<()> {
self.state = WorkState::AwaitingInput;
self.write()
}
pub fn update_progress(
&mut self,
iteration: usize,
last_tool: Option<String>,
last_event: Option<String>,
) -> std::io::Result<()> {
self.progress = Some(ProgressSnapshot {
iteration,
last_tool,
last_event,
updated_at: Some(now_rfc3339()),
});
self.write()
}
pub fn mark_completed(&mut self, output_full: String) -> std::io::Result<()> {
self.state = WorkState::Completed;
let mut finish = finish_now();
finish.output_full = Some(output_full);
self.finish = Some(finish);
self.write()
}
pub fn mark_failed(&mut self, error: String) -> std::io::Result<()> {
self.state = WorkState::Failed;
let mut finish = finish_now();
finish.error = Some(error);
self.finish = Some(finish);
self.write()
}
pub fn mark_interrupted(&mut self) -> std::io::Result<()> {
self.state = WorkState::Interrupted;
let subject = match self.kind {
WorkKind::Command => "this command",
WorkKind::Agent => "this agent",
};
let mut finish = finish_now();
finish.error = Some(format!(
"OpenCrabs restarted while {subject} was running, so it was killed before finishing"
));
self.finish = Some(finish);
self.write()
}
pub fn read(id: &str) -> Option<Self> {
let path = status_path(id);
if !path.exists() {
return None;
}
let data = fs::read_to_string(&path).ok()?;
serde_json::from_str(&data).ok()
}
fn write(&self) -> std::io::Result<()> {
let path = status_path(&self.id);
ensure_dir()?;
let tmp = path.with_extension("json.tmp");
let data = serde_json::to_string_pretty(self).map_err(std::io::Error::other)?;
let mut f = fs::File::create(&tmp)?;
f.write_all(data.as_bytes())?;
f.sync_all()?;
fs::rename(tmp, path)
}
pub fn list_all() -> std::io::Result<Vec<String>> {
let dir = status_dir();
if !dir.exists() {
return Ok(Vec::new());
}
let mut ids = Vec::new();
for entry in fs::read_dir(&dir)? {
let entry = entry?;
if let Some(name) = entry.file_name().to_str()
&& let Some(id) = name.strip_suffix(".json")
{
ids.push(id.to_string());
}
}
ids.sort();
Ok(ids)
}
pub fn find_agent_by_session(session_id: &str) -> Option<WorkStatus> {
let ids = Self::list_all().ok()?;
ids.into_iter().filter_map(|id| Self::read(&id)).find(|s| {
matches!(s.kind, WorkKind::Agent)
&& s.session_id == session_id
&& !matches!(s.state, WorkState::Completed | WorkState::Failed)
})
}
}
#[derive(Deserialize)]
struct LegacySubagentStatus {
id: String,
label: String,
parent_session_id: String,
state: WorkState,
prompt: String,
started_at: String,
#[serde(default)]
progress: Option<ProgressSnapshot>,
#[serde(default)]
completed_at: Option<String>,
#[serde(default)]
error: Option<String>,
}
pub fn migrate_legacy_dir(legacy: &Path) -> usize {
let entries = match fs::read_dir(legacy) {
Ok(entries) => entries,
Err(_) => return 0,
};
let mut migrated = 0usize;
for entry in entries.flatten() {
let path = entry.path();
if path.extension().is_none_or(|e| e != "json") {
continue;
}
let data = match fs::read_to_string(&path) {
Ok(data) => data,
Err(e) => {
tracing::warn!(
target: "subagent",
"Could not read legacy sub-agent status file {}: {e}",
path.display()
);
continue;
}
};
let old: LegacySubagentStatus = match serde_json::from_str(&data) {
Ok(old) => old,
Err(e) => {
tracing::warn!(
target: "subagent",
"Could not parse legacy sub-agent status file {}: {e}",
path.display()
);
continue;
}
};
if status_path(&old.id).exists() {
let _ = fs::remove_file(&path);
continue;
}
let status = WorkStatus {
id: old.id.clone(),
kind: WorkKind::Agent,
session_id: old.parent_session_id,
parent_session_id: None,
label: old.label,
task: old.prompt,
spawned_at: old.started_at,
state: old.state,
progress: old.progress,
finish: old.completed_at.map(|completed_at| {
let mut finish = WorkFinish {
completed_at,
success: None,
code: None,
elapsed_secs: None,
output_bytes: None,
error: None,
output_full: None,
};
finish.error = old.error.clone();
finish
}),
};
match status.write() {
Ok(()) => {
let _ = fs::remove_file(&path);
migrated += 1;
}
Err(e) => {
tracing::warn!(
target: "subagent",
"Could not migrate legacy sub-agent status {} into the unified dir: {e}",
old.id
);
}
}
}
if migrated > 0 {
let _ = fs::remove_dir(legacy);
}
migrated
}
pub fn cleanup_stale(max_age: Duration) -> std::io::Result<(usize, usize)> {
let dir = status_dir();
if !dir.exists() {
return Ok((0, 0));
}
let cutoff = SystemTime::now()
.checked_sub(max_age)
.unwrap_or(SystemTime::UNIX_EPOCH);
let mut scanned = 0usize;
let mut removed = 0usize;
for entry in fs::read_dir(&dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().is_none_or(|e| e != "json") {
continue;
}
scanned += 1;
let should_delete = if let Ok(data) = fs::read_to_string(&path) {
if let Ok(status) = serde_json::from_str::<WorkStatus>(&data) {
status
.finish
.as_ref()
.is_some_and(|f| parse_completed_at(&cutoff, &f.completed_at))
|| status.finish.is_none() && file_stale(&path, &cutoff)
} else {
file_stale(&path, &cutoff)
}
} else {
file_stale(&path, &cutoff)
};
if should_delete {
fs::remove_file(&path)?;
removed += 1;
}
}
Ok((scanned, removed))
}
fn parse_completed_at(cutoff: &SystemTime, ts: &str) -> bool {
let Ok(dt) = chrono::DateTime::parse_from_rfc3339(ts) else {
return false; };
let completed = SystemTime::UNIX_EPOCH
.checked_add(Duration::from_secs(dt.timestamp() as u64))
.unwrap_or(SystemTime::UNIX_EPOCH);
completed < *cutoff
}
fn file_stale(path: &Path, cutoff: &SystemTime) -> bool {
path.metadata()
.and_then(|m| m.modified())
.map(|mtime| mtime < *cutoff)
.unwrap_or(true) }
fn now_rfc3339() -> String {
chrono::Utc::now().to_rfc3339()
}