use crate::daemon::{now, Daemon, EXTERNAL_DETAIL_CAP};
use autofork_core::glob;
use autofork_core::hooks::HookOn;
use autofork_core::moments::ExternalKind;
use autofork_core::store::SessionRow;
use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
const MAX_WALK_DEPTH: usize = 24;
const SKIP_DIRS: &[&str] = &[".git", "node_modules", "target", ".venv", "__pycache__"];
struct Subscriber {
session: SessionRow,
raw: String,
forks: bool,
hooks: bool,
}
type Snapshot = HashMap<String, (u128, u64)>;
#[derive(Default)]
pub struct Watcher {
snapshots: HashMap<String, Snapshot>,
pending: HashMap<String, (Vec<String>, i64)>,
warned: HashSet<String>,
}
pub async fn watch_loop(daemon: Arc<Daemon>) {
let mut watcher = Watcher::default();
loop {
let cfg = daemon.cfg_for(None);
if cfg.watch_interval_secs == 0 {
tokio::select! {
_ = tokio::time::sleep(Duration::from_secs(30)) => continue,
_ = daemon.shutdown.notified() => return,
}
}
tokio::select! {
_ = tokio::time::sleep(Duration::from_secs(cfg.watch_interval_secs)) => {}
_ = daemon.shutdown.notified() => return,
}
sweep(&daemon, &mut watcher, &cfg);
}
}
pub fn sweep(daemon: &Arc<Daemon>, watcher: &mut Watcher, cfg: &autofork_core::config::Config) {
let registry = build_registry(daemon);
watcher.snapshots.retain(|k, _| registry.contains_key(k));
watcher.pending.retain(|k, _| registry.contains_key(k));
let t = now();
for pattern in registry.keys() {
let (snap, truncated) = scan(pattern, cfg.watch_max_files);
if truncated && watcher.warned.insert(pattern.clone()) {
tracing::warn!(
pattern = %pattern, cap = cfg.watch_max_files,
"watched pattern matches more files than watch_max_files; \
watching the first {} and ignoring the rest (narrow the glob, \
or raise watch_max_files)", cfg.watch_max_files
);
}
let Some(prev) = watcher.snapshots.get(pattern) else {
tracing::debug!(pattern = %pattern, files = snap.len(), "watch baseline");
watcher.snapshots.insert(pattern.clone(), snap);
continue;
};
let changed = diff(prev, &snap);
watcher.snapshots.insert(pattern.clone(), snap);
if changed.is_empty() {
continue;
}
let entry = watcher
.pending
.entry(pattern.clone())
.or_insert_with(|| (Vec::new(), t));
for c in changed {
if !entry.0.contains(&c) {
entry.0.push(c);
}
}
if t - entry.1 > 2 * cfg.watch_debounce_secs as i64 {
entry.1 = t - cfg.watch_debounce_secs as i64;
}
}
let settled: Vec<String> = watcher
.pending
.iter()
.filter(|(_, (_, started))| t - *started >= cfg.watch_debounce_secs as i64)
.map(|(k, _)| k.clone())
.collect();
for pattern in settled {
let Some((paths, _)) = watcher.pending.remove(&pattern) else {
continue;
};
let Some(subs) = registry.get(&pattern) else {
continue;
};
dispatch(daemon, subs, &paths, t);
}
}
fn dispatch(daemon: &Arc<Daemon>, subs: &[Subscriber], paths: &[String], t: i64) {
let detail = paths.join("\n");
for sub in subs {
tracing::info!(
session = %sub.session.session_id, pattern = %sub.raw, files = paths.len(),
"watched paths changed"
);
if sub.forks {
let store = daemon.store.lock().unwrap();
let _ = store.record_pending_trigger(
&sub.session.session_id,
ExternalKind::Changed.label(),
&sub.raw,
&detail,
EXTERNAL_DETAIL_CAP,
t,
);
}
if sub.hooks {
crate::hooks::fire_external(
daemon,
&crate::hooks::HookCtx::from_row(&sub.session),
ExternalKind::Changed,
&sub.raw,
&detail,
);
}
daemon.nudge(&sub.session.session_id);
}
}
fn build_registry(daemon: &Arc<Daemon>) -> HashMap<String, Vec<Subscriber>> {
let sessions = {
let store = daemon.store.lock().unwrap();
store.list_open_sessions().unwrap_or_default()
};
let home = std::env::var_os("HOME").map(PathBuf::from);
let mut out: HashMap<String, Vec<Subscriber>> = HashMap::new();
for session in sessions {
let base = &session.project_root;
let mut wanted: HashMap<String, (bool, bool)> = HashMap::new();
let (forks, _) = autofork_core::discovery::discover_forks(
&session.cwd,
Some(&daemon.user_forks_root()),
daemon.claude_dir().as_deref(),
daemon.agents_dir().as_deref(),
);
for entry in &forks {
for trigger in &entry.parsed.def.run_on {
if let autofork_core::frontmatter::ForkRunOn::Changed { pattern } = trigger {
wanted.entry(pattern.clone()).or_default().0 = true;
}
}
}
let (hooks, _) =
autofork_core::hooks::discover_hooks(&session.cwd, Some(&daemon.user_hooks_root()));
for entry in &hooks {
if entry.parsed.def.command.is_empty() {
continue;
}
for on in &entry.parsed.def.on {
if let HookOn::Changed { pattern } = on {
wanted.entry(pattern.clone()).or_default().1 = true;
}
}
}
for (raw, (forks, hooks)) in wanted {
let abs = glob::absolutize(&raw, base, home.as_deref());
out.entry(abs).or_default().push(Subscriber {
session: session.clone(),
raw,
forks,
hooks,
});
}
}
out
}
fn scan(pattern: &str, cap: usize) -> (Snapshot, bool) {
let mut snap = Snapshot::new();
let root = glob::literal_prefix(pattern);
if !glob::has_wildcard(pattern) {
if let Some(id) = file_id(&root) {
snap.insert(root.to_string_lossy().into_owned(), id);
}
return (snap, false);
}
let mut truncated = false;
walk(&root, 0, pattern, cap, &mut snap, &mut truncated);
(snap, truncated)
}
fn walk(
dir: &Path,
depth: usize,
pattern: &str,
cap: usize,
snap: &mut Snapshot,
truncated: &mut bool,
) {
if depth > MAX_WALK_DEPTH || snap.len() >= cap {
*truncated = *truncated || snap.len() >= cap;
return;
}
let Ok(read) = std::fs::read_dir(dir) else {
return;
};
for item in read.filter_map(|e| e.ok()) {
if snap.len() >= cap {
*truncated = true;
return;
}
let path = item.path();
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
let is_dir = item.file_type().map(|t| t.is_dir()).unwrap_or(false);
if is_dir {
if SKIP_DIRS.contains(&name) {
continue;
}
walk(&path, depth + 1, pattern, cap, snap, truncated);
continue;
}
let s = path.to_string_lossy();
if glob::matches(pattern, &s) {
if let Some(id) = file_id(&path) {
snap.insert(s.into_owned(), id);
}
}
}
}
fn file_id(path: &Path) -> Option<(u128, u64)> {
let md = std::fs::metadata(path).ok()?;
let mtime = md
.modified()
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map(|d| d.as_nanos())
.unwrap_or(0);
Some((mtime, md.len()))
}
fn diff(prev: &Snapshot, next: &Snapshot) -> Vec<String> {
let mut out: Vec<String> = Vec::new();
for (path, id) in next {
if prev.get(path) != Some(id) {
out.push(path.clone());
}
}
for path in prev.keys() {
if !next.contains_key(path) {
out.push(path.clone());
}
}
out.sort();
out
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs;
#[test]
fn diff_reports_adds_edits_and_deletes() {
let mut a = Snapshot::new();
a.insert("/x/1.md".into(), (1, 10));
a.insert("/x/2.md".into(), (1, 10));
let mut b = Snapshot::new();
b.insert("/x/2.md".into(), (2, 11)); b.insert("/x/3.md".into(), (1, 5)); let d = diff(&a, &b);
assert_eq!(d, vec!["/x/1.md", "/x/2.md", "/x/3.md"]);
}
#[test]
fn scan_matches_the_glob_and_skips_noise() {
let tmp = tempfile::tempdir().unwrap();
let root = tmp.path().canonicalize().unwrap();
fs::create_dir_all(root.join("h/2026/08")).unwrap();
fs::create_dir_all(root.join(".git/objects")).unwrap();
fs::write(root.join("h/2026/08/a.md"), "a").unwrap();
fs::write(root.join("h/2026/08/b.txt"), "b").unwrap();
fs::write(root.join(".git/objects/c.md"), "c").unwrap();
let pattern = format!("{}/h/**/*.md", root.display());
let (snap, truncated) = scan(&pattern, 1000);
assert!(!truncated);
let keys: Vec<&String> = snap.keys().collect();
assert_eq!(keys.len(), 1, "{keys:?}");
assert!(keys[0].ends_with("h/2026/08/a.md"));
}
#[test]
fn scan_honors_the_file_cap() {
let tmp = tempfile::tempdir().unwrap();
let root = tmp.path().canonicalize().unwrap();
for i in 0..10 {
fs::write(root.join(format!("f{i}.md")), "x").unwrap();
}
let pattern = format!("{}/*.md", root.display());
let (snap, truncated) = scan(&pattern, 4);
assert!(truncated);
assert!(snap.len() <= 4);
}
#[test]
fn literal_pattern_needs_no_walk() {
let tmp = tempfile::tempdir().unwrap();
let root = tmp.path().canonicalize().unwrap();
let f = root.join("one.md");
fs::write(&f, "x").unwrap();
let (snap, _) = scan(&f.to_string_lossy(), 1000);
assert_eq!(snap.len(), 1);
let (snap, _) = scan(&root.join("later.md").to_string_lossy(), 1000);
assert!(snap.is_empty());
}
}