use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use arc_swap::ArcSwap;
use super::session::{ClaudeSessionFile, read};
#[derive(Debug, Default)]
pub struct WatcherSnapshot {
pub candidate_pids: HashSet<i32>,
pub pid_metadata: HashMap<i32, ClaudeSessionFile>,
}
pub type Snapshot = Arc<ArcSwap<WatcherSnapshot>>;
#[must_use]
pub fn spawn(sessions_dir: PathBuf) -> Snapshot {
let initial = scan(&sessions_dir);
let snapshot: Snapshot = Arc::new(ArcSwap::from_pointee(initial));
let weak = Arc::downgrade(&snapshot);
tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(1));
interval.tick().await;
loop {
interval.tick().await;
let Some(slot) = weak.upgrade() else {
break;
};
let dir = sessions_dir.clone();
let next = match tokio::task::spawn_blocking(move || scan(&dir)).await {
Ok(next) => next,
Err(err) => {
tracing::warn!(
error = %err,
"session-watcher: scan task failed",
);
continue;
}
};
slot.store(Arc::new(next));
}
});
snapshot
}
fn scan(dir: &Path) -> WatcherSnapshot {
let mut snapshot = WatcherSnapshot::default();
let entries = match std::fs::read_dir(dir) {
Ok(e) => e,
Err(e) => {
tracing::debug!(
dir = %dir.display(),
error = %e,
"session-watcher: could not read sessions dir",
);
return snapshot;
}
};
for entry in entries.flatten() {
let name = entry.file_name();
let Some(stem) = name
.to_string_lossy()
.strip_suffix(".json")
.map(str::to_owned)
else {
continue;
};
let Ok(pid) = stem.parse::<i32>() else {
continue;
};
snapshot.candidate_pids.insert(pid);
if let Some(file) = read(dir, pid) {
snapshot.pid_metadata.insert(pid, file);
}
}
snapshot
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
#[test]
fn scan_parses_pid_filenames_and_metadata() {
let dir = tempfile::tempdir().unwrap();
let p = dir.path();
std::fs::write(
p.join("123.json"),
r#"{"pid":123,"sessionId":"abc","cwd":"/x"}"#,
)
.unwrap();
std::fs::write(
p.join("456.json"),
r#"{"pid":456,"sessionId":"def","cwd":"/y"}"#,
)
.unwrap();
std::fs::write(p.join("not-a-pid.json"), "{}").unwrap();
std::fs::write(p.join("789.txt"), "{}").unwrap();
let snap = scan(p);
assert!(snap.candidate_pids.contains(&123));
assert!(snap.candidate_pids.contains(&456));
assert_eq!(snap.candidate_pids.len(), 2);
assert_eq!(snap.pid_metadata.get(&123).unwrap().session_id, "abc");
assert_eq!(snap.pid_metadata.get(&456).unwrap().session_id, "def");
}
#[test]
fn scan_unparseable_metadata_drops_from_meta_keeps_in_pids() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("777.json"), "not json").unwrap();
let snap = scan(dir.path());
assert!(snap.candidate_pids.contains(&777));
assert!(!snap.pid_metadata.contains_key(&777));
}
#[test]
fn scan_missing_dir_returns_empty() {
let dir = tempfile::tempdir().unwrap();
let missing = dir.path().join("does-not-exist");
let snap = scan(&missing);
assert!(snap.candidate_pids.is_empty());
assert!(snap.pid_metadata.is_empty());
}
#[tokio::test]
async fn spawn_starts_with_initial_snapshot() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(
dir.path().join("42.json"),
r#"{"pid":42,"sessionId":"hello"}"#,
)
.unwrap();
let snap = spawn(dir.path().to_path_buf());
let loaded = snap.load();
assert!(loaded.candidate_pids.contains(&42));
assert_eq!(loaded.pid_metadata.get(&42).unwrap().session_id, "hello");
}
#[cfg(test)]
fn make_snapshot(pid: i32, tick: i64) -> WatcherSnapshot {
let mut next = WatcherSnapshot::default();
next.candidate_pids.insert(pid);
let raw = format!(r#"{{"pid":{pid},"sessionId":"sid-{tick}","cwd":"/"}}"#);
let meta: ClaudeSessionFile = serde_json::from_str(&raw).unwrap();
next.pid_metadata.insert(pid, meta);
next
}
#[cfg(test)]
async fn drive_writer(snap: Snapshot) {
for tick in 0..200i64 {
let pid = 1000 + (tick as i32 % 5);
snap.store(Arc::new(make_snapshot(pid, tick)));
tokio::task::yield_now().await;
}
}
#[cfg(test)]
async fn drive_reader(snap: Snapshot) {
for _ in 0..2_000 {
let loaded = snap.load();
assert_snapshot_consistent(&loaded);
tokio::task::yield_now().await;
}
}
#[cfg(test)]
fn assert_snapshot_consistent(loaded: &WatcherSnapshot) {
for pid in &loaded.candidate_pids {
assert!(
loaded.pid_metadata.contains_key(pid),
"candidate {pid} present without metadata — torn snapshot read",
);
}
}
#[cfg(unix)]
fn wait_until(limit: Duration, mut cond: impl FnMut() -> bool) -> bool {
let deadline = std::time::Instant::now() + limit;
loop {
if cond() {
return true;
}
if std::time::Instant::now() >= deadline {
return false;
}
std::thread::sleep(Duration::from_millis(10));
}
}
#[cfg(unix)]
fn mkfifo(path: &Path) {
use std::ffi::CString;
use std::os::unix::ffi::OsStrExt;
let c_path = CString::new(path.as_os_str().as_bytes()).unwrap();
let rc = unsafe { libc::mkfifo(c_path.as_ptr(), 0o600) };
assert_eq!(
rc,
0,
"mkfifo({}) failed: {}",
path.display(),
std::io::Error::last_os_error(),
);
}
#[cfg(unix)]
fn wait_for_fifo_reader(path: &Path) -> Option<std::fs::File> {
use std::os::unix::fs::OpenOptionsExt;
let mut opened = None;
wait_until(Duration::from_secs(15), || {
opened = std::fs::OpenOptions::new()
.write(true)
.custom_flags(libc::O_NONBLOCK)
.open(path)
.ok();
opened.is_some()
});
opened
}
#[cfg(unix)]
#[test]
fn periodic_scan_does_not_block_the_runtime_worker() {
use std::sync::atomic::{AtomicU64, Ordering};
let dir = tempfile::tempdir().unwrap();
let fifo = dir.path().join("4242.json");
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(1)
.enable_all()
.build()
.unwrap();
let snap = rt.block_on(async { spawn(dir.path().to_path_buf()) });
mkfifo(&fifo);
let beats = Arc::new(AtomicU64::new(0));
let counter = Arc::clone(&beats);
rt.spawn(async move {
loop {
counter.fetch_add(1, Ordering::Relaxed);
tokio::task::yield_now().await;
}
});
let writer = wait_for_fifo_reader(&fifo)
.expect("watcher never opened the FIFO — no periodic scan ran");
let before = beats.load(Ordering::Relaxed);
let progressed = wait_until(Duration::from_secs(5), || {
beats.load(Ordering::Relaxed) > before
});
assert!(
progressed,
"the runtime's only worker made no progress while a sessions-dir \
scan was in flight — the scan is running inline on the worker",
);
drop(snap);
std::fs::remove_file(&fifo).unwrap();
drop(writer);
}
#[tokio::test]
async fn periodic_scan_publishes_sessions_that_appear_after_startup() {
let dir = tempfile::tempdir().unwrap();
let snap = spawn(dir.path().to_path_buf());
assert!(snap.load().candidate_pids.is_empty());
std::fs::write(
dir.path().join("4321.json"),
r#"{"pid":4321,"sessionId":"appeared-later"}"#,
)
.unwrap();
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
while !snap.load().candidate_pids.contains(&4321) {
assert!(
tokio::time::Instant::now() < deadline,
"periodic scan never published the new session file",
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert_eq!(
snap.load().pid_metadata.get(&4321).unwrap().session_id,
"appeared-later",
);
}
#[tokio::test]
async fn snapshot_swap_is_atomic_across_candidate_pids_and_pid_metadata() {
let snap: Snapshot = Arc::new(ArcSwap::from_pointee(WatcherSnapshot::default()));
let writer = tokio::spawn(drive_writer(snap.clone()));
let reader = tokio::spawn(drive_reader(snap.clone()));
writer.await.unwrap();
reader.await.unwrap();
}
}