Skip to content
Open
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
8 changes: 6 additions & 2 deletions fact/src/bpf/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -302,6 +302,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! {
Expand Down Expand Up @@ -339,8 +340,11 @@ impl Bpf {
}
guard.clear_ready();
},
_ = self.paths_config.changed() => {
self.load_paths().context("Failed to load paths")?;
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() {
Expand Down
8 changes: 6 additions & 2 deletions fact/src/endpoints.rs
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ impl Server {
/// Wait for configuration changes or fact to stop.
async fn idle(&mut self) -> anyhow::Result<bool> {
tokio::select! {
_ = self.config.changed() => Ok(true),
_ = self.config.changed(), if self.config.has_changed().is_ok() => Ok(true),
Comment thread
Molter73 marked this conversation as resolved.
_ = self.running.changed() => Ok(*self.running.borrow()),
}
}
Expand All @@ -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()),
}
}
Expand Down
14 changes: 10 additions & 4 deletions fact/src/host_scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -578,7 +578,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,
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
}
}
Expand Down Expand Up @@ -611,6 +611,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! {
Expand Down Expand Up @@ -708,10 +709,15 @@ You can increase this limit with:
}
}
_ = scan_trigger.notified() => self.scan()?,
_ = self.paths.changed() => {
self.reload_paths_config()?;
self.scan()?;
res = self.paths.changed(), if !config_is_closed => {
match res {
Ok(()) => {
self.reload_paths_config()?;
self.scan()?;
}
Err(_) => config_is_closed = true,
}
}
}
}

Expand Down
4 changes: 2 additions & 2 deletions fact/src/output/grpc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()),
}
}
Expand All @@ -274,7 +274,7 @@ impl Client {

async fn idle(&mut self) -> anyhow::Result<bool> {
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()),
}
}
Expand Down
10 changes: 8 additions & 2 deletions fact/src/output/otel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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! {
Expand Down Expand Up @@ -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()),
}
};
Expand All @@ -124,7 +130,7 @@ impl Client {

async fn idle(&mut self) -> anyhow::Result<bool> {
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()),
}
}
Expand Down
8 changes: 6 additions & 2 deletions fact/src/rate_limiter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ impl RateLimiter {
pub fn start(mut self, task_set: &mut JoinSet<anyhow::Result<()>>) {
task_set.spawn(async move {
debug!("Starting rate limiter...");
let mut config_is_closed = false;
loop {
tokio::select! {
event = self.rx.recv() => {
Expand All @@ -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,
}
},
}
}
Expand Down
Loading