From 893e6f223feebfc7dcda528e9f106f6f91d9d591 Mon Sep 17 00:00:00 2001 From: guillaumezin Date: Sat, 12 Sep 2026 13:35:45 +0200 Subject: [PATCH 1/3] Add afk or not-akf status in window watcher to ease parsing with Grafana for instance --- watchers/src/report_client.rs | 34 ++++++++++++++++++++++++++--- watchers/src/watchers/x11_window.rs | 3 ++- 2 files changed, 33 insertions(+), 4 deletions(-) diff --git a/watchers/src/report_client.rs b/watchers/src/report_client.rs index 838882b..9098d31 100644 --- a/watchers/src/report_client.rs +++ b/watchers/src/report_client.rs @@ -7,12 +7,15 @@ use serde_json::{Map, Value}; use std::collections::HashMap; use std::error::Error; use std::future::Future; +use std::sync::atomic::{AtomicBool, AtomicI64, Ordering}; pub struct ReportClient { pub client: AwClient, pub config: Config, idle_bucket_name: String, active_window_bucket_name: String, + is_idle: AtomicBool, + last_idle_update: AtomicI64, } impl ReportClient { @@ -37,6 +40,8 @@ impl ReportClient { config, idle_bucket_name, active_window_bucket_name, + is_idle: AtomicBool::new(true), + last_idle_update: AtomicI64::new(0), }) } @@ -98,8 +103,22 @@ impl ReportClient { .with_context(|| "Failed to send heartbeat") } + pub fn afk_status(&self) -> String { + let staleness_limit = self.config.idle_timeout.num_seconds() * 2; + let last_update = self.last_idle_update.load(Ordering::Relaxed); + let age = Utc::now().timestamp() - last_update; + + if age > staleness_limit || self.is_idle.load(Ordering::Relaxed) { + "afk".to_string() + } else { + "not-afk".to_string() + } + } + pub async fn send_active_window(&self, app_id: &str, title: &str) -> anyhow::Result<()> { - self.send_active_window_with_extra(app_id, title, None) + let mut extra = HashMap::new(); + extra.insert("status".to_string(), self.afk_status()); + self.send_active_window_with_extra(app_id, title, Some(extra)) .await } @@ -191,16 +210,25 @@ impl ReportClient { } pub async fn handle_idle_status(&self, status: Status) -> anyhow::Result<()> { + self.last_idle_update + .store(Utc::now().timestamp(), Ordering::Relaxed); + match status { Status::Idle { changed, last_input_time, duration, - } => self.idle(changed, last_input_time, duration).await, + } => { + self.is_idle.store(true, Ordering::Relaxed); + self.idle(changed, last_input_time, duration).await + } Status::Active { changed, last_input_time, - } => self.non_idle(changed, last_input_time).await, + } => { + self.is_idle.store(false, Ordering::Relaxed); + self.non_idle(changed, last_input_time).await + } } } diff --git a/watchers/src/watchers/x11_window.rs b/watchers/src/watchers/x11_window.rs index 1534a6e..1f00117 100644 --- a/watchers/src/watchers/x11_window.rs +++ b/watchers/src/watchers/x11_window.rs @@ -22,10 +22,11 @@ impl WindowWatcher { ) -> anyhow::Result<()> { let mut extra_data = HashMap::new(); extra_data.insert("wm_instance".to_string(), wm_instance.to_string()); + extra_data.insert("status".to_string(), client.afk_status()); // NOUVEAU client .send_active_window_with_extra(app_id, title, Some(extra_data)) .await - } + } async fn send_active_window(&mut self, client: &ReportClient) -> anyhow::Result<()> { let data = self.client.active_window_data()?; From eb56c8eb6a26c3951e2a7752c5503c3bf171b90a Mon Sep 17 00:00:00 2001 From: guillaumezin Date: Sat, 12 Sep 2026 14:50:36 +0200 Subject: [PATCH 2/3] Remove comment --- watchers/src/watchers/x11_window.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/watchers/src/watchers/x11_window.rs b/watchers/src/watchers/x11_window.rs index 1f00117..72aedcf 100644 --- a/watchers/src/watchers/x11_window.rs +++ b/watchers/src/watchers/x11_window.rs @@ -22,11 +22,11 @@ impl WindowWatcher { ) -> anyhow::Result<()> { let mut extra_data = HashMap::new(); extra_data.insert("wm_instance".to_string(), wm_instance.to_string()); - extra_data.insert("status".to_string(), client.afk_status()); // NOUVEAU + extra_data.insert("status".to_string(), client.afk_status()); client .send_active_window_with_extra(app_id, title, Some(extra_data)) .await - } + } async fn send_active_window(&mut self, client: &ReportClient) -> anyhow::Result<()> { let data = self.client.active_window_data()?; From c96298f074993a2d80edb7a66363edf3d3a318fd Mon Sep 17 00:00:00 2001 From: guillaumezin Date: Fri, 25 Sep 2026 22:07:48 +0200 Subject: [PATCH 3/3] fix(watchers): give up and let the process restart after repeated run_iteration failures MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit run_first_supported's loop caught every error from run_iteration (including timeouts) and just logged it, then kept looping forever on the same broken state. In particular, once a Wayland idle watcher's event queue socket dies (e.g. "Broken pipe" after the KWin compositor restarts, crashes, or a suspend/resume cycle invalidates the connection), there is no reconnection logic anywhere — every subsequent iteration fails identically, indefinitely, without the process ever exiting or self-healing. Since main.rs races the idle and active-window watcher tasks with tokio::select! and exits the whole process as soon as either task completes, we can use that as the recovery mechanism: after too many consecutive failures, give up on the current watcher instance and return, which ends the process. An external supervisor (systemd, launch script) is expected to restart it, establishing a fresh connection to the compositor/X server from scratch. Changes: - Track consecutive_failures across loop iterations, reset to 0 on success. - After MAX_CONSECUTIVE_FAILURES (10, i.e. ~50s at the default 5s idle poll interval) consecutive errors or timeouts, log a final message and return, ending run_first_supported's loop. Applies uniformly to both WatcherType::Idle and WatcherType::ActiveWindow since they share this loop. This complements (does not replace) the staleness-based afk fallback in report_client.rs: that fallback keeps status reporting safe while the watcher is stuck and during the brief restart window, while this change actually restores idle detection instead of leaving it permanently stuck after a broken connection. --- watchers/src/watchers.rs | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/watchers/src/watchers.rs b/watchers/src/watchers.rs index 1add035..8148b92 100644 --- a/watchers/src/watchers.rs +++ b/watchers/src/watchers.rs @@ -137,19 +137,34 @@ pub async fn run_first_supported(client: Arc, watcher_type: &Watch let supported_watcher = filter_first_supported(&client, watcher_type).await; if let Some(mut watcher) = supported_watcher { info!("Starting {watcher_type} watcher"); + + const MAX_CONSECUTIVE_FAILURES: u32 = 10; + let mut consecutive_failures: u32 = 0; + loop { let sleep_time = watcher_type.sleep_time(&client.config); match timeout(sleep_time, watcher.run_iteration(&client)).await { - Ok(Ok(())) => { /* Successfully completed. */ } + Ok(Ok(())) => { + consecutive_failures = 0; + } Ok(Err(e)) => { error!("Error on {watcher_type} iteration: {e}"); + consecutive_failures += 1; } Err(_) => { error!("Timeout on {watcher_type} iteration after {sleep_time:?}"); + consecutive_failures += 1; } } + if consecutive_failures >= MAX_CONSECUTIVE_FAILURES { + error!( + "{watcher_type} watcher failed {consecutive_failures} times in a row, giving up to let the process restart" + ); + return false; + } + sleep(sleep_time).await; } }