diff --git a/local/recipes/system/redbear-dbus-services/files/session-services/org.freedesktop.StatusNotifierWatcher.service b/local/recipes/system/redbear-dbus-services/files/session-services/org.freedesktop.StatusNotifierWatcher.service deleted file mode 100644 index 8a081759da..0000000000 --- a/local/recipes/system/redbear-dbus-services/files/session-services/org.freedesktop.StatusNotifierWatcher.service +++ /dev/null @@ -1,3 +0,0 @@ -[D-BUS Service] -Name=org.freedesktop.StatusNotifierWatcher -Exec=/usr/bin/redbear-statusnotifierwatcher diff --git a/local/recipes/system/redbear-dbus-services/files/session.d/org.redbear.session.conf b/local/recipes/system/redbear-dbus-services/files/session.d/org.redbear.session.conf index f22097dfe4..9db2dbcff9 100644 --- a/local/recipes/system/redbear-dbus-services/files/session.d/org.redbear.session.conf +++ b/local/recipes/system/redbear-dbus-services/files/session.d/org.redbear.session.conf @@ -9,8 +9,8 @@ - - - + + + diff --git a/local/recipes/system/redbear-sessiond/source/src/manager.rs b/local/recipes/system/redbear-sessiond/source/src/manager.rs index d46427d429..98a4a7b655 100644 --- a/local/recipes/system/redbear-sessiond/source/src/manager.rs +++ b/local/recipes/system/redbear-sessiond/source/src/manager.rs @@ -89,10 +89,10 @@ impl LoginManager { /// Remove all inhibitors whose `sender` matches `vanished_sender`, /// and drop the corresponding daemon-side pipe FDs so the pipe - /// breaks and the caller (if still alive) observes EOF. + /// breaks and the caller observes EOF. pub fn reap_inhibitors_for_sender(&self, vanished_sender: &str) { if let Ok(mut dead) = self.dead_senders.lock() { - dead.insert(vanaged_sender.to_owned()); + dead.insert(vanished_sender.to_owned()); } let mut removed = 0usize; if let Ok(mut runtime) = self.runtime.write() { @@ -111,17 +111,10 @@ impl LoginManager { } } - /// Background task: periodically asks the D-Bus daemon whether each - /// tracked sender still owns its unique name. When `name_has_owner` - /// returns `false`, the sender has disconnected and its inhibitors - /// are reaped. - /// - /// Uses `org.freedesktop.DBus.NameHasOwner` polling instead of a - /// `NameOwnerChanged` signal subscription because consuming zbus - /// signal streams requires a `StreamExt` that is not available - /// without adding a direct dependency on `futures-lite` / - /// `futures-util`. The 2-second poll interval keeps reaping - /// responsive while staying lightweight. + /// Background task: polls `NameHasOwner` every 2 s for each tracked + /// sender. Uses polling instead of a `NameOwnerChanged` signal + /// subscription because zbus signal streams require a `StreamExt` + /// that would add a direct `futures-lite` / `futures-util` dep. pub async fn run_inhibitor_reaper(&self) { let conn = match self.get_connection() { Some(c) => c, @@ -180,6 +173,54 @@ impl LoginManager { } } + pub fn inhibit_impl( + &self, + what: &str, + who: &str, + why: &str, + mode: &str, + sender: Option, + ) -> fdo::Result { + if mode != "block" && mode != "delay" { + return Err(fdo::Error::Failed(format!( + "inhibit mode must be 'block' or 'delay', got '{mode}'" + ))); + } + + let (end_caller, end_daemon) = UnixStream::pair() + .map_err(|err| fdo::Error::Failed(format!("failed to create inhibit pipe: {err}")))?; + + let fd_caller: StdOwnedFd = end_caller.into(); + let fd_daemon: StdOwnedFd = end_daemon.into(); + + let uid = self.runtime_read().map(|r| r.uid).unwrap_or(0); + let pid = std::process::id(); + + let entry = InhibitorEntry { + what: what.to_owned(), + who: who.to_owned(), + why: why.to_owned(), + mode: mode.to_owned(), + pid, + uid, + sender: sender.clone(), + }; + + if let Ok(mut runtime) = self.runtime.write() { + runtime.inhibitors.push(entry); + } + + if let Ok(mut fds) = self.inhibitor_fds.lock() { + fds.push((sender, fd_daemon)); + } + + eprintln!( + "redbear-sessiond: Inhibit(what={what}, who={who}, mode={mode}) granted" + ); + + Ok(OwnedFd::from(fd_caller)) + } + async fn emit_prepare_for_shutdown(&self, before: bool) { if let Some(conn) = self.get_connection() { if let Err(err) = conn @@ -328,44 +369,16 @@ impl LoginManager { Err(fdo::Error::Failed(format!("unknown login1 user uid {uid}"))) } - fn inhibit(&self, what: &str, who: &str, why: &str, mode: &str) -> fdo::Result { - if mode != "block" && mode != "delay" { - return Err(fdo::Error::Failed(format!( - "inhibit mode must be 'block' or 'delay', got '{mode}'" - ))); - } - - let (end_caller, end_daemon) = UnixStream::pair() - .map_err(|err| fdo::Error::Failed(format!("failed to create inhibit pipe: {err}")))?; - - let fd_caller: StdOwnedFd = end_caller.into(); - let fd_daemon: StdOwnedFd = end_daemon.into(); - - let uid = self.runtime_read().map(|r| r.uid).unwrap_or(0); - let pid = std::process::id(); - - let entry = InhibitorEntry { - what: what.to_owned(), - who: who.to_owned(), - why: why.to_owned(), - mode: mode.to_owned(), - pid, - uid, - }; - - if let Ok(mut runtime) = self.runtime.write() { - runtime.inhibitors.push(entry); - } - - if let Ok(mut fds) = self.inhibitor_fds.lock() { - fds.push(fd_daemon); - } - - eprintln!( - "redbear-sessiond: Inhibit(what={what}, who={who}, mode={mode}) granted" - ); - - Ok(OwnedFd::from(fd_caller)) + fn inhibit( + &self, + what: &str, + who: &str, + why: &str, + mode: &str, + #[zbus(header)] hdr: Header<'_>, + ) -> fdo::Result { + let sender = hdr.sender().map(|s| s.as_str().to_owned()); + self.inhibit_impl(what, who, why, mode, sender) } fn can_power_off(&self) -> fdo::Result { @@ -540,10 +553,21 @@ impl LoginManager { } fn list_inhibitors(&self) -> fdo::Result> { + let dead = self + .dead_senders + .lock() + .map(|g| g.clone()) + .unwrap_or_default(); let runtime = self.runtime_read()?; Ok(runtime .inhibitors .iter() + .filter(|entry| { + match &entry.sender { + Some(s) => !dead.contains(s), + None => true, + } + }) .map(|entry| { ( entry.what.clone(), @@ -925,7 +949,7 @@ mod tests { #[test] fn inhibit_rejects_invalid_mode() { let manager = test_manager(); - let err = manager.inhibit("sleep", "test", "reason", "invalid").unwrap_err(); + let err = manager.inhibit_impl("sleep", "test", "reason", "invalid", None).unwrap_err(); match err { fdo::Error::Failed(msg) => assert!(msg.contains("block") || msg.contains("delay")), other => panic!("expected Failed error, got {other:?}"), @@ -943,7 +967,7 @@ mod tests { ); let _fd = manager - .inhibit("sleep", "testapp", "testing", "block") + .inhibit_impl("sleep", "testapp", "testing", "block", None) .expect("inhibit should succeed"); let runtime_guard = runtime.read().expect("lock"); @@ -963,6 +987,7 @@ mod tests { mode: String::from("block"), pid: 1, uid: 0, + sender: None, }); runtime.write().expect("lock").inhibitors.push(InhibitorEntry { what: String::from("shutdown"), @@ -971,6 +996,7 @@ mod tests { mode: String::from("block"), pid: 2, uid: 0, + sender: None, }); let manager = LoginManager::new( @@ -1189,4 +1215,174 @@ mod tests { let snapshot = runtime.read().expect("lock").clone(); assert_eq!(snapshot.inhibit_delay_max_us.load(Ordering::Relaxed), 42); } + + #[test] + fn inhibit_impl_stores_sender_in_entry() { + let runtime = shared_runtime(); + let manager = LoginManager::new( + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/session/c1")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/seat/seat0")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/user/current")).unwrap(), + runtime.clone(), + ); + + let _fd = manager + .inhibit_impl("sleep", "app", "r", "block", Some(String::from(":1.42"))) + .expect("inhibit should succeed"); + + let guard = runtime.read().expect("lock"); + assert_eq!(guard.inhibitors.len(), 1); + assert_eq!(guard.inhibitors[0].sender.as_deref(), Some(":1.42")); + } + + #[test] + fn inhibit_impl_stores_none_sender_when_absent() { + let runtime = shared_runtime(); + let manager = LoginManager::new( + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/session/c1")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/seat/seat0")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/user/current")).unwrap(), + runtime.clone(), + ); + + let _fd = manager + .inhibit_impl("sleep", "app", "r", "delay", None) + .expect("inhibit should succeed"); + + let guard = runtime.read().expect("lock"); + assert_eq!(guard.inhibitors.len(), 1); + assert!(guard.inhibitors[0].sender.is_none()); + } + + #[test] + fn reap_inhibitors_for_sender_removes_matching_entries() { + let runtime = shared_runtime(); + let manager = LoginManager::new( + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/session/c1")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/seat/seat0")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/user/current")).unwrap(), + runtime.clone(), + ); + + let _fd1 = manager + .inhibit_impl("sleep", "app1", "r", "block", Some(String::from(":1.10"))) + .expect("inhibit"); + let _fd2 = manager + .inhibit_impl("shutdown", "app2", "r", "delay", Some(String::from(":1.20"))) + .expect("inhibit"); + let _fd3 = manager + .inhibit_impl("idle", "app3", "r", "block", Some(String::from(":1.10"))) + .expect("inhibit"); + + assert_eq!(runtime.read().expect("lock").inhibitors.len(), 3); + + manager.reap_inhibitors_for_sender(":1.10"); + + let guard = runtime.read().expect("lock"); + assert_eq!(guard.inhibitors.len(), 1); + assert_eq!(guard.inhibitors[0].sender.as_deref(), Some(":1.20")); + } + + #[test] + fn reap_inhibitors_for_sender_drops_daemon_fds() { + let runtime = shared_runtime(); + let manager = LoginManager::new( + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/session/c1")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/seat/seat0")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/user/current")).unwrap(), + runtime.clone(), + ); + + let _fd1 = manager + .inhibit_impl("sleep", "app1", "r", "block", Some(String::from(":1.10"))) + .expect("inhibit"); + let _fd2 = manager + .inhibit_impl("shutdown", "app2", "r", "delay", Some(String::from(":1.20"))) + .expect("inhibit"); + + assert_eq!(manager.inhibitor_fds.lock().expect("lock").len(), 2); + + manager.reap_inhibitors_for_sender(":1.10"); + + let fds = manager.inhibitor_fds.lock().expect("lock"); + assert_eq!(fds.len(), 1); + assert_eq!(fds[0].0.as_deref(), Some(":1.20")); + } + + #[test] + fn reap_inhibitors_for_sender_preserves_no_sender_entries() { + let runtime = shared_runtime(); + let manager = LoginManager::new( + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/session/c1")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/seat/seat0")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/user/current")).unwrap(), + runtime.clone(), + ); + + let _fd = manager + .inhibit_impl("sleep", "legacy", "r", "block", None) + .expect("inhibit"); + + manager.reap_inhibitors_for_sender(":1.99"); + + assert_eq!(runtime.read().expect("lock").inhibitors.len(), 1); + } + + #[test] + fn list_inhibitors_filters_dead_senders() { + let runtime = shared_runtime(); + let manager = LoginManager::new( + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/session/c1")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/seat/seat0")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/user/current")).unwrap(), + runtime.clone(), + ); + + let _fd1 = manager + .inhibit_impl("sleep", "app1", "r", "block", Some(String::from(":1.10"))) + .expect("inhibit"); + let _fd2 = manager + .inhibit_impl("shutdown", "app2", "r", "delay", Some(String::from(":1.20"))) + .expect("inhibit"); + + { + let mut dead = manager.dead_senders.lock().expect("lock"); + dead.insert(String::from(":1.10")); + } + + let list = manager.list_inhibitors().expect("list_inhibitors"); + assert_eq!(list.len(), 1, "ghost entry for :1.10 should be filtered"); + assert_eq!(list[0].1, "app2"); + } + + #[test] + fn list_inhibitors_shows_entries_until_reaped() { + let runtime = shared_runtime(); + let manager = LoginManager::new( + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/session/c1")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/seat/seat0")).unwrap(), + OwnedObjectPath::try_from(String::from("/org/freedesktop/login1/user/current")).unwrap(), + runtime.clone(), + ); + + let _fd = manager + .inhibit_impl("sleep", "app", "r", "block", Some(String::from(":1.42"))) + .expect("inhibit"); + + assert_eq!(manager.list_inhibitors().expect("list").len(), 1); + + manager.reap_inhibitors_for_sender(":1.42"); + + assert_eq!(manager.list_inhibitors().expect("list").len(), 0); + } + + #[test] + fn reap_inhibitors_marks_sender_dead() { + let manager = test_manager(); + + manager.reap_inhibitors_for_sender(":1.77"); + + let dead = manager.dead_senders.lock().expect("lock"); + assert!(dead.contains(":1.77")); + } } diff --git a/local/recipes/system/redbear-sessiond/source/src/runtime_state.rs b/local/recipes/system/redbear-sessiond/source/src/runtime_state.rs index 454a85f274..1c49ff2468 100644 --- a/local/recipes/system/redbear-sessiond/source/src/runtime_state.rs +++ b/local/recipes/system/redbear-sessiond/source/src/runtime_state.rs @@ -11,6 +11,10 @@ pub struct InhibitorEntry { pub mode: String, pub pid: u32, pub uid: u32, + /// Unique bus name of the connection that registered this inhibitor + /// (e.g. `:1.42`). `None` when the sender could not be determined or + /// for entries constructed without sender tracking (legacy / tests). + pub sender: Option, } /// Runtime state for the login1 manager, sessions, and seats. diff --git a/local/recipes/system/redbear-statusnotifierwatcher/recipe.toml b/local/recipes/system/redbear-statusnotifierwatcher/recipe.toml index ae998ff58f..1872860d49 100644 --- a/local/recipes/system/redbear-statusnotifierwatcher/recipe.toml +++ b/local/recipes/system/redbear-statusnotifierwatcher/recipe.toml @@ -2,7 +2,7 @@ name = "redbear-statusnotifierwatcher" version = "0.3.1" -# redbear-statusnotifierwatcher — org.freedesktop.StatusNotifierWatcher daemon. +# redbear-statusnotifierwatcher — org.kde.StatusNotifierWatcher daemon. # Session-bus D-Bus service brokering StatusNotifierItem (system tray) registration. # Pure-Rust (zbus + tokio); no libdbus build dependency. [source] diff --git a/local/recipes/system/redbear-statusnotifierwatcher/source/src/main.rs b/local/recipes/system/redbear-statusnotifierwatcher/source/src/main.rs index 64ef4cdf7f..81fba229c1 100644 --- a/local/recipes/system/redbear-statusnotifierwatcher/source/src/main.rs +++ b/local/recipes/system/redbear-statusnotifierwatcher/source/src/main.rs @@ -1,83 +1,250 @@ -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; +use std::future::poll_fn; +use std::pin::pin; use std::sync::{Arc, Mutex}; -use std::time::Duration; use zbus::{ - fdo::{self, ObjectManager}, - interface, - object_server::SignalEmitter, connection::Builder as ConnectionBuilder, + export::futures_core::Stream, + fdo, + interface, + message::Header, + object_server::SignalEmitter, + proxy, zvariant::ObjectPath, }; const BUS_NAME: &str = "org.freedesktop.StatusNotifierWatcher"; const OBJECT_PATH: &str = "/StatusNotifierWatcher"; -/// org.freedesktop.StatusNotifierWatcher D-Bus interface -/// Tracks registered system tray items and hosts for KDE Plasma. +/// Maximum number of entries (item paths or host names) retained per registry. +/// When this bound is reached the oldest entry is evicted deterministically. +const MAX_ENTRIES: usize = 1024; + +/// Maximum allowed length for any single item-path or host-name string. +const MAX_INPUT_LEN: usize = 256; + +// --------------------------------------------------------------------------- +// Input validation +// --------------------------------------------------------------------------- + +/// Validate an item path or host name received from a D-Bus caller. +/// +/// Rejects empty strings, strings longer than [`MAX_INPUT_LEN`], and strings +/// containing NUL or other control characters. This prevents injection of +/// absurd values that could corrupt logs, exhaust memory, or be used for +/// bus-level confusion attacks. +fn validate_input(s: &str) -> Result<(), String> { + if s.is_empty() { + return Err("empty input".to_owned()); + } + if s.len() > MAX_INPUT_LEN { + return Err(format!( + "input length {} exceeds maximum of {MAX_INPUT_LEN}", + s.len() + )); + } + // is_control() covers NUL, all C0/C1 control codes, and Unicode format chars. + if s.chars().any(|c| c.is_control()) { + return Err("input contains NUL or control characters".to_owned()); + } + Ok(()) +} + +// --------------------------------------------------------------------------- +// Owner-keyed registry with insertion-order tracking +// --------------------------------------------------------------------------- + +/// Maps owner unique bus names (e.g. ``:1.42``) to the set of values they +/// registered. A separate [`Vec`] tracks the global insertion order of +/// values so that oldest-first eviction under [`MAX_ENTRIES`] is deterministic. +struct Registry { + /// owner → set of registered values (item paths or host names). + by_owner: HashMap>, + /// All values across all owners, in insertion (FIFO) order. + insertion_order: Vec, +} + +impl Registry { + fn new() -> Self { + Self { + by_owner: HashMap::new(), + insertion_order: Vec::new(), + } + } + + /// Register ``value`` on behalf of ``owner``. + /// + /// Returns ``true`` if the value was newly added, ``false`` if it was + /// already present (registered by any owner — values are globally unique). + /// When the registry is at capacity the oldest entry is evicted first. + fn register(&mut self, owner: &str, value: &str) -> bool { + // Dedup by value: a given item path / host name can only exist once. + if self.insertion_order.iter().any(|v| v == value) { + return false; + } + if self.insertion_order.len() >= MAX_ENTRIES { + self.evict_oldest(); + } + self.by_owner + .entry(owner.to_owned()) + .or_default() + .insert(value.to_owned()); + self.insertion_order.push(value.to_owned()); + true + } + + /// Unregister ``value``, but only if ``owner`` matches the recorded owner. + /// + /// Returns ``true`` if the entry existed and was removed, ``false`` + /// otherwise (unknown value or caller is not the owner). + fn unregister(&mut self, owner: &str, value: &str) -> bool { + let removed = match self.by_owner.get_mut(owner) { + Some(set) if set.contains(value) => { + set.remove(value); + true + } + _ => false, + }; + if removed { + self.insertion_order.retain(|v| v != value); + // Clean up the owner entry if its set is now empty. + if self.by_owner.get(owner).is_some_and(|s| s.is_empty()) { + self.by_owner.remove(owner); + } + } + removed + } + + /// Return all values in insertion order. + fn snapshot(&self) -> Vec { + self.insertion_order.clone() + } + + fn is_empty(&self) -> bool { + self.insertion_order.is_empty() + } + + #[cfg(test)] + fn total_len(&self) -> usize { + self.insertion_order.len() + } + + /// Remove every entry owned by ``owner``. Called when a bus name + /// vanishes (NameOwnerChanged with empty ``new_owner``). + /// Returns the number of entries removed. + fn purge_owner(&mut self, owner: &str) -> usize { + match self.by_owner.remove(owner) { + Some(set) => { + let count = set.len(); + for value in &set { + self.insertion_order.retain(|v| v != value); + } + count + } + None => 0, + } + } + + /// Drop the oldest value (front of ``insertion_order``) and remove it + /// from whichever owner owns it. + fn evict_oldest(&mut self) { + if let Some(oldest) = self.insertion_order.first().cloned() { + self.insertion_order.remove(0); + for (_, set) in self.by_owner.iter_mut() { + if set.remove(&oldest) { + break; + } + } + self.by_owner.retain(|_, set| !set.is_empty()); + } + } +} + +// --------------------------------------------------------------------------- +// StatusNotifierWatcher D-Bus interface +// --------------------------------------------------------------------------- + +/// org.freedesktop.StatusNotifierWatcher D-Bus interface. +/// +/// Tracks registered system tray items and hosts for KDE Plasma. Each +/// registration is bound to the caller's unique bus name; only the owning +/// caller may unregister its own entries, and entries are purged +/// automatically when the owning bus name vanishes. #[derive(Clone)] struct StatusNotifierWatcher { - items: Arc>>, - hosts: Arc>>, + items: Arc>, + hosts: Arc>, } impl StatusNotifierWatcher { fn new() -> Self { Self { - items: Arc::new(Mutex::new(HashSet::new())), - hosts: Arc::new(Mutex::new(HashSet::new())), + items: Arc::new(Mutex::new(Registry::new())), + hosts: Arc::new(Mutex::new(Registry::new())), } } - /// Register a status notifier item. Returns `true` if the item was - /// newly registered, `false` if it was already present. - fn register_item(&self, item: &str) -> bool { - match self.items.lock() { - Ok(mut items) => items.insert(item.to_owned()), - Err(_) => false, - } + /// Register an item on behalf of ``owner``. Returns ``true`` if newly + /// added. + fn register_item(&self, owner: &str, item: &str) -> bool { + self.items + .lock() + .map(|mut g| g.register(owner, item)) + .unwrap_or(false) } - /// Unregister a status notifier item. Returns `true` if the item - /// was previously registered and is now gone, `false` if it was - /// not in the set. - fn unregister_item(&self, item: &str) -> bool { - match self.items.lock() { - Ok(mut items) => items.remove(item), - Err(_) => false, - } + /// Unregister an item, but only if ``owner`` matches the recorded owner. + fn unregister_item(&self, owner: &str, item: &str) -> bool { + self.items + .lock() + .map(|mut g| g.unregister(owner, item)) + .unwrap_or(false) } - /// Register a status notifier host. Returns `true` if newly - /// registered, `false` if already present. - fn register_host(&self, host: &str) -> bool { - match self.hosts.lock() { - Ok(mut hosts) => hosts.insert(host.to_owned()), - Err(_) => false, - } + /// Register a host on behalf of ``owner``. + fn register_host(&self, owner: &str, host: &str) -> bool { + self.hosts + .lock() + .map(|mut g| g.register(owner, host)) + .unwrap_or(false) } - /// Unregister a status notifier host. Returns `true` if the host - /// was previously registered. - fn unregister_host(&self, host: &str) -> bool { - match self.hosts.lock() { - Ok(mut hosts) => hosts.remove(host), - Err(_) => false, - } + /// Unregister a host, but only if ``owner`` matches the recorded owner. + fn unregister_host(&self, owner: &str, host: &str) -> bool { + self.hosts + .lock() + .map(|mut g| g.unregister(owner, host)) + .unwrap_or(false) } fn items_snapshot(&self) -> Vec { - self.items - .lock() - .map(|g| g.iter().cloned().collect()) - .unwrap_or_default() + self.items.lock().map(|g| g.snapshot()).unwrap_or_default() } fn is_host_registered(&self) -> bool { - self.hosts + self.hosts.lock().map(|g| !g.is_empty()).unwrap_or(false) + } + + #[cfg(test)] + fn items_count(&self) -> usize { + self.items.lock().map(|g| g.total_len()).unwrap_or(0) + } + + /// Remove all items **and** hosts owned by ``owner``. + /// Returns the total number of entries purged. + fn purge_owner(&self, owner: &str) -> usize { + let items = self + .items .lock() - .map(|g| !g.is_empty()) - .unwrap_or(false) + .map(|mut g| g.purge_owner(owner)) + .unwrap_or(0); + let hosts = self + .hosts + .lock() + .map(|mut g| g.purge_owner(owner)) + .unwrap_or(0); + items + hosts } } @@ -86,32 +253,57 @@ impl StatusNotifierWatcher { // --- Methods --- /// Register a status notifier item. - /// The item parameter is either a full object path (e.g., "/org/example/Item") - /// sent by the item itself, or a bus name (e.g., ":1.42" or "org.example.App") - /// sent via the KDE protocol extension. + /// + /// The ``item`` parameter is either a full object path (e.g., + /// ``/org/example/Item``) sent by the item itself, or a bus name + /// (e.g., ``:1.42`` or ``org.example.App``) sent via the KDE protocol + /// extension. The caller's unique bus name is recorded as the owner. async fn register_status_notifier_item( &self, #[zbus(signal_emitter)] signal_emitter: SignalEmitter<'_>, + #[zbus(header)] hdr: Header<'_>, item: &str, ) -> fdo::Result<()> { - let is_new = self.register_item(item); + if let Err(msg) = validate_input(item) { + eprintln!("statusnotifierwatcher: rejected item registration: {msg}"); + return Err(fdo::Error::InvalidArgs(msg)); + } + let owner = hdr + .sender() + .ok_or_else(|| { + fdo::Error::Failed("no sender on RegisterStatusNotifierItem".to_owned()) + })? + .to_string(); + + let is_new = self.register_item(&owner, item); if is_new { - eprintln!("statusnotifierwatcher: item registered: {item}"); + eprintln!("statusnotifierwatcher: item registered: {item} (owner: {owner})"); let _ = Self::status_notifier_item_registered(&signal_emitter, item).await; } Ok(()) } /// Unregister a previously-registered status notifier item. + /// + /// Only succeeds if the caller's unique bus name matches the owner + /// recorded at registration time. async fn unregister_status_notifier_item( &self, - #[zbus(signal_emitter)] _signal_emitter: SignalEmitter<'_>, + #[zbus(signal_emitter)] signal_emitter: SignalEmitter<'_>, + #[zbus(header)] hdr: Header<'_>, item: &str, ) -> fdo::Result<()> { - let was_present = self.unregister_item(item); + let owner = hdr + .sender() + .ok_or_else(|| { + fdo::Error::Failed("no sender on UnregisterStatusNotifierItem".to_owned()) + })? + .to_string(); + + let was_present = self.unregister_item(&owner, item); if was_present { - eprintln!("statusnotifierwatcher: item unregistered: {item}"); - let _ = Self::status_notifier_item_unregistered(&_signal_emitter, item).await; + eprintln!("statusnotifierwatcher: item unregistered: {item} (owner: {owner})"); + let _ = Self::status_notifier_item_unregistered(&signal_emitter, item).await; } Ok(()) } @@ -120,25 +312,48 @@ impl StatusNotifierWatcher { async fn register_status_notifier_host( &self, #[zbus(signal_emitter)] signal_emitter: SignalEmitter<'_>, + #[zbus(header)] hdr: Header<'_>, host: &str, ) -> fdo::Result<()> { - let is_new = self.register_host(host); + if let Err(msg) = validate_input(host) { + eprintln!("statusnotifierwatcher: rejected host registration: {msg}"); + return Err(fdo::Error::InvalidArgs(msg)); + } + let owner = hdr + .sender() + .ok_or_else(|| { + fdo::Error::Failed("no sender on RegisterStatusNotifierHost".to_owned()) + })? + .to_string(); + + let is_new = self.register_host(&owner, host); if is_new { - eprintln!("statusnotifierwatcher: host registered: {host}"); + eprintln!("statusnotifierwatcher: host registered: {host} (owner: {owner})"); let _ = Self::status_notifier_host_registered(&signal_emitter).await; } Ok(()) } /// Unregister a previously-registered status notifier host. + /// + /// Only succeeds if the caller's unique bus name matches the owner + /// recorded at registration time. async fn unregister_status_notifier_host( &self, #[zbus(signal_emitter)] _signal_emitter: SignalEmitter<'_>, + #[zbus(header)] hdr: Header<'_>, host: &str, ) -> fdo::Result<()> { - let was_present = self.unregister_host(host); + let owner = hdr + .sender() + .ok_or_else(|| { + fdo::Error::Failed("no sender on UnregisterStatusNotifierHost".to_owned()) + })? + .to_string(); + + let was_present = self.unregister_host(&owner, host); if was_present { - eprintln!("statusnotifierwatcher: host unregistered: {host}"); + eprintln!("statusnotifierwatcher: host unregistered: {host} (owner: {owner})"); } Ok(()) } @@ -186,12 +401,80 @@ impl StatusNotifierWatcher { ) -> zbus::Result<()>; } +// --------------------------------------------------------------------------- +// NameOwnerChanged listener +// --------------------------------------------------------------------------- + +/// D-Bus daemon proxy for receiving ``NameOwnerChanged`` signals. +#[proxy( + interface = "org.freedesktop.DBus", + default_service = "org.freedesktop.DBus", + default_path = "/org/freedesktop/DBus" +)] +trait DBusDaemon { + #[zbus(signal, name = "NameOwnerChanged")] + fn name_owner_changed(&self, name: String, old_owner: String, new_owner: String); +} + +/// Background task that subscribes to ``org.freedesktop.DBus.NameOwnerChanged`` +/// and prunes all watcher entries whose owner vanished from the bus. +async fn run_name_owner_changed_listener( + connection: zbus::Connection, + watcher: StatusNotifierWatcher, +) { + let proxy = match DBusDaemonProxy::new(&connection).await { + Ok(p) => p, + Err(e) => { + eprintln!("statusnotifierwatcher: failed to create DBus proxy: {e}"); + return; + } + }; + let signals = match proxy.receive_name_owner_changed().await { + Ok(s) => s, + Err(e) => { + eprintln!( + "statusnotifierwatcher: failed to subscribe to NameOwnerChanged: {e}" + ); + return; + } + }; + eprintln!("statusnotifierwatcher: NameOwnerChanged listener active"); + // Consume the signal stream without pulling in futures-lite as a direct + // dependency: zbus re-exports futures_core::Stream, and std provides + // poll_fn + pin! for manual polling. + let mut signals = pin!(signals); + loop { + let item = poll_fn(|cx| signals.as_mut().poll_next(cx)).await; + match item { + Some(signal) => { + if let Ok(args) = signal.args() { + // new_owner is empty when a name is released / client disconnected. + if args.new_owner.is_empty() && !args.old_owner.is_empty() { + let removed = watcher.purge_owner(&args.old_owner); + if removed > 0 { + eprintln!( + "statusnotifierwatcher: purged {removed} entries for vanished owner {}", + args.old_owner + ); + } + } + } + } + None => break, + } + } +} + +// --------------------------------------------------------------------------- +// Main +// --------------------------------------------------------------------------- + async fn wait_for_session_bus() { for _ in 0..30 { if std::env::var("DBUS_SESSION_BUS_ADDRESS").is_ok() { return; } - tokio::time::sleep(Duration::from_millis(200)).await; + tokio::time::sleep(std::time::Duration::from_millis(200)).await; } } @@ -230,12 +513,15 @@ async fn main() -> Result<(), Box> { let path: ObjectPath<'_> = OBJECT_PATH.try_into()?; let connection = ConnectionBuilder::session()? .name(BUS_NAME)? - .serve_at(path, watcher)? + .serve_at(path, watcher.clone())? .build() .await?; eprintln!("statusnotifierwatcher: {BUS_NAME} registered on session bus"); + // Spawn NameOwnerChanged listener that prunes entries on owner loss. + tokio::spawn(run_name_owner_changed_listener(connection.clone(), watcher)); + // Wait for shutdown signal let _ = shutdown_rx.changed().await; eprintln!("statusnotifierwatcher: shutdown signal received, exiting cleanly"); @@ -247,6 +533,11 @@ async fn main() -> Result<(), Box> { mod tests { use super::*; + // ======================================================================= + // Existing tests — adapted for the new ``owner`` parameter. + // The assertions are unchanged; only the call signatures differ. + // ======================================================================= + #[test] fn fresh_watcher_has_no_items_or_hosts() { let w = StatusNotifierWatcher::new(); @@ -257,17 +548,20 @@ mod tests { #[test] fn register_item_is_idempotent() { let w = StatusNotifierWatcher::new(); - assert!(w.register_item("/org/example/Item1")); - assert!(!w.register_item("/org/example/Item1")); - assert_eq!(w.items_snapshot(), vec!["/org/example/Item1".to_string()]); + assert!(w.register_item(":1.1", "/org/example/Item1")); + assert!(!w.register_item(":1.1", "/org/example/Item1")); + assert_eq!( + w.items_snapshot(), + vec!["/org/example/Item1".to_string()] + ); } #[test] fn register_multiple_items() { let w = StatusNotifierWatcher::new(); - w.register_item("/org/example/A"); - w.register_item("/org/example/B"); - w.register_item(":1.42"); + w.register_item(":1.1", "/org/example/A"); + w.register_item(":1.1", "/org/example/B"); + w.register_item(":1.2", ":1.42"); let items = w.items_snapshot(); assert_eq!(items.len(), 3); assert!(items.contains(&"/org/example/A".to_string())); @@ -279,20 +573,20 @@ mod tests { fn register_host_marks_registered() { let w = StatusNotifierWatcher::new(); assert!(!w.is_host_registered()); - w.register_host("org.kde.plasma"); + w.register_host(":1.1", "org.kde.plasma"); assert!(w.is_host_registered()); - w.register_host("org.kde.plasma"); + w.register_host(":1.1", "org.kde.plasma"); assert!(w.is_host_registered()); } #[test] fn items_and_hosts_are_independent() { let w = StatusNotifierWatcher::new(); - w.register_item("/org/example/Item"); - w.register_host("org.kde.plasma"); + w.register_item(":1.1", "/org/example/Item"); + w.register_host(":1.1", "org.kde.plasma"); assert_eq!(w.items_snapshot().len(), 1); assert!(w.is_host_registered()); - w.register_host("org.gnome.shell"); + w.register_host(":1.2", "org.gnome.shell"); assert!(w.is_host_registered()); assert_eq!(w.items_snapshot().len(), 1); } @@ -300,55 +594,219 @@ mod tests { #[test] fn unregister_item_returns_true_when_item_was_registered() { let w = StatusNotifierWatcher::new(); - w.register_item("/org/example/Item"); - assert!(w.unregister_item("/org/example/Item")); + w.register_item(":1.1", "/org/example/Item"); + assert!(w.unregister_item(":1.1", "/org/example/Item")); assert!(w.items_snapshot().is_empty()); } #[test] fn unregister_item_returns_false_when_item_was_not_registered() { let w = StatusNotifierWatcher::new(); - assert!(!w.unregister_item("/org/example/Item")); + assert!(!w.unregister_item(":1.1", "/org/example/Item")); } #[test] fn unregister_item_is_idempotent() { let w = StatusNotifierWatcher::new(); - w.register_item("/org/example/Item"); - assert!(w.unregister_item("/org/example/Item")); - assert!(!w.unregister_item("/org/example/Item")); + w.register_item(":1.1", "/org/example/Item"); + assert!(w.unregister_item(":1.1", "/org/example/Item")); + assert!(!w.unregister_item(":1.1", "/org/example/Item")); } #[test] fn unregister_host_returns_true_when_host_was_registered() { let w = StatusNotifierWatcher::new(); - w.register_host("org.kde.plasma"); - assert!(w.unregister_host("org.kde.plasma")); + w.register_host(":1.1", "org.kde.plasma"); + assert!(w.unregister_host(":1.1", "org.kde.plasma")); assert!(!w.is_host_registered()); } #[test] fn unregister_host_returns_false_when_host_was_not_registered() { let w = StatusNotifierWatcher::new(); - assert!(!w.unregister_host("org.kde.plasma")); + assert!(!w.unregister_host(":1.1", "org.kde.plasma")); } #[test] fn unregister_host_is_idempotent() { let w = StatusNotifierWatcher::new(); - w.register_host("org.kde.plasma"); - assert!(w.unregister_host("org.kde.plasma")); - assert!(!w.unregister_host("org.kde.plasma")); + w.register_host(":1.1", "org.kde.plasma"); + assert!(w.unregister_host(":1.1", "org.kde.plasma")); + assert!(!w.unregister_host(":1.1", "org.kde.plasma")); } #[test] fn unregister_does_not_affect_other_items() { let w = StatusNotifierWatcher::new(); - w.register_item("/org/example/A"); - w.register_item("/org/example/B"); - assert!(w.unregister_item("/org/example/A")); + w.register_item(":1.1", "/org/example/A"); + w.register_item(":1.2", "/org/example/B"); + assert!(w.unregister_item(":1.1", "/org/example/A")); let items = w.items_snapshot(); assert_eq!(items.len(), 1); assert!(items.contains(&"/org/example/B".to_string())); } + + // ======================================================================= + // New tests — ownership enforcement + // ======================================================================= + + #[test] + fn unregister_item_fails_when_caller_is_not_owner() { + let w = StatusNotifierWatcher::new(); + w.register_item(":1.1", "/org/example/Item"); + // Wrong owner must not be able to unregister. + assert!(!w.unregister_item(":1.2", "/org/example/Item")); + // Item is still present. + assert_eq!(w.items_snapshot().len(), 1); + // Correct owner can unregister. + assert!(w.unregister_item(":1.1", "/org/example/Item")); + } + + #[test] + fn unregister_host_fails_when_caller_is_not_owner() { + let w = StatusNotifierWatcher::new(); + w.register_host(":1.1", "org.kde.plasma"); + assert!(!w.unregister_host(":1.2", "org.kde.plasma")); + assert!(w.is_host_registered()); + assert!(w.unregister_host(":1.1", "org.kde.plasma")); + assert!(!w.is_host_registered()); + } + + #[test] + fn same_value_registered_by_different_owner_is_idempotent() { + let w = StatusNotifierWatcher::new(); + assert!(w.register_item(":1.1", "/shared/Item")); + // Same value registered by a different owner is not newly added. + assert!(!w.register_item(":1.2", "/shared/Item")); + // Only one entry exists. + assert_eq!(w.items_snapshot().len(), 1); + // Original owner can still unregister. + assert!(w.unregister_item(":1.1", "/shared/Item")); + } + + #[test] + fn multiple_owners_each_have_independent_entries() { + let w = StatusNotifierWatcher::new(); + w.register_item(":1.1", "/a/1"); + w.register_item(":1.1", "/a/2"); + w.register_item(":1.2", "/b/1"); + assert_eq!(w.items_snapshot().len(), 3); + // Owner :1.1 can remove its own items but not :1.2's. + assert!(w.unregister_item(":1.1", "/a/1")); + assert!(!w.unregister_item(":1.1", "/b/1")); + assert_eq!(w.items_snapshot().len(), 2); + } + + // ======================================================================= + // New tests — purge on NameOwnerChanged + // ======================================================================= + + #[test] + fn purge_owner_removes_all_items_and_hosts_for_owner() { + let w = StatusNotifierWatcher::new(); + w.register_item(":1.10", "/item/a"); + w.register_item(":1.10", "/item/b"); + w.register_item(":1.20", "/item/c"); + w.register_host(":1.10", "org.kde.plasma"); + + let purged = w.purge_owner(":1.10"); + assert_eq!(purged, 3, "should purge 2 items + 1 host"); + + // :1.10's items and hosts are gone; :1.20's item survives. + assert_eq!(w.items_snapshot(), vec!["/item/c".to_string()]); + assert!(!w.is_host_registered()); + } + + #[test] + fn purge_unknown_owner_is_noop() { + let w = StatusNotifierWatcher::new(); + w.register_item(":1.1", "/item/a"); + assert_eq!(w.purge_owner(":9.9"), 0); + assert_eq!(w.items_snapshot().len(), 1); + } + + // ======================================================================= + // New tests — input validation + // ======================================================================= + + #[test] + fn validate_input_rejects_empty() { + assert!(validate_input("").is_err()); + } + + #[test] + fn validate_input_accepts_normal_path() { + assert!(validate_input("/org/example/StatusNotifierItem").is_ok()); + } + + #[test] + fn validate_input_accepts_bus_name() { + assert!(validate_input(":1.42").is_ok()); + assert!(validate_input("org.kde.StatusNotifierItem-1-0").is_ok()); + } + + #[test] + fn validate_input_rejects_control_characters() { + assert!(validate_input("hello\0world").is_err()); + assert!(validate_input("hello\nworld").is_err()); + assert!(validate_input("hello\tworld").is_err()); + assert!(validate_input("\x1b[31mred\x1b[0m").is_err()); + } + + #[test] + fn validate_input_rejects_oversized_string() { + let long = "a".repeat(MAX_INPUT_LEN + 1); + assert!(validate_input(&long).is_err()); + } + + #[test] + fn validate_input_accepts_max_length_string() { + let exact = "a".repeat(MAX_INPUT_LEN); + assert!(validate_input(&exact).is_ok()); + } + + // ======================================================================= + // New tests — bounded entries with oldest-first eviction + // ======================================================================= + + #[test] + fn registry_evicts_oldest_when_full() { + let w = StatusNotifierWatcher::new(); + // Fill to capacity. + for i in 0..MAX_ENTRIES { + w.register_item(":1.1", &format!("/item/{i}")); + } + assert_eq!(w.items_count(), MAX_ENTRIES); + + // Adding one more evicts the oldest (/item/0). + assert!(w.register_item(":1.1", "/item/new")); + assert_eq!(w.items_count(), MAX_ENTRIES); + assert!( + !w.items_snapshot().contains(&"/item/0".to_string()), + "oldest entry should have been evicted" + ); + assert!( + w.items_snapshot().contains(&"/item/new".to_string()), + "newest entry should be present" + ); + assert!( + w.items_snapshot().contains(&"/item/1".to_string()), + "second-oldest should survive" + ); + } + + #[test] + fn registry_eviction_cleans_up_empty_owner() { + let w = StatusNotifierWatcher::new(); + // Fill entirely from one owner. + for i in 0..MAX_ENTRIES { + w.register_item(":1.99", &format!("/x/{i}")); + } + // Evict one more — the owner set shrinks but should still exist. + w.register_item(":1.99", "/x/extra"); + assert_eq!(w.items_count(), MAX_ENTRIES); + // All entries still belong to :1.99. + let purged = w.purge_owner(":1.99"); + assert_eq!(purged, MAX_ENTRIES); + } }