Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions crates/agent/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
96 changes: 79 additions & 17 deletions crates/agent/src/daemon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down Expand Up @@ -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<F, Fut>(
sender: watch::Sender<Arc<agentdesktop_core::model::Discovery>>,
interval: Duration,
mut changes: mpsc::Receiver<()>,
mut watch: Option<InventoryWatch>,
mut discover: F,
) where
F: FnMut() -> Fut,
Expand All @@ -786,19 +802,30 @@ async fn refresh_inventory_with<F, Fut>(
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;
}
let discovery = discover().await;
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;
}
Expand Down Expand Up @@ -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 {
Expand All @@ -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;
Expand All @@ -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
Expand All @@ -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);
Expand Down
Loading