diff --git a/Cargo.lock b/Cargo.lock index a25c9ac..f3401a2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -66,6 +66,8 @@ dependencies = [ "keyring-core", "libc", "memchr", + "notify", + "notify-debouncer-full", "open", "plist", "rand", diff --git a/crates/agent/Cargo.toml b/crates/agent/Cargo.toml index b3c1ebe..8b256c0 100644 --- a/crates/agent/Cargo.toml +++ b/crates/agent/Cargo.toml @@ -37,6 +37,8 @@ http-body-util.workspace = true json5.workspace = true jsonc-parser.workspace = true memchr.workspace = true +notify.workspace = true +notify-debouncer-full.workspace = true open.workspace = true rand.workspace = true rcgen.workspace = true diff --git a/crates/agent/src/daemon.rs b/crates/agent/src/daemon.rs index bd01c19..cc04752 100644 --- a/crates/agent/src/daemon.rs +++ b/crates/agent/src/daemon.rs @@ -35,7 +35,8 @@ use windows_sys::Win32::Security::SECURITY_ATTRIBUTES; #[cfg(windows)] use crate::windows_security::SecurityDescriptor; use crate::{ - api, enrollment::EnrollmentState, gateway_oidc, llm_proxy, reconcile, remote, secure_fs, + api, enrollment::EnrollmentState, gateway_oidc, inventory_watch::InventoryWatch, llm_proxy, + reconcile, remote, secure_fs, }; #[cfg(unix)] @@ -769,12 +770,27 @@ async fn refresh_inventory( interval: Duration, reconciler: reconcile::Reconciler, ) { - refresh_inventory_with(sender, interval, || reconciler.discover()).await; + // The watcher owns its change channel so the capacity-one coalescing it + // relies on cannot be undone here. When it cannot start, a dropped sender + // leaves a receiver that never yields, and the interval carries on alone. + let (watch, changes) = match InventoryWatch::new() { + Ok((watch, changes)) => (Some(watch), changes), + Err(error) => { + tracing::warn!( + error = %format!("{error:#}"), + "inventory file watching unavailable; refreshing on the interval only" + ); + (None, mpsc::channel(1).1) + } + }; + refresh_inventory_with(sender, interval, changes, watch, || reconciler.discover()).await; } async fn refresh_inventory_with( sender: watch::Sender>, interval: Duration, + mut changes: mpsc::Receiver<()>, + mut watch: Option, mut discover: F, ) where F: FnMut() -> Fut, @@ -786,10 +802,18 @@ async fn refresh_inventory_with( tracing::error!("inventory interval must be greater than zero; not refreshing inventory"); return; } + if let Some(watch) = watch.as_mut() { + watch.sync(&sender.borrow().clone()); + } let mut ticker = time::interval_at(time::Instant::now() + interval, interval); ticker.set_missed_tick_behavior(time::MissedTickBehavior::Delay); loop { - ticker.tick().await; + // A closed change channel disables its branch rather than ending the + // loop, so a watcher that could not start leaves the interval running. + let trigger = tokio::select! { + _ = ticker.tick() => "interval", + Some(()) = changes.recv() => "change", + }; if sender.is_closed() { return; } @@ -797,8 +821,11 @@ async fn refresh_inventory_with( if **sender.borrow() == discovery { continue; } - tracing::info!("inventory changed"); + tracing::info!(trigger, "inventory changed"); log_discovery(&discovery); + if let Some(watch) = watch.as_mut() { + watch.sync(&discovery); + } if sender.send(Arc::new(discovery)).is_err() { return; } @@ -1233,7 +1260,7 @@ mod tests { use agentdesktop_core::DEFAULT_SOCKET_PATH; use agentdesktop_core::config::{self, parse_daemon}; use agentdesktop_core::model::{Agent, Discovery, LlmProxyInfo}; - use tokio::sync::watch; + use tokio::sync::{mpsc, watch}; fn discovery(kinds: &[&str]) -> Discovery { Discovery { @@ -1254,7 +1281,8 @@ mod tests { #[tokio::test(start_paused = true)] async fn refuses_to_schedule_a_zero_interval() { let (sender, receiver) = watch::channel(Arc::new(discovery(&["codex"]))); - super::refresh_inventory_with(sender, Duration::ZERO, || async { + let (_changes_sender, changes) = mpsc::channel(1); + super::refresh_inventory_with(sender, Duration::ZERO, changes, None, || async { panic!("a zero interval must not schedule a scan") }) .await; @@ -1265,20 +1293,27 @@ mod tests { async fn publishes_inventory_only_when_it_changes() { let (sender, mut receiver) = watch::channel(Arc::new(discovery(&["codex"]))); let scans = Arc::new(AtomicUsize::new(0)); + let (_changes_sender, changes) = mpsc::channel(1); let refresher = tokio::spawn({ let scans = Arc::clone(&scans); - super::refresh_inventory_with(sender, Duration::from_secs(60), move || { - let scan = scans.fetch_add(1, Ordering::SeqCst); - // The first scan repeats the boot snapshot; every later scan - // reports a newly installed tool. - async move { - if scan == 0 { - discovery(&["codex"]) - } else { - discovery(&["codex", "cursor"]) + super::refresh_inventory_with( + sender, + Duration::from_secs(60), + changes, + None, + move || { + let scan = scans.fetch_add(1, Ordering::SeqCst); + // The first scan repeats the boot snapshot; every later scan + // reports a newly installed tool. + async move { + if scan == 0 { + discovery(&["codex"]) + } else { + discovery(&["codex", "cursor"]) + } } - } - }) + }, + ) }); // Paused time auto-advances between ticks, so this resolves on the @@ -1299,12 +1334,39 @@ mod tests { refresher.abort(); } + /// Real time, and an interval far longer than the test will wait, so only + /// the change signal can account for the refresh. + #[tokio::test] + async fn refreshes_on_a_change_signal_without_waiting_for_the_interval() { + let (sender, mut receiver) = watch::channel(Arc::new(discovery(&["codex"]))); + let (changes_sender, changes) = mpsc::channel(1); + let refresher = tokio::spawn(super::refresh_inventory_with( + sender, + Duration::from_secs(3600), + changes, + None, + || async { discovery(&["codex", "cursor"]) }, + )); + + changes_sender.try_send(()).expect("refresher is listening"); + tokio::time::timeout(Duration::from_secs(5), receiver.changed()) + .await + .expect("a change signal must refresh before the interval elapses") + .expect("refresher is alive"); + assert_eq!(receiver.borrow_and_update().agents.len(), 2); + + refresher.abort(); + } + #[tokio::test(start_paused = true)] async fn stops_refreshing_once_every_reader_is_gone() { let (sender, receiver) = watch::channel(Arc::new(discovery(&["codex"]))); + let (_changes_sender, changes) = mpsc::channel(1); let refresher = tokio::spawn(super::refresh_inventory_with( sender, Duration::from_secs(60), + changes, + None, || async { discovery(&["codex", "cursor"]) }, )); drop(receiver); diff --git a/crates/agent/src/inventory_watch.rs b/crates/agent/src/inventory_watch.rs new file mode 100644 index 0000000..0ba8732 --- /dev/null +++ b/crates/agent/src/inventory_watch.rs @@ -0,0 +1,363 @@ +//! Filesystem watching for the files an inventory was built from. +//! +//! Interval refreshes bound how stale the inventory can get; this bounds how +//! long a change takes to show up. A developer adding an MCP server should not +//! wait out the interval before the fleet view reflects it. + +use std::{ + collections::BTreeMap, + path::{Path, PathBuf}, + time::Duration, +}; + +use agentdesktop_core::model::Discovery; +use notify::RecursiveMode; +use notify_debouncer_full::{DebounceEventResult, Debouncer, RecommendedCache}; +use tokio::sync::mpsc; +use tracing::{debug, warn}; + +/// Matches the window the controller uses for its own configuration watch. +/// +/// Editors frequently write a settings file as a burst of rename and write +/// events, and a scan walks user home directories, so coalescing is worth more +/// than reacting to the first event. +const DEBOUNCE: Duration = Duration::from_millis(250); + +/// Watches the directories holding the files the current inventory came from. +pub(crate) struct InventoryWatch { + debouncer: Debouncer, + watched: BTreeMap, +} + +impl InventoryWatch { + /// Builds a watcher and the channel it signals on. + /// + /// The channel deliberately has capacity one, and the watcher owns that + /// choice rather than its caller. A scan walks user home directories and + /// answers every change that arrived before it started, so activity while a + /// signal is already pending is dropped rather than queued: otherwise a + /// broad watched directory could enqueue batches faster than they are + /// consumed and run one full scan per stale token. + /// + /// Events are not inspected beyond debouncing. A scan is the only way to + /// know whether the inventory actually changed, and the caller already + /// discards a scan that matches the published snapshot. + pub(crate) fn new() -> anyhow::Result<(Self, mpsc::Receiver<()>)> { + let (changes, receiver) = mpsc::channel(1); + let debouncer = notify_debouncer_full::new_debouncer( + DEBOUNCE, + None, + move |result: DebounceEventResult| match result { + Ok(_) => { + let _ = changes.try_send(()); + } + Err(errors) => warn!(?errors, "inventory watch error"), + }, + )?; + Ok(( + Self { + debouncer, + watched: BTreeMap::new(), + }, + receiver, + )) + } + + /// Points the watch at the directories behind `discovery`. + /// + /// Called again after every published snapshot, so a tool that gains or + /// loses configuration files is followed without restarting the daemon. + pub(crate) fn sync(&mut self, discovery: &Discovery) { + let wanted = source_directories(discovery); + let stale: Vec<_> = self + .watched + .iter() + .filter(|(path, mode)| wanted.get(*path) != Some(mode)) + .map(|(path, _)| path.clone()) + .collect(); + let added: Vec<_> = wanted + .iter() + .filter(|(path, mode)| self.watched.get(*path) != Some(mode)) + .map(|(path, mode)| (path.clone(), *mode)) + .collect(); + + for path in stale { + let _ = self.debouncer.unwatch(&path); + self.watched.remove(&path); + } + for (path, mode) in added { + match self.debouncer.watch(&path, mode) { + Ok(()) => { + debug!(path = %path.display(), ?mode, "watching inventory source directory"); + self.watched.insert(path, mode); + } + Err(error) => { + // A directory can disappear between the scan and this call, + // and a daemon does not necessarily have access to every + // user's home. Neither is fatal; the interval still runs. + debug!(path = %path.display(), %error, "cannot watch inventory source directory"); + } + } + } + } + + #[cfg(test)] + pub(crate) fn watched(&self) -> &BTreeMap { + &self.watched + } +} + +/// Directories whose contents produced the inventory, and how deeply to watch. +/// +/// Parent directories rather than the files themselves, so that an editor +/// replacing a file atomically, and a sibling file appearing next to one +/// already discovered, both register. Executables are deliberately not +/// watched: an upgrade changes a version rather than a configuration, and +/// install directories are noisy enough to make every scan a wasted one. +fn source_directories(discovery: &Discovery) -> BTreeMap { + let mut directories = BTreeMap::new(); + for agent in &discovery.agents { + for server in &agent.mcp_servers { + if let Some(parent) = existing_directory(server.source.parent()) { + insert(&mut directories, parent, RecursiveMode::NonRecursive); + } + } + for skill in &agent.skills { + // A skill walk recurses, so the discovered file can sit well below + // the root it was found from; watching only its immediate parent + // would miss a skill added anywhere else in that tree. + match skill_root(&skill.path) { + Some(root) => insert(&mut directories, root, RecursiveMode::Recursive), + None => { + if let Some(parent) = existing_directory(skill.path.parent()) { + insert(&mut directories, parent, RecursiveMode::NonRecursive); + } + } + } + } + } + directories +} + +/// Root a recursive skill walk started from, when it can be identified. +/// +/// Every root the providers scan ends in a `skills` component — `.claude/skills`, +/// `.codex/skills`, `.agents/skills`, `.cursor/skills`, `.github/skills` — so the +/// nearest such ancestor is the tree a new sibling skill would appear in. When a +/// path does not match that shape the caller falls back to the file's own +/// directory, which is what this module did before and is never worse. +fn skill_root(file: &Path) -> Option { + let mut directory = file.parent(); + while let Some(candidate) = directory { + if candidate.file_name().is_some_and(|name| name == "skills") && candidate.is_dir() { + return Some(candidate.to_path_buf()); + } + directory = candidate.parent(); + } + None +} + +fn existing_directory(directory: Option<&Path>) -> Option { + directory + .filter(|directory| directory.is_dir()) + .map(Path::to_path_buf) +} + +/// Keeps the deeper watch when a directory is wanted both ways. +fn insert(directories: &mut BTreeMap, path: PathBuf, mode: RecursiveMode) { + let entry = directories.entry(path).or_insert(mode); + if mode == RecursiveMode::Recursive { + *entry = RecursiveMode::Recursive; + } +} + +#[cfg(test)] +mod tests { + use std::{fs, time::Duration}; + + use agentdesktop_core::model::{Agent, McpServer, Skill}; + + use super::*; + + fn temp_root(name: &str) -> PathBuf { + let root = std::env::temp_dir().join(format!( + "agentdesktop-watch-{name}-{}-{:?}", + std::process::id(), + std::thread::current().id() + )); + let _ = fs::remove_dir_all(&root); + fs::create_dir_all(&root).unwrap(); + root + } + + fn inventory(servers: Vec, skills: Vec) -> Discovery { + Discovery { + agents: vec![Agent { + kind: "cursor".to_owned(), + executable: PathBuf::from("/usr/local/bin/cursor"), + version: None, + mcp_servers: servers, + skills, + }], + model_runtimes: Vec::new(), + } + } + + fn server(source: &Path) -> McpServer { + McpServer { + name: "docs".to_owned(), + transport: "http".to_owned(), + command: None, + url: Some("https://example.test/mcp".to_owned()), + enabled: true, + source: source.to_path_buf(), + } + } + + fn skill(path: &Path) -> Skill { + Skill { + path: path.to_path_buf(), + front_matter: Default::default(), + } + } + + #[test] + fn watches_the_directory_holding_a_discovered_mcp_configuration() { + let root = temp_root("mcp"); + let config = root.join("mcp.json"); + fs::write(&config, "{}").unwrap(); + + let directories = source_directories(&inventory(vec![server(&config)], Vec::new())); + + assert_eq!(directories.get(&root), Some(&RecursiveMode::NonRecursive)); + // The executable's directory is deliberately not watched. + assert!(!directories.contains_key(&PathBuf::from("/usr/local/bin"))); + let _ = fs::remove_dir_all(&root); + } + + /// A skill walk recurses, so a discovered `SKILL.md` can sit several levels + /// below the root it came from. Watching only its own directory would miss + /// a skill added anywhere else in that tree, which is the ordinary case. + #[test] + fn watches_the_whole_skill_tree_rather_than_one_skill_directory() { + let root = temp_root("skills"); + let skills = root.join(".codex/skills"); + let nested = skills.join(".system/review-agent"); + fs::create_dir_all(&nested).unwrap(); + let file = nested.join("SKILL.md"); + fs::write(&file, "---\nname: review-agent\n---\n").unwrap(); + + let directories = source_directories(&inventory(Vec::new(), vec![skill(&file)])); + + assert_eq!(directories.get(&skills), Some(&RecursiveMode::Recursive)); + assert!( + !directories.contains_key(&nested), + "the tree is watched recursively, so the leaf needs no watch of its own" + ); + let _ = fs::remove_dir_all(&root); + } + + #[test] + fn falls_back_to_the_skill_directory_when_there_is_no_skills_root() { + let root = temp_root("rootless"); + let file = root.join("SKILL.md"); + fs::write(&file, "---\nname: loose\n---\n").unwrap(); + + let directories = source_directories(&inventory(Vec::new(), vec![skill(&file)])); + + assert_eq!(directories.get(&root), Some(&RecursiveMode::NonRecursive)); + let _ = fs::remove_dir_all(&root); + } + + #[test] + fn ignores_sources_whose_directory_is_gone() { + let directories = source_directories(&inventory( + vec![server(Path::new("/nonexistent-agentdesktop/mcp.json"))], + Vec::new(), + )); + assert!(directories.is_empty()); + } + + #[test] + fn follows_the_inventory_when_sources_appear_and_disappear() { + let root = temp_root("resync"); + let first = root.join("first"); + let second = root.join("second"); + fs::create_dir_all(&first).unwrap(); + fs::create_dir_all(&second).unwrap(); + fs::write(first.join("mcp.json"), "{}").unwrap(); + fs::write(second.join("mcp.json"), "{}").unwrap(); + + let (mut watch, _changes) = InventoryWatch::new().expect("create watcher"); + + watch.sync(&inventory( + vec![server(&first.join("mcp.json"))], + Vec::new(), + )); + assert_eq!(watch.watched().keys().collect::>(), vec![&first]); + + watch.sync(&inventory( + vec![server(&second.join("mcp.json"))], + Vec::new(), + )); + assert_eq!( + watch.watched().keys().collect::>(), + vec![&second], + "a source that left the inventory must stop being watched" + ); + + watch.sync(&inventory(Vec::new(), Vec::new())); + assert!(watch.watched().is_empty()); + let _ = fs::remove_dir_all(&root); + } + + /// The timeout is generous because macOS delivers through FSEvents, which + /// batches with its own latency on top of the debounce window. + #[tokio::test] + async fn reports_a_write_to_a_watched_directory() { + let root = temp_root("events"); + let config = root.join("mcp.json"); + fs::write(&config, "{}").unwrap(); + + let (mut watch, mut changes) = InventoryWatch::new().expect("create watcher"); + watch.sync(&inventory(vec![server(&config)], Vec::new())); + + fs::write(&config, r#"{"mcpServers":{}}"#).unwrap(); + + let change = tokio::time::timeout(Duration::from_secs(20), changes.recv()).await; + let _ = fs::remove_dir_all(&root); + assert_eq!( + change.expect("a write to a watched directory must be reported"), + Some(()) + ); + } + + /// Bursts must not queue a scan each. The watcher owns its channel so this + /// cannot be undone by a caller passing a deeper one. + #[tokio::test] + async fn coalesces_activity_into_a_single_pending_scan() { + let root = temp_root("burst"); + let config = root.join("mcp.json"); + fs::write(&config, "{}").unwrap(); + + let (mut watch, mut changes) = InventoryWatch::new().expect("create watcher"); + watch.sync(&inventory(vec![server(&config)], Vec::new())); + + for index in 0..25 { + fs::write(root.join(format!("noise-{index}.json")), "{}").unwrap(); + } + + let first = tokio::time::timeout(Duration::from_secs(20), changes.recv()).await; + assert_eq!(first.expect("the burst must be reported"), Some(())); + // Whatever the burst produced, at most one signal was ever held. + let mut pending = 0; + while changes.try_recv().is_ok() { + pending += 1; + } + let _ = fs::remove_dir_all(&root); + assert!( + pending <= 1, + "a burst queued {pending} extra scans; capacity-one coalescing is not working" + ); + } +} diff --git a/crates/agent/src/lib.rs b/crates/agent/src/lib.rs index b961d56..1c5ef6f 100644 --- a/crates/agent/src/lib.rs +++ b/crates/agent/src/lib.rs @@ -6,6 +6,7 @@ pub mod enrollment; pub mod gateway_oidc; mod github_oauth; pub mod identity; +mod inventory_watch; mod llm_proxy; pub mod oidc; pub mod provider;