diff --git a/fact/src/bpf/mod.rs b/fact/src/bpf/mod.rs index e141dbdd..10362574 100644 --- a/fact/src/bpf/mod.rs +++ b/fact/src/bpf/mod.rs @@ -293,6 +293,7 @@ impl Bpf { task_set.spawn(async move { let rb = self.take_ringbuffer()?; let mut fd = AsyncFd::new(rb)?; + let mut config_is_closed = false; loop { tokio::select! { @@ -330,8 +331,19 @@ impl Bpf { } guard.clear_ready(); }, - _ = self.paths_config.changed() => { - self.load_paths().context("Failed to load paths")?; + // The precondition here could directly use `has_changed().is_err()`, + // however, since this is a tight loop processing events from the + // kernel, using a local variable should be more performant since + // `has_changed` reads an atomic variable. + // + // This approach also handles the case in which a value is sent + // on the channel and then closed, which should not happen in our + // code at the point this comment was written. + res = self.paths_config.changed(), if !config_is_closed => { + match res { + Ok(()) => self.load_paths().context("Failed to load paths")?, + Err(_) => config_is_closed = true, + } }, _ = self.running.changed() => { if !*self.running.borrow() { diff --git a/fact/src/endpoints.rs b/fact/src/endpoints.rs index 34d95334..e1d8b6e3 100644 --- a/fact/src/endpoints.rs +++ b/fact/src/endpoints.rs @@ -74,7 +74,7 @@ impl Server { /// Wait for configuration changes or fact to stop. async fn idle(&mut self) -> anyhow::Result { tokio::select! { - _ = self.config.changed() => Ok(true), + _ = self.config.changed(), if self.config.has_changed().is_ok() => Ok(true), _ = self.running.changed() => Ok(*self.running.borrow()), } } @@ -98,7 +98,11 @@ impl Server { } }); }, - _ = self.config.changed() => return Ok(true), + res = self.config.changed(), if self.config.has_changed().is_ok() => { + if res.is_ok() { + return Ok(true); + } + } _ = self.running.changed() => return Ok(*self.running.borrow()), } } diff --git a/fact/src/host_scanner.rs b/fact/src/host_scanner.rs index c8c72eb2..86658de2 100644 --- a/fact/src/host_scanner.rs +++ b/fact/src/host_scanner.rs @@ -553,7 +553,7 @@ You can increase this limit with: tokio::select! { _ = interval.tick() => scan_trigger.notify_one(), _ = running.changed() => break, - _ = scan_interval.changed() => break, + _ = scan_interval.changed(), if scan_interval.has_changed().is_ok() => break, } } } @@ -587,6 +587,7 @@ You can increase this limit with: task_set.spawn(async move { info!("Starting host scanner..."); + let mut config_is_closed = false; loop { tokio::select! { @@ -684,9 +685,12 @@ You can increase this limit with: } } _ = scan_trigger.notified() => self.scan()?, - _ = self.paths.changed() => { - self.scan()?; + res = self.paths.changed(), if !config_is_closed => { + match res { + Ok(()) => self.scan()?, + Err(_) => config_is_closed = true, } + } } } diff --git a/fact/src/output/grpc.rs b/fact/src/output/grpc.rs index e91eb897..08ad5fe6 100644 --- a/fact/src/output/grpc.rs +++ b/fact/src/output/grpc.rs @@ -262,7 +262,7 @@ impl Client { Err(e) => warn!("gRPC stream error: {e:?}"), } } - _ = self.config.changed() => return Ok(true), + _ = self.config.changed(), if self.config.has_changed().is_ok() => return Ok(true), _ = self.running.changed() => return Ok(*self.running.borrow()), } } @@ -274,7 +274,7 @@ impl Client { async fn idle(&mut self) -> anyhow::Result { tokio::select! { - _ = self.config.changed() => Ok(true), + _ = self.config.changed(), if self.config.has_changed().is_ok() => Ok(true), _ = self.running.changed() => Ok(*self.running.borrow()), } } diff --git a/fact/src/output/otel.rs b/fact/src/output/otel.rs index 3af202df..e02dea12 100644 --- a/fact/src/output/otel.rs +++ b/fact/src/output/otel.rs @@ -77,6 +77,7 @@ impl Client { let (tx, rx) = oneshot::channel(); self.subscriber.send(tx).await?; let mut rx = rx.await?; + let mut config_is_closed = false; let res = loop { tokio::select! { @@ -109,7 +110,12 @@ impl Client { } } } - _ = self.config.changed() => break Ok(true), + res = self.config.changed(), if !config_is_closed => { + match res { + Ok(()) => break Ok(true), + Err(_) => config_is_closed = true, + } + } _ = self.running.changed() => break Ok(*self.running.borrow()), } }; @@ -124,7 +130,7 @@ impl Client { async fn idle(&mut self) -> anyhow::Result { tokio::select! { - _ = self.config.changed() => Ok(true), + _ = self.config.changed(), if self.config.has_changed().is_ok() => Ok(true), _ = self.running.changed() => Ok(*self.running.borrow()), } } diff --git a/fact/src/rate_limiter.rs b/fact/src/rate_limiter.rs index e4bd9978..f5cbffac 100644 --- a/fact/src/rate_limiter.rs +++ b/fact/src/rate_limiter.rs @@ -65,6 +65,7 @@ impl RateLimiter { pub fn start(mut self, task_set: &mut JoinSet>) { task_set.spawn(async move { debug!("Starting rate limiter..."); + let mut config_is_closed = false; loop { tokio::select! { event = self.rx.recv() => { @@ -81,8 +82,11 @@ impl RateLimiter { self.metrics.errored(); } }, - _ = self.rate_config.changed() => { - self.reload_limiter()?; + res = self.rate_config.changed(), if !config_is_closed => { + match res { + Ok(()) => self.reload_limiter()?, + Err(_) => config_is_closed = true, + } }, } }