sessiond+statusnotifierwatcher: inhibitor lifecycle, kstop checks, bus name

Three related fixes for redbear-sessiond and redbear-statusnotifierwatcher:

1. StatusNotifierWatcher: change well-known D-Bus name from
   'org.freedesktop.StatusNotifierWatcher' to 'org.kde.StatusNotifierWatcher'.

   Qt tray clients (qdbustrayicon, qdbusmenuconnection) explicitly watch
   for the KDE-prefixed name; the freedesktop-prefixed name left the
   service invisible to any Qt-based system tray. Both the daemon
   BUS_NAME / #[interface(name)] and the D-Bus activation
   /etc/dbus-1/session-services/ file are renamed, and the session
   policy file's <allow own=...> entry is updated. The daemon-level
   doc comment is updated to document why the name is KDE-prefixed.

2. sessiond: replace host-dependent can_methods_return_na test.

   The test asserted hardcoded 'yes' for can_power_off()/reboot()/
   suspend(), which fail on Linux hosts where /scheme/sys/kstop does
   not exist (kstop_writable() returns false). Switch the assertion to
   runtime-detect: compute the expected value from kstop_writable()
   inside the test, so it passes on both the Redox target (yes)
   and a Linux host (na). All 52 sessiond tests now pass on host.

3. sessiond: implement inhibitor lifecycle reaping.

   Inhibit() now captures the caller's unique bus name via zbus
   #[zbus(header)] hdr: Header<'_> + hdr.sender(). The InhibitorEntry
   gains an Option<String> sender field; the daemon-side FDs are now
   tracked with their owner. set_connection() spawns a background
   task that subscribes to org.freedesktop.DBus.NameOwnerChanged via
   zbus::fdo::DBusProxy; when a sender vanishes, all inhibitors and
   FDs owned by it are removed. list_inhibitors() defensively filters
   out entries with dead senders. The test suite gains 8 new tests
   covering sender tracking, reap-by-sender, dead-sender filtering,
   and FD ownership.

   Split inhibit() into the D-Bus-facing method (header-capturing)
   and inhibit_impl() (the testable core). Tests call inhibit_impl
   directly with an explicit sender argument.

Verified: 60/60 sessiond tests pass; 12/12 statusnotifierwatcher
tests pass.
This commit is contained in:
2026-07-27 22:19:01 +09:00
parent 16f74ab87c
commit 4522bc39ca
6 changed files with 803 additions and 148 deletions
@@ -1,3 +0,0 @@
[D-BUS Service]
Name=org.freedesktop.StatusNotifierWatcher
Exec=/usr/bin/redbear-statusnotifierwatcher
@@ -9,8 +9,8 @@
<allow own="org.freedesktop.Notifications"/>
<allow send_destination="org.freedesktop.Notifications"/>
<allow receive_sender="org.freedesktop.Notifications"/>
<allow own="org.freedesktop.StatusNotifierWatcher"/>
<allow send_destination="org.freedesktop.StatusNotifierWatcher"/>
<allow receive_sender="org.freedesktop.StatusNotifierWatcher"/>
<allow own="org.kde.StatusNotifierWatcher"/>
<allow send_destination="org.kde.StatusNotifierWatcher"/>
<allow receive_sender="org.kde.StatusNotifierWatcher"/>
</policy>
</busconfig>
@@ -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<String>,
) -> fdo::Result<OwnedFd> {
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<OwnedFd> {
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<OwnedFd> {
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<String> {
@@ -540,10 +553,21 @@ impl LoginManager {
}
fn list_inhibitors(&self) -> fdo::Result<Vec<(String, String, String, String, u32, u32)>> {
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"));
}
}
@@ -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<String>,
}
/// Runtime state for the login1 manager, sessions, and seats.
@@ -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]
@@ -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<String, HashSet<String>>,
/// All values across all owners, in insertion (FIFO) order.
insertion_order: Vec<String>,
}
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<String> {
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<Mutex<HashSet<String>>>,
hosts: Arc<Mutex<HashSet<String>>>,
items: Arc<Mutex<Registry>>,
hosts: Arc<Mutex<Registry>>,
}
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<String> {
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<dyn std::error::Error>> {
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<dyn std::error::Error>> {
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);
}
}