use anyhow::{Context, Result};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use std::fs::{self, OpenOptions};
use std::path::{Path, PathBuf};
use std::thread;
use std::time::{Duration, Instant};
use crate::config::Config;
const LOCK_RETRY: Duration = Duration::from_millis(10);
const LOCK_TIMEOUT: Duration = Duration::from_secs(10);
const LOCK_STALE_AFTER: Duration = Duration::from_secs(60);
struct ActivityFileLock {
path: PathBuf,
}
impl ActivityFileLock {
fn acquire(path: &Path) -> Result<Self> {
let started = Instant::now();
let lock_path = lock_path(path);
loop {
match OpenOptions::new()
.write(true)
.create_new(true)
.open(&lock_path)
{
Ok(_) => return Ok(Self { path: lock_path }),
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
if lock_is_stale(&lock_path) {
let _ = fs::remove_file(&lock_path);
continue;
}
if started.elapsed() >= LOCK_TIMEOUT {
anyhow::bail!(
"timed out waiting for activity lock {}",
lock_path.display()
);
}
thread::sleep(LOCK_RETRY);
}
Err(error) => {
return Err(error).with_context(|| format!("create {}", lock_path.display()));
}
}
}
}
}
impl Drop for ActivityFileLock {
fn drop(&mut self) {
let _ = fs::remove_file(&self.path);
}
}
fn lock_path(path: &Path) -> PathBuf {
path.with_file_name(format!(
"{}.lock",
path.file_name()
.and_then(|name| name.to_str())
.unwrap_or("activity.json")
))
}
fn lock_is_stale(path: &Path) -> bool {
path.metadata()
.and_then(|metadata| metadata.modified())
.ok()
.and_then(|modified| modified.elapsed().ok())
.is_some_and(|age| age > LOCK_STALE_AFTER)
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum ActivityKind {
Shell,
Subagent,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum ActivityState {
Running,
Done,
Failed,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Activity {
pub id: String,
pub kind: ActivityKind,
pub label: String,
pub detail: String,
pub state: ActivityState,
pub started_at: DateTime<Utc>,
pub finished_at: Option<DateTime<Utc>>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
struct ActivityFile {
items: Vec<Activity>,
}
pub struct ActivityLog {
path: PathBuf,
file: ActivityFile,
}
impl ActivityLog {
pub fn load() -> Result<Self> {
Config::ensure_home()?;
let path = Self::path()?;
let file = if path.exists() {
let raw =
fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
if raw.trim().is_empty() {
ActivityFile::default()
} else {
serde_json::from_str(&raw).unwrap_or_default()
}
} else {
ActivityFile::default()
};
Ok(Self { path, file })
}
pub fn path() -> Result<PathBuf> {
Config::ensure_home()?;
Ok(Config::home_dir().join("activity.json"))
}
pub fn list(&self) -> &[Activity] {
&self.file.items
}
pub fn visible(&self) -> Vec<Activity> {
let mut items = self.file.items.clone();
items.sort_by(|a, b| {
let ar = matches!(a.state, ActivityState::Running);
let br = matches!(b.state, ActivityState::Running);
br.cmp(&ar).then(b.started_at.cmp(&a.started_at))
});
items.truncate(40);
items
}
pub fn start(
kind: ActivityKind,
label: impl Into<String>,
detail: impl Into<String>,
) -> Result<String> {
let path = Self::path()?;
let _lock = ActivityFileLock::acquire(&path)?;
let mut log = Self::load()?;
let id = uuid::Uuid::new_v4().to_string();
log.file.items.push(Activity {
id: id.clone(),
kind,
label: clip(&label.into(), 72),
detail: clip(&detail.into(), 120),
state: ActivityState::Running,
started_at: Utc::now(),
finished_at: None,
});
log.prune();
log.save()?;
Ok(id)
}
pub fn finish(id: &str, ok: bool, detail: impl Into<String>) -> Result<()> {
let path = Self::path()?;
let _lock = ActivityFileLock::acquire(&path)?;
let mut log = Self::load()?;
if let Some(item) = log.file.items.iter_mut().find(|item| item.id == id) {
item.state = if ok {
ActivityState::Done
} else {
ActivityState::Failed
};
item.finished_at = Some(Utc::now());
let detail = detail.into();
if !detail.trim().is_empty() {
item.detail = clip(&detail, 120);
}
}
log.prune();
log.save()
}
fn prune(&mut self) {
let running: Vec<Activity> = self
.file
.items
.iter()
.filter(|item| matches!(item.state, ActivityState::Running))
.cloned()
.collect();
let mut done: Vec<Activity> = self
.file
.items
.iter()
.filter(|item| !matches!(item.state, ActivityState::Running))
.cloned()
.collect();
done.sort_by_key(|item| item.started_at);
let keep_done = done.len().saturating_sub(24);
done.drain(..keep_done);
self.file.items = running.into_iter().chain(done).collect();
}
fn save(&self) -> Result<()> {
if let Some(parent) = self.path.parent() {
fs::create_dir_all(parent)?;
}
let raw = serde_json::to_string_pretty(&self.file)?;
let temp = self.path.with_file_name(format!(
".{}.{}.tmp",
self.path
.file_name()
.and_then(|name| name.to_str())
.unwrap_or("activity.json"),
std::process::id()
));
fs::write(&temp, raw).with_context(|| format!("write {}", temp.display()))?;
fs::rename(&temp, &self.path).with_context(|| format!("write {}", self.path.display()))
}
}
fn clip(s: &str, n: usize) -> String {
let one = s.split_whitespace().collect::<Vec<_>>().join(" ");
one.chars().take(n).collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn visible_puts_running_before_finished() {
let log = ActivityLog {
path: PathBuf::from("activity.json"),
file: ActivityFile {
items: vec![
sample("a", ActivityState::Done, "2020-01-01T00:00:02Z"),
sample("b", ActivityState::Running, "2020-01-01T00:00:01Z"),
sample("c", ActivityState::Failed, "2020-01-01T00:00:03Z"),
],
},
};
let ids: Vec<_> = log.visible().into_iter().map(|item| item.id).collect();
assert_eq!(ids[0], "b");
}
#[test]
fn clip_collapses_whitespace() {
assert_eq!(clip(" cargo test\n--lib ", 12), "cargo test -");
}
fn sample(id: &str, state: ActivityState, at: &str) -> Activity {
Activity {
id: id.into(),
kind: ActivityKind::Shell,
label: id.into(),
detail: String::new(),
state,
started_at: at.parse().unwrap(),
finished_at: None,
}
}
#[test]
fn stale_lock_is_broken_and_replaced_by_a_fresh_one() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("activity.json");
let lock = lock_path(&path);
fs::write(&lock, b"").unwrap();
let stale_at = std::time::SystemTime::now()
.checked_sub(LOCK_STALE_AFTER + Duration::from_secs(5))
.unwrap();
let file = OpenOptions::new().write(true).open(&lock).unwrap();
file.set_modified(stale_at).unwrap();
drop(file);
let guard = ActivityFileLock::acquire(&path).unwrap();
assert_eq!(guard.path, lock);
assert!(!lock_is_stale(&lock));
drop(guard);
assert!(!lock.exists());
}
#[test]
fn save_round_trips_and_leaves_no_temp_files() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("activity.json");
let log = ActivityLog {
path: path.clone(),
file: ActivityFile {
items: vec![sample("a", ActivityState::Running, "2020-01-01T00:00:00Z")],
},
};
log.save().unwrap();
let file: ActivityFile = serde_json::from_str(&fs::read_to_string(&path).unwrap()).unwrap();
assert_eq!(file.items.len(), 1);
assert_eq!(file.items[0].id, "a");
let leftovers: Vec<_> = fs::read_dir(dir.path())
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.file_name())
.collect();
assert_eq!(leftovers.len(), 1, "temp files left behind: {leftovers:?}");
}
}