From 3b9d50de03b4e8d7b11bcd7dba3251b9682ae8e0 Mon Sep 17 00:00:00 2001 From: Jaganmohan Reddy Sanikommu Date: Tue, 15 Sep 2026 17:38:42 -0400 Subject: [PATCH] feat(agent): refresh the inventory when its source files change The interval added in #75 bounds how stale the inventory can get. This bounds how long a change takes to appear: a developer adding an MCP server no longer waits out `inventoryInterval` before the fleet view reflects it. After every published snapshot the daemon watches the directories the files in that snapshot came from, and rescans when one reports activity. It reuses the debounced-watch approach the controller already uses for its own configuration file, including the same 250ms window. The watch set is derived from the inventory rather than a fixed path list, so no provider code is involved and a new provider is followed by virtue of reporting its sources. A discovered MCP configuration contributes its own directory, watched shallowly, so an atomic replace and a sibling configuration both register. A discovered skill contributes the root its walk started from, watched recursively, because that walk descends and a skill added elsewhere in the tree would otherwise go unnoticed until the interval. Skill paths that do not sit under a `skills` root fall back to the file's own directory. Executable directories are deliberately not watched. An upgrade changes a version rather than a configuration, and install directories are noisy enough to make most of those scans wasted. The watcher owns the channel it signals on, which has capacity one: a scan answers every change that arrived before it began, so activity while a signal is already pending is dropped rather than queued. Without that, a broad watched directory could enqueue batches faster than they are consumed and run one full scan per stale token. A watcher that cannot start is not fatal. Its sender drops, the receiver never yields, that select branch disables itself, and the interval carries on alone. Refs #62. Signed-off-by: Jaganmohan Reddy Sanikommu --- Cargo.lock | 2 + crates/agent/Cargo.toml | 2 + crates/agent/src/daemon.rs | 96 ++++++-- crates/agent/src/inventory_watch.rs | 363 ++++++++++++++++++++++++++++ crates/agent/src/lib.rs | 1 + 5 files changed, 447 insertions(+), 17 deletions(-) create mode 100644 crates/agent/src/inventory_watch.rs diff --git a/Cargo.lock b/Cargo.lock index a25c9ac3..f3401a22 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 b3c1ebe8..8b256c00 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 bd01c190..cc04752d 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 00000000..0ba87329 --- /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 b961d56a..1c5ef6f6 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;