Skip to main content

rlmctl_core/guard/
notify.rs

1//! Desktop notifications for guard interventions.
2//!
3//! The daemon loop hands every applied [`Action`] to a [`Notifier`], which
4//! keeps one notification per app: "<App> paused" after a freeze, replaced
5//! in place by "<App> slowed down" after a cap, and closed when the app is
6//! released. A thaw that fails replaces it with "<App> could not be
7//! resumed", which stays until a later release of that cgroup succeeds, the
8//! cgroup is gone or no longer frozen, or the guard shuts down. Failed freezes and caps show
9//! nothing (they are in the history and the journal). An optional early
10//! warning ("Memory is running low") is sent at most once a minute while
11//! pressure is High or Critical and no app is held.
12//!
13//! Notifications never feed back into guard decisions: the [`Notifier`]
14//! only reads what the effector already did. The desktop transport,
15//! [`DesktopSink`], runs on its own thread so a slow or missing
16//! notification server can never hold up the guard loop.
17
18use super::effector::Applied;
19use super::policy::is_scarce;
20use super::types::{Action, Level, Sample};
21use common::GuardConfig;
22use std::collections::{HashMap, HashSet};
23use std::path::Path;
24use std::process::Command;
25use std::sync::mpsc;
26use std::time::{Duration, Instant};
27
28/// Icon and desktop entry the notifications carry (the GTK app's).
29pub const APP_ID: &str = "io.github.rlm.gtk";
30/// Application name shown by the notification server.
31pub const APP_NAME: &str = "rlm";
32
33/// Title and body of the early warning.
34pub const PRESSURE_TITLE: &str = "Memory is running low";
35pub const PRESSURE_BODY: &str = "rlm will step in if an app keeps growing.";
36
37/// At most one early warning per this many ms.
38pub const PRESSURE_INTERVAL_MS: u64 = 60_000;
39
40/// Sink key of the early warning. App keys are exe basenames, which never
41/// contain a `/`, so this cannot collide with one.
42pub const PRESSURE_KEY: &str = "/pressure";
43
44/// How the memory looked on one tick, for the early warning.
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub enum Memory {
47    /// The policy did not run this tick.
48    Unknown,
49    /// Not both under High or Critical pressure and short of memory.
50    Fine,
51    /// Pressure is High or Critical and memory is short, as the policy
52    /// judges it before acting ([`is_scarce`]).
53    Low,
54}
55
56/// [`Memory`] for a tick whose policy computed `level` from `sample`.
57pub fn memory_state(level: Level, sample: &Sample, trigger: &common::GuardTrigger) -> Memory {
58    if matches!(level, Level::High | Level::Critical) && is_scarce(sample, trigger) {
59        Memory::Low
60    } else {
61        Memory::Fine
62    }
63}
64
65/// Where notifications go. `key` names one notification: a second `show`
66/// with the same key replaces it in place, `close` removes it.
67pub trait NotifySink {
68    fn show(&mut self, key: &str, title: &str, body: &str);
69    fn close(&mut self, key: &str);
70}
71
72/// Title and body for an app that was just frozen for `hold_secs`.
73pub fn paused_text(app: &str, hold_secs: u64) -> (String, String) {
74    let unit = if hold_secs == 1 { "second" } else { "seconds" };
75    (
76        format!("{app} paused"),
77        // GNOME shows one line of body text in the popup, so keep it short;
78        // the title already names the app.
79        format!("Paused for {hold_secs} {unit} while memory is low."),
80    )
81}
82
83/// Title and body for an app held under a `memory.high` of `cap_bytes`.
84pub fn slowed_text(app: &str, cap_bytes: u64) -> (String, String) {
85    (
86        format!("{app} slowed down"),
87        format!(
88            "Held to about {} until memory frees up.",
89            format_size(cap_bytes)
90        ),
91    )
92}
93
94/// Title and body for an app whose thaw failed, so it may still be frozen.
95pub fn failed_resume_text(app: &str) -> (String, String) {
96    (
97        format!("{app} could not be resumed"),
98        "It may stay paused. Run rlm guard status for details.".to_string(),
99    )
100}
101
102/// A byte count in the decimal units desktops show: "3.2 GB", "512 MB".
103pub fn format_size(bytes: u64) -> String {
104    const GB: f64 = 1e9;
105    const MB: f64 = 1e6;
106    let b = bytes as f64;
107    if b >= GB {
108        format!("{:.1} GB", b / GB)
109    } else {
110        format!("{:.0} MB", (b / MB).max(1.0))
111    }
112}
113
114pub use crate::appname::display_name;
115use crate::appname::DesktopNames;
116
117/// Resolves app keys to [`display_name`]s from the system: installed
118/// desktop entries (read once) and a member of the cgroup, for the
119/// directory its executable is in and, for versioned binaries, its process
120/// name.
121pub struct AppNames {
122    desktop: DesktopNames,
123    /// Desktop entries still being read on a background thread.
124    loading: Option<mpsc::Receiver<DesktopNames>>,
125}
126
127impl AppNames {
128    /// Read the desktop entries now.
129    pub fn new() -> Self {
130        Self {
131            desktop: crate::desktop::installed_names(),
132            loading: None,
133        }
134    }
135
136    /// Read the desktop entries on a background thread, so the caller (the
137    /// guard loop) never waits on the disk. Until they are read, names fall
138    /// back to the basename rules of [`display_name`].
139    pub fn in_background() -> Self {
140        let (tx, rx) = mpsc::channel();
141        let spawned = std::thread::Builder::new()
142            .name("rlm-app-names".into())
143            .spawn(move || {
144                let _ = tx.send(crate::desktop::installed_names());
145            });
146        Self {
147            desktop: DesktopNames::default(),
148            loading: spawned.ok().map(|_| rx),
149        }
150    }
151
152    /// The display name of `key`, whose processes live under `cgroup`.
153    pub fn name(&mut self, key: &str, cgroup: &str) -> String {
154        if let Some(rx) = &self.loading {
155            match rx.try_recv() {
156                Ok(desktop) => {
157                    self.desktop = desktop;
158                    self.loading = None;
159                }
160                Err(mpsc::TryRecvError::Disconnected) => self.loading = None,
161                Err(mpsc::TryRecvError::Empty) => {}
162            }
163        }
164        let base = key.split('@').next().unwrap_or(key);
165        let member = member_in(cgroup, base);
166        let comm = if base.chars().any(char::is_alphabetic) {
167            None
168        } else {
169            member.and_then(comm_of)
170        };
171        let dir = member.and_then(crate::appname::exe_dir_of_pid);
172        display_name(key, &self.desktop, comm.as_deref(), dir.as_deref())
173    }
174}
175
176impl Default for AppNames {
177    fn default() -> Self {
178        Self::new()
179    }
180}
181
182/// The first process under `cgroup` running `exe`.
183fn member_in(cgroup: &str, exe: &str) -> Option<u32> {
184    super::cgfs::pids_under(cgroup)
185        .into_iter()
186        .find(|&p| super::cgfs::exe_basename(p).as_deref() == Some(exe))
187}
188
189/// The process name of `pid`.
190fn comm_of(pid: u32) -> Option<String> {
191    std::fs::read_to_string(format!("/proc/{pid}/comm"))
192        .ok()
193        .map(|c| c.trim().to_string())
194        .filter(|c| !c.is_empty())
195}
196
197/// Whether the cgroup directory `dir` may still be frozen: a freeze is
198/// requested (`cgroup.freeze` is 1) or `cgroup.events` says `frozen 1`.
199/// False when neither can be read.
200fn frozen_at(dir: &Path) -> bool {
201    let read = |file: &str| std::fs::read_to_string(dir.join(file)).ok();
202    read("cgroup.freeze").and_then(|c| super::cgfs::parse_freeze(&c)) == Some(true)
203        || read("cgroup.events").and_then(|c| super::cgfs::parse_frozen(&c)) == Some(true)
204}
205
206/// The two notification settings, the only ones the guard applies live.
207#[derive(Debug, Clone, Copy, PartialEq, Eq)]
208pub struct NotifyFlags {
209    /// `guard.notify`: send notifications at all.
210    pub notify: bool,
211    /// `guard.notify_pressure`: also send the early warning.
212    pub pressure: bool,
213}
214
215impl NotifyFlags {
216    pub fn from_config(g: &GuardConfig) -> Self {
217        Self {
218            notify: g.notify,
219            pressure: g.notify_pressure,
220        }
221    }
222}
223
224/// The flags to use after the config file changed on disk: the reloaded
225/// file's two notification flags when it is valid, else the `current` ones.
226/// Nothing else in the reloaded config is used.
227pub fn flags_after_reload(
228    current: NotifyFlags,
229    reloaded: std::result::Result<&GuardConfig, &common::Error>,
230) -> NotifyFlags {
231    match reloaded {
232        Ok(g) => NotifyFlags::from_config(g),
233        Err(_) => current,
234    }
235}
236
237/// One cgroup the guard holds, as far as notifications care.
238#[derive(Debug, Clone)]
239struct Held {
240    app: String,
241    /// `Some` once capped: the `memory.high` written, in bytes.
242    cap_bytes: Option<u64>,
243}
244
245/// Keeps one notification per held app in step with what the effector did.
246///
247/// Call [`record`](Self::record) for every applied action, then
248/// [`end_tick`](Self::end_tick) once per loop iteration; notifications are
249/// only sent from `end_tick`, so a thaw and a cap of the same app in one tick
250/// replace the notification instead of closing and reopening it. A release
251/// by thaw also waits one more tick before closing, since the policy caps a
252/// still-hot app on the tick after its thaw.
253pub struct Notifier<S: NotifySink> {
254    sink: S,
255    flags: NotifyFlags,
256    freeze_hold_secs: u64,
257    /// Held cgroups, keyed by cgroup path.
258    held: HashMap<String, Held>,
259    /// Display name of each held or shown app, keyed by app.
260    names: HashMap<String, String>,
261    /// What each app's notification shows now, keyed by app.
262    shown: HashMap<String, (String, String)>,
263    /// Apps released by a thaw during the current tick.
264    thawed: HashSet<String>,
265    /// Cgroups whose thaw failed, so they may still be frozen, with their app.
266    stuck: HashMap<String, String>,
267    /// Whether a cgroup may still be frozen: it exists and either a freeze
268    /// is requested (`cgroup.freeze` is 1; tasks in uninterruptible sleep can
269    /// keep `frozen` at 0 for a while) or `cgroup.events` says `frozen 1`.
270    /// A `stuck` entry for which this is
271    /// false is dropped: the cgroup is gone, or something (a later thaw by
272    /// hand, or the app's cgroup being recreated) left it running.
273    still_frozen: Box<dyn Fn(&str) -> bool>,
274    pressure_shown: bool,
275    last_pressure_ms: Option<u64>,
276}
277
278impl<S: NotifySink> Notifier<S> {
279    pub fn new(sink: S, cfg: &GuardConfig) -> Self {
280        Self {
281            sink,
282            flags: NotifyFlags::from_config(cfg),
283            freeze_hold_secs: cfg.timing.freeze_hold_secs,
284            held: HashMap::new(),
285            names: HashMap::new(),
286            shown: HashMap::new(),
287            thawed: HashSet::new(),
288            stuck: HashMap::new(),
289            still_frozen: Box::new(|cg| frozen_at(&super::cgfs::abs(cg))),
290            pressure_shown: false,
291            last_pressure_ms: None,
292        }
293    }
294
295    pub fn sink(&self) -> &S {
296        &self.sink
297    }
298
299    pub fn flags(&self) -> NotifyFlags {
300        self.flags
301    }
302
303    /// Apply new notification settings. Turning `notify` off closes what is
304    /// shown; turning the early warning off closes it.
305    pub fn set_flags(&mut self, flags: NotifyFlags) {
306        if !flags.notify {
307            self.close_all();
308        } else if !flags.pressure {
309            self.close_pressure();
310        }
311        self.flags = flags;
312    }
313
314    /// Note one action the effector applied. `applied` is `None` when it
315    /// failed. Cheap: nothing is sent or looked up until
316    /// [`end_tick`](Self::end_tick).
317    pub fn record(&mut self, action: &Action, applied: Option<Applied>) {
318        match action {
319            Action::Freeze { res, name: app } | Action::Cap { res, name: app } => {
320                let Some(applied) = applied else {
321                    return;
322                };
323                let cap_bytes = match action {
324                    Action::Cap { .. } => Some(applied.cap_bytes.unwrap_or(0)),
325                    _ => None,
326                };
327                self.held.insert(
328                    res.cgroup.clone(),
329                    Held {
330                        app: app.clone(),
331                        cap_bytes,
332                    },
333                );
334            }
335            // A failed thaw may leave the app frozen, so its notification
336            // says so instead of closing.
337            Action::Thaw { res } => {
338                if let Some(h) = self.held.remove(&res.cgroup) {
339                    if applied.is_some() {
340                        self.thawed.insert(h.app);
341                    } else {
342                        self.stuck.insert(res.cgroup.clone(), h.app);
343                    }
344                }
345                if applied.is_some() {
346                    self.stuck.remove(&res.cgroup);
347                }
348            }
349            // A cap lift closes the notification even if it reported an
350            // error: the guard no longer holds the cgroup either way.
351            Action::LiftCap { res } => {
352                self.held.remove(&res.cgroup);
353                if applied.is_some() {
354                    self.stuck.remove(&res.cgroup);
355                }
356            }
357        }
358    }
359
360    /// Send what changed this tick. `memory` is how the memory looked on this
361    /// tick (see [`memory_state`]). `name` maps an app key
362    /// and one of its cgroups to a display name; it is called once per newly
363    /// held app, after every action of the tick was applied.
364    pub fn end_tick(
365        &mut self,
366        now_ms: u64,
367        memory: Memory,
368        name: &mut dyn FnMut(&str, &str) -> String,
369    ) {
370        let thawed = std::mem::take(&mut self.thawed);
371        // The policy has dropped a cgroup whose thaw failed, so no later
372        // release will clear it; stop saying it may be paused once it is
373        // gone or no longer frozen.
374        let frozen = &self.still_frozen;
375        self.stuck.retain(|cg, _| frozen(cg));
376        if !self.flags.notify {
377            return;
378        }
379
380        let mut held: Vec<(&String, &Held)> = self.held.iter().collect();
381        held.sort_by(|a, b| a.0.cmp(b.0));
382        for (cgroup, h) in &held {
383            if !self.names.contains_key(&h.app) {
384                self.names.insert(h.app.clone(), name(&h.app, cgroup));
385            }
386        }
387        let mut want: HashMap<String, (String, String)> = HashMap::new();
388        for (_, h) in held {
389            let display = self
390                .names
391                .get(&h.app)
392                .map_or(h.app.as_str(), |n| n.as_str());
393            let capped: u64 = self
394                .held
395                .values()
396                .filter(|o| o.app == h.app)
397                .filter_map(|o| o.cap_bytes)
398                .sum();
399            let any_capped = self
400                .held
401                .values()
402                .any(|o| o.app == h.app && o.cap_bytes.is_some());
403            let text = if any_capped {
404                slowed_text(display, capped)
405            } else {
406                paused_text(display, self.freeze_hold_secs)
407            };
408            want.entry(h.app.clone()).or_insert(text);
409        }
410        let mut stuck: Vec<(&String, &String)> = self.stuck.iter().collect();
411        stuck.sort();
412        for (cgroup, app) in stuck {
413            if !self.names.contains_key(app) {
414                self.names.insert(app.clone(), name(app, cgroup));
415            }
416            let display = self.names.get(app).map_or(app.as_str(), |n| n.as_str());
417            want.entry(app.clone())
418                .or_insert_with(|| failed_resume_text(display));
419        }
420
421        let mut gone: Vec<String> = self
422            .shown
423            .keys()
424            .filter(|app| !want.contains_key(*app) && !thawed.contains(*app))
425            .cloned()
426            .collect();
427        gone.sort();
428        for app in gone {
429            self.sink.close(&app);
430            self.shown.remove(&app);
431        }
432        let (held, shown, stuck) = (&self.held, &self.shown, &self.stuck);
433        self.names.retain(|app, _| {
434            shown.contains_key(app)
435                || held.values().any(|h| h.app == *app)
436                || stuck.values().any(|a| a == app)
437        });
438        let mut changed: Vec<(String, (String, String))> = want
439            .into_iter()
440            .filter(|(app, text)| self.shown.get(app) != Some(text))
441            .collect();
442        changed.sort();
443        for (app, (title, body)) in changed {
444            self.sink.show(&app, &title, &body);
445            self.shown.insert(app, (title, body));
446        }
447
448        let pressing = memory == Memory::Low;
449        if self.held.is_empty() && self.shown.is_empty() && pressing {
450            let due = self
451                .last_pressure_ms
452                .is_none_or(|last| now_ms.saturating_sub(last) >= PRESSURE_INTERVAL_MS);
453            if self.flags.pressure && due {
454                self.sink.show(PRESSURE_KEY, PRESSURE_TITLE, PRESSURE_BODY);
455                self.pressure_shown = true;
456                self.last_pressure_ms = Some(now_ms);
457            }
458        } else if memory != Memory::Unknown || !self.held.is_empty() {
459            self.close_pressure();
460        }
461    }
462
463    /// Close every notification this notifier has open (guard shutdown).
464    pub fn close_all(&mut self) {
465        self.stuck.clear();
466        let mut apps: Vec<String> = self.shown.drain().map(|(app, _)| app).collect();
467        apps.sort();
468        for app in apps {
469            self.sink.close(&app);
470        }
471        self.close_pressure();
472    }
473
474    fn close_pressure(&mut self) {
475        if self.pressure_shown {
476            self.sink.close(PRESSURE_KEY);
477            self.pressure_shown = false;
478        }
479    }
480}
481
482/// Longest wait for one call to the notification server.
483const CALL_TIMEOUT: Duration = Duration::from_secs(2);
484
485/// The `expire_timeout` sent with every notification: -1, the server's
486/// default. A finite value is the popup's on-screen lifetime on KDE, dunst,
487/// mako and xfce4-notifyd, and a cap can outlast any fixed limit; the guard
488/// closes its notifications itself when it releases an app.
489const EXPIRE_TIMEOUT: i32 = -1;
490/// Requests queued for the sender thread; more are dropped.
491const QUEUE: usize = 64;
492
493const NOTIFY_DEST: &str = "org.freedesktop.Notifications";
494const NOTIFY_PATH: &str = "/org/freedesktop/Notifications";
495
496enum Cmd {
497    Show {
498        key: String,
499        title: String,
500        body: String,
501    },
502    Close {
503        key: String,
504    },
505    Flush(mpsc::Sender<()>),
506}
507
508/// Sends notifications to the desktop from a background thread, through the
509/// freedesktop Notifications D-Bus API on the session bus (`Notify` with
510/// `replaces_id` to update in place, `CloseNotification` to clear). The
511/// thread keeps the server's notification id per key. Without a session
512/// bus it falls back to `notify-send`, which can neither update nor close,
513/// and tries the bus again after 30 s; a connection that breaks is dropped
514/// and rebuilt on the next request. Requests never block the caller: they are queued, and dropped when the
515/// queue is full. Every failure is logged at debug and otherwise ignored.
516pub struct DesktopSink {
517    tx: Option<mpsc::SyncSender<Cmd>>,
518}
519
520impl DesktopSink {
521    /// Start the sender thread. It connects to the session bus on the first
522    /// request, so a guard that never notifies never connects.
523    pub fn spawn() -> Self {
524        let (tx, rx) = mpsc::sync_channel(QUEUE);
525        let started = std::thread::Builder::new()
526            .name("rlm-notify".into())
527            .spawn(move || sender(rx));
528        match started {
529            Ok(_) => Self { tx: Some(tx) },
530            Err(e) => {
531                tracing::debug!(error = %e, "cannot start the notification thread");
532                Self { tx: None }
533            }
534        }
535    }
536
537    /// Wait up to `timeout` for every queued request to be sent. Returns
538    /// false on timeout.
539    pub fn flush(&self, timeout: Duration) -> bool {
540        let (done, wait) = mpsc::channel();
541        if !self.send(Cmd::Flush(done)) {
542            return false;
543        }
544        wait.recv_timeout(timeout).is_ok()
545    }
546
547    fn send(&self, cmd: Cmd) -> bool {
548        match &self.tx {
549            Some(tx) => match tx.try_send(cmd) {
550                Ok(()) => true,
551                Err(e) => {
552                    tracing::debug!(error = %e, "notification dropped");
553                    false
554                }
555            },
556            None => false,
557        }
558    }
559}
560
561impl NotifySink for DesktopSink {
562    fn show(&mut self, key: &str, title: &str, body: &str) {
563        self.send(Cmd::Show {
564            key: key.to_string(),
565            title: title.to_string(),
566            body: body.to_string(),
567        });
568    }
569
570    fn close(&mut self, key: &str) {
571        self.send(Cmd::Close {
572            key: key.to_string(),
573        });
574    }
575}
576
577/// The sender thread: hands every request to a [`Sender`] over the real
578/// session bus.
579fn sender(rx: mpsc::Receiver<Cmd>) {
580    let mut sender = Sender::new(SessionBus);
581    for cmd in rx {
582        sender.handle(cmd, Instant::now());
583    }
584}
585
586/// How long a failed connect is remembered before the next request tries
587/// again. Requests in between use `notify-send`.
588const RECONNECT_AFTER: Duration = Duration::from_secs(30);
589
590/// Why a call to the notification server failed.
591#[derive(Debug)]
592enum BusError {
593    /// The connection or the server is gone: reconnect, and forget the ids.
594    Gone(String),
595    /// Only this call failed (for example it timed out).
596    Call(String),
597}
598
599/// One live connection to the notification server.
600trait Bus {
601    fn notify(&self, replaces_id: u32, title: &str, body: &str) -> Result<u32, BusError>;
602    fn close(&self, id: u32) -> Result<(), BusError>;
603}
604
605/// Makes [`Bus`] connections, and sends without one.
606trait Connector {
607    type Bus: Bus;
608    fn connect(&mut self) -> Option<Self::Bus>;
609    /// Show a notification without a bus (it cannot be updated or closed).
610    fn fallback(&mut self, title: &str, body: &str);
611}
612
613/// The sender's state: the connection (or when to retry one) and the
614/// server's id for each shown notification.
615struct Sender<C: Connector> {
616    connector: C,
617    bus: Option<C::Bus>,
618    /// After a failed connect: no new attempt before this instant.
619    retry_at: Option<Instant>,
620    ids: HashMap<String, u32>,
621}
622
623impl<C: Connector> Sender<C> {
624    fn new(connector: C) -> Self {
625        Self {
626            connector,
627            bus: None,
628            retry_at: None,
629            ids: HashMap::new(),
630        }
631    }
632
633    /// Connect if there is no connection and the retry delay has passed.
634    fn ensure_bus(&mut self, now: Instant) {
635        if self.bus.is_some() || self.retry_at.is_some_and(|t| now < t) {
636            return;
637        }
638        self.bus = self.connector.connect();
639        self.retry_at = if self.bus.is_none() {
640            Some(now + RECONNECT_AFTER)
641        } else {
642            None
643        };
644    }
645
646    /// Drop a connection that is gone. The ids belonged to it (or to a
647    /// server that is gone), so they are forgotten too.
648    fn reset(&mut self, why: &str) {
649        tracing::debug!(error = why, "notification connection lost; reconnecting");
650        self.bus = None;
651        self.retry_at = None;
652        self.ids.clear();
653    }
654
655    fn handle(&mut self, cmd: Cmd, now: Instant) {
656        if let Cmd::Flush(done) = cmd {
657            let _ = done.send(());
658            return;
659        }
660        self.ensure_bus(now);
661        match cmd {
662            Cmd::Show { key, title, body } => {
663                let Some(bus) = &self.bus else {
664                    self.connector.fallback(&title, &body);
665                    return;
666                };
667                let replaces = self.ids.get(&key).copied().unwrap_or(0);
668                match bus.notify(replaces, &title, &body) {
669                    Ok(id) => {
670                        self.ids.insert(key, id);
671                    }
672                    Err(BusError::Gone(e)) => {
673                        self.reset(&e);
674                        // Show it anyway, on a new connection if one comes up.
675                        self.ensure_bus(now);
676                        match &self.bus {
677                            Some(bus) => match bus.notify(0, &title, &body) {
678                                Ok(id) => {
679                                    self.ids.insert(key, id);
680                                }
681                                Err(e) => tracing::debug!(error = ?e, "Notify failed"),
682                            },
683                            None => self.connector.fallback(&title, &body),
684                        }
685                    }
686                    Err(BusError::Call(e)) => tracing::debug!(error = e, "Notify failed"),
687                }
688            }
689            Cmd::Close { key } => {
690                let (Some(bus), Some(id)) = (&self.bus, self.ids.remove(&key)) else {
691                    return;
692                };
693                match bus.close(id) {
694                    Ok(()) => {}
695                    Err(BusError::Gone(e)) => self.reset(&e),
696                    Err(BusError::Call(e)) => {
697                        tracing::debug!(error = e, "CloseNotification failed")
698                    }
699                }
700            }
701            Cmd::Flush(_) => {}
702        }
703    }
704}
705
706/// The real [`Connector`]: the session bus, with `notify-send` as fallback.
707struct SessionBus;
708
709impl Connector for SessionBus {
710    type Bus = zbus::blocking::Connection;
711
712    fn connect(&mut self) -> Option<Self::Bus> {
713        connect()
714    }
715
716    fn fallback(&mut self, title: &str, body: &str) {
717        notify_send(title, body);
718    }
719}
720
721impl Bus for zbus::blocking::Connection {
722    fn notify(&self, replaces_id: u32, title: &str, body: &str) -> Result<u32, BusError> {
723        dbus_notify(self, replaces_id, title, body, EXPIRE_TIMEOUT).map_err(classify)
724    }
725
726    fn close(&self, id: u32) -> Result<(), BusError> {
727        self.call_method(
728            Some(NOTIFY_DEST),
729            NOTIFY_PATH,
730            Some(NOTIFY_DEST),
731            "CloseNotification",
732            &(id,),
733        )
734        .map(|_| ())
735        .map_err(classify)
736    }
737}
738
739/// Sort a zbus error: an I/O failure other than a timeout means the
740/// connection is broken, and a missing or disconnected service means the
741/// notification server went away; both are [`BusError::Gone`]. Anything
742/// else, including a call that timed out, fails only that call.
743fn classify(e: zbus::Error) -> BusError {
744    let gone = match &e {
745        zbus::Error::InputOutput(io) => io.kind() != std::io::ErrorKind::TimedOut,
746        zbus::Error::MethodError(name, _, _) => matches!(
747            name.as_str(),
748            "org.freedesktop.DBus.Error.ServiceUnknown"
749                | "org.freedesktop.DBus.Error.Disconnected"
750                | "org.freedesktop.DBus.Error.NoServer"
751        ),
752        _ => false,
753    };
754    if gone {
755        BusError::Gone(e.to_string())
756    } else {
757        BusError::Call(e.to_string())
758    }
759}
760
761/// Connect to the session bus, giving up after [`CALL_TIMEOUT`] so a
762/// wedged bus cannot stall the sender; `None` means use `notify-send`.
763fn connect() -> Option<zbus::blocking::Connection> {
764    let conn = within(CALL_TIMEOUT, || {
765        zbus::blocking::connection::Builder::session()
766            .map(|b| b.method_timeout(CALL_TIMEOUT))
767            .and_then(|b| b.build())
768            .map_err(|e| tracing::debug!(error = %e, "cannot connect to the session bus"))
769            .ok()
770    });
771    if conn.is_none() {
772        tracing::debug!("no session bus; notifications use notify-send");
773    }
774    conn
775}
776
777/// Run `f` on a helper thread and wait at most `timeout` for it. On timeout
778/// the thread is left to finish on its own and its result is dropped.
779fn within<T: Send + 'static>(
780    timeout: Duration,
781    f: impl FnOnce() -> Option<T> + Send + 'static,
782) -> Option<T> {
783    let (tx, rx) = mpsc::channel();
784    std::thread::Builder::new()
785        .name("rlm-notify-connect".into())
786        .spawn(move || {
787            let _ = tx.send(f());
788        })
789        .ok()?;
790    rx.recv_timeout(timeout).ok().flatten()
791}
792
793fn dbus_notify(
794    conn: &zbus::blocking::Connection,
795    replaces_id: u32,
796    title: &str,
797    body: &str,
798    expire_timeout: i32,
799) -> zbus::Result<u32> {
800    use zbus::zvariant::Value;
801    let mut hints: HashMap<&str, Value<'_>> = HashMap::new();
802    hints.insert("desktop-entry", Value::from(APP_ID));
803    // Urgency normal.
804    hints.insert("urgency", Value::U8(1));
805    let actions: Vec<&str> = Vec::new();
806    let reply = conn.call_method(
807        Some(NOTIFY_DEST),
808        NOTIFY_PATH,
809        Some(NOTIFY_DEST),
810        "Notify",
811        &(
812            APP_NAME,
813            replaces_id,
814            APP_ID,
815            title,
816            body,
817            actions,
818            hints,
819            expire_timeout,
820        ),
821    )?;
822    reply.body().deserialize::<u32>()
823}
824
825/// Fallback without a session bus: a plain `notify-send`, reaped on its own
826/// thread so a hung one cannot stall the sender.
827fn notify_send(title: &str, body: &str) {
828    match Command::new("notify-send")
829        .args(["-a", APP_NAME, "-i", APP_ID, title, body])
830        .spawn()
831    {
832        Ok(mut child) => {
833            std::thread::spawn(move || {
834                let _ = child.wait();
835            });
836        }
837        Err(e) => tracing::debug!(error = %e, "notify-send unavailable; skipping notification"),
838    }
839}
840
841#[cfg(test)]
842mod tests {
843    use super::super::resolve::{Coverage, Mechanism, Resolution, Verdict};
844    use super::*;
845
846    #[derive(Debug, Clone, PartialEq, Eq)]
847    enum Call {
848        Show(String, String, String),
849        Close(String),
850    }
851
852    #[derive(Default)]
853    struct Fake(Vec<Call>);
854
855    impl NotifySink for Fake {
856        fn show(&mut self, key: &str, title: &str, body: &str) {
857            self.0
858                .push(Call::Show(key.into(), title.into(), body.into()));
859        }
860        fn close(&mut self, key: &str) {
861            self.0.push(Call::Close(key.into()));
862        }
863    }
864
865    impl Notifier<Fake> {
866        fn take(&mut self) -> Vec<Call> {
867            std::mem::take(&mut self.sink.0)
868        }
869    }
870
871    const GB: u64 = 1_000_000_000;
872
873    fn res(cg: &str) -> Resolution {
874        Resolution {
875            cgroup: cg.into(),
876            unit: None,
877            verdict: Verdict::Freeze,
878            coverage: Coverage::Full,
879            mechanism: Mechanism::Raw,
880        }
881    }
882
883    fn freeze(cg: &str) -> Action {
884        Action::Freeze {
885            res: res(cg),
886            name: "firefox".into(),
887        }
888    }
889
890    fn cap(cg: &str) -> Action {
891        Action::Cap {
892            res: res(cg),
893            name: "firefox".into(),
894        }
895    }
896
897    fn ok() -> Option<Applied> {
898        Some(Applied { cap_bytes: None })
899    }
900
901    fn capped(bytes: u64) -> Option<Applied> {
902        Some(Applied {
903            cap_bytes: Some(bytes),
904        })
905    }
906
907    fn notifier(notify: bool, pressure: bool) -> Notifier<Fake> {
908        let cfg = GuardConfig {
909            notify,
910            notify_pressure: pressure,
911            ..GuardConfig::default()
912        };
913        let mut n = Notifier::new(Fake::default(), &cfg);
914        n.still_frozen = Box::new(|_| true);
915        n
916    }
917
918    fn names(key: &str, _cg: &str) -> String {
919        display_name(key, &DesktopNames::default(), None, None)
920    }
921
922    fn paused() -> Call {
923        let (t, b) = paused_text("Firefox", 5);
924        Call::Show("firefox".into(), t, b)
925    }
926
927    fn slowed(bytes: u64) -> Call {
928        let (t, b) = slowed_text("Firefox", bytes);
929        Call::Show("firefox".into(), t, b)
930    }
931
932    #[test]
933    fn texts_read_as_specified() {
934        assert_eq!(
935            paused_text("Firefox", 5),
936            (
937                "Firefox paused".to_string(),
938                "Paused for 5 seconds while memory is low.".to_string()
939            )
940        );
941        assert_eq!(
942            slowed_text("Firefox", 3_200_000_000),
943            (
944                "Firefox slowed down".to_string(),
945                "Held to about 3.2 GB until memory frees up.".to_string()
946            )
947        );
948        assert!(paused_text("X", 1).1.contains("for 1 second while"));
949        assert_eq!(format_size(512_000_000), "512 MB");
950        assert_eq!(format_size(268_435_456), "268 MB");
951    }
952
953    #[test]
954    fn freeze_shows_paused_once_for_all_of_an_apps_cgroups() {
955        let mut n = notifier(true, false);
956        n.record(&freeze("/a"), ok());
957        n.record(&freeze("/b"), ok());
958        n.end_tick(0, Memory::Low, &mut names);
959        assert_eq!(n.take(), vec![paused()]);
960        n.end_tick(1_000, Memory::Low, &mut names);
961        assert!(n.take().is_empty(), "nothing changed, nothing sent");
962    }
963
964    #[test]
965    fn thaw_then_cap_in_one_tick_replaces_without_close() {
966        let mut n = notifier(true, false);
967        n.record(&freeze("/a"), ok());
968        n.end_tick(0, Memory::Low, &mut names);
969        n.take();
970        n.record(&Action::Thaw { res: res("/a") }, ok());
971        n.record(&cap("/a"), capped(3 * GB));
972        n.end_tick(5_000, Memory::Low, &mut names);
973        assert_eq!(n.take(), vec![slowed(3 * GB)]);
974    }
975
976    #[test]
977    fn cap_on_the_tick_after_a_thaw_still_replaces() {
978        let mut n = notifier(true, false);
979        n.record(&freeze("/a"), ok());
980        n.end_tick(0, Memory::Low, &mut names);
981        n.take();
982        n.record(&Action::Thaw { res: res("/a") }, ok());
983        n.end_tick(5_000, Memory::Low, &mut names);
984        assert!(n.take().is_empty(), "the close waits a tick");
985        n.record(&cap("/a"), capped(2 * GB));
986        n.end_tick(6_000, Memory::Low, &mut names);
987        assert_eq!(n.take(), vec![slowed(2 * GB)]);
988    }
989
990    #[test]
991    fn thaw_alone_closes() {
992        let mut n = notifier(true, false);
993        n.record(&freeze("/a"), ok());
994        n.end_tick(0, Memory::Low, &mut names);
995        n.take();
996        n.record(&Action::Thaw { res: res("/a") }, ok());
997        n.end_tick(5_000, Memory::Fine, &mut names);
998        n.end_tick(6_000, Memory::Fine, &mut names);
999        assert_eq!(n.take(), vec![Call::Close("firefox".into())]);
1000    }
1001
1002    #[test]
1003    fn lift_closes_in_the_same_tick() {
1004        let mut n = notifier(true, false);
1005        n.record(&cap("/a"), capped(GB));
1006        n.end_tick(0, Memory::Low, &mut names);
1007        assert_eq!(n.take(), vec![slowed(GB)]);
1008        n.record(&Action::LiftCap { res: res("/a") }, ok());
1009        n.end_tick(40_000, Memory::Fine, &mut names);
1010        assert_eq!(n.take(), vec![Call::Close("firefox".into())]);
1011    }
1012
1013    #[test]
1014    fn a_release_that_reported_an_error_still_closes() {
1015        let mut n = notifier(true, false);
1016        n.record(&cap("/a"), capped(GB));
1017        n.end_tick(0, Memory::Low, &mut names);
1018        n.take();
1019        n.record(&Action::LiftCap { res: res("/a") }, None);
1020        n.end_tick(1_000, Memory::Fine, &mut names);
1021        assert_eq!(n.take(), vec![Call::Close("firefox".into())]);
1022    }
1023
1024    fn not_resumed() -> Call {
1025        let (t, b) = failed_resume_text("Firefox");
1026        Call::Show("firefox".into(), t, b)
1027    }
1028
1029    #[test]
1030    fn a_failed_thaw_replaces_the_notification_instead_of_closing() {
1031        assert_eq!(
1032            failed_resume_text("Firefox"),
1033            (
1034                "Firefox could not be resumed".to_string(),
1035                "It may stay paused. Run rlm guard status for details.".to_string()
1036            )
1037        );
1038        let mut n = notifier(true, false);
1039        n.record(&freeze("/a"), ok());
1040        n.end_tick(0, Memory::Low, &mut names);
1041        n.take();
1042        n.record(&Action::Thaw { res: res("/a") }, None);
1043        n.end_tick(5_000, Memory::Fine, &mut names);
1044        assert_eq!(n.take(), vec![not_resumed()]);
1045        n.end_tick(6_000, Memory::Fine, &mut names);
1046        assert!(n.take().is_empty(), "stays up");
1047        n.close_all();
1048        assert_eq!(n.take(), vec![Call::Close("firefox".into())]);
1049    }
1050
1051    #[test]
1052    fn a_failed_thaw_notice_closes_once_its_cgroup_is_gone_or_thawed() {
1053        let frozen = std::rc::Rc::new(std::cell::Cell::new(true));
1054        let mut n = notifier(true, false);
1055        let f = frozen.clone();
1056        n.still_frozen = Box::new(move |cg| cg != "/a" || f.get());
1057        n.record(&freeze("/a"), ok());
1058        n.end_tick(0, Memory::Low, &mut names);
1059        n.record(&Action::Thaw { res: res("/a") }, None);
1060        n.end_tick(5_000, Memory::Fine, &mut names);
1061        assert_eq!(n.take(), vec![paused(), not_resumed()]);
1062        n.end_tick(5_500, Memory::Fine, &mut names);
1063        assert!(n.take().is_empty(), "stays up while still frozen");
1064        frozen.set(false);
1065        n.end_tick(6_000, Memory::Fine, &mut names);
1066        assert_eq!(n.take(), vec![Call::Close("firefox".into())]);
1067        n.end_tick(7_000, Memory::Fine, &mut names);
1068        assert!(n.take().is_empty());
1069    }
1070
1071    #[test]
1072    fn the_default_check_reads_the_frozen_state() {
1073        let n = Notifier::new(Fake::default(), &GuardConfig::default());
1074        assert!(
1075            !(n.still_frozen)("/rlm-test-no-such-cgroup-for-notify"),
1076            "a missing cgroup is not frozen"
1077        );
1078    }
1079
1080    #[test]
1081    fn a_requested_or_finished_freeze_counts_as_frozen() {
1082        // A temp directory with the two files stands in for the cgroup.
1083        let dir = tempfile::tempdir().unwrap();
1084        let set = |freeze: &str, frozen: &str| {
1085            std::fs::write(dir.path().join("cgroup.freeze"), freeze).unwrap();
1086            std::fs::write(
1087                dir.path().join("cgroup.events"),
1088                format!("populated 1\nfrozen {frozen}\n"),
1089            )
1090            .unwrap();
1091        };
1092        set("1\n", "0");
1093        assert!(
1094            frozen_at(dir.path()),
1095            "freeze requested, tasks still stopping"
1096        );
1097        set("0\n", "1");
1098        assert!(frozen_at(dir.path()), "frozen, the request since dropped");
1099        set("1\n", "1");
1100        assert!(frozen_at(dir.path()));
1101        set("0\n", "0");
1102        assert!(!frozen_at(dir.path()), "thawed");
1103        set("x", "0");
1104        assert!(!frozen_at(dir.path()), "an unreadable request is no freeze");
1105        std::fs::remove_file(dir.path().join("cgroup.freeze")).unwrap();
1106        std::fs::write(dir.path().join("cgroup.events"), "frozen 1\n").unwrap();
1107        assert!(frozen_at(dir.path()), "events alone still say frozen");
1108    }
1109
1110    #[test]
1111    fn a_later_successful_release_closes_the_failed_thaw_notice() {
1112        let mut n = notifier(true, false);
1113        n.record(&freeze("/a"), ok());
1114        n.end_tick(0, Memory::Low, &mut names);
1115        n.record(&Action::Thaw { res: res("/a") }, None);
1116        n.end_tick(5_000, Memory::Low, &mut names);
1117        n.take();
1118        n.record(&freeze("/a"), ok());
1119        n.end_tick(6_000, Memory::Low, &mut names);
1120        assert_eq!(n.take(), vec![paused()]);
1121        n.record(&Action::Thaw { res: res("/a") }, ok());
1122        n.end_tick(11_000, Memory::Fine, &mut names);
1123        n.end_tick(12_000, Memory::Fine, &mut names);
1124        assert_eq!(n.take(), vec![Call::Close("firefox".into())]);
1125    }
1126
1127    #[test]
1128    fn failed_interventions_send_nothing() {
1129        let mut n = notifier(true, false);
1130        n.record(&freeze("/a"), None);
1131        n.record(&cap("/b"), None);
1132        n.end_tick(0, Memory::Low, &mut names);
1133        assert!(n.take().is_empty());
1134    }
1135
1136    #[test]
1137    fn caps_of_one_app_add_up() {
1138        let mut n = notifier(true, false);
1139        n.record(&cap("/a"), capped(GB));
1140        n.record(&cap("/b"), capped(2 * GB));
1141        n.end_tick(0, Memory::Low, &mut names);
1142        assert_eq!(n.take(), vec![slowed(3 * GB)]);
1143        n.record(&Action::LiftCap { res: res("/a") }, ok());
1144        n.end_tick(1_000, Memory::Low, &mut names);
1145        assert_eq!(n.take(), vec![slowed(2 * GB)], "updated, not closed");
1146    }
1147
1148    #[test]
1149    fn disabled_notify_sends_nothing() {
1150        let mut n = notifier(false, true);
1151        n.record(&freeze("/a"), ok());
1152        n.end_tick(0, Memory::Low, &mut names);
1153        n.record(&Action::Thaw { res: res("/a") }, ok());
1154        n.end_tick(5_000, Memory::Low, &mut names);
1155        n.end_tick(6_000, Memory::Low, &mut names);
1156        n.close_all();
1157        assert!(n.take().is_empty());
1158    }
1159
1160    #[test]
1161    fn turning_notify_off_clears_what_is_shown() {
1162        let mut n = notifier(true, false);
1163        n.record(&freeze("/a"), ok());
1164        n.end_tick(0, Memory::Low, &mut names);
1165        n.take();
1166        n.set_flags(NotifyFlags {
1167            notify: false,
1168            pressure: false,
1169        });
1170        assert_eq!(n.take(), vec![Call::Close("firefox".into())]);
1171        n.end_tick(1_000, Memory::Low, &mut names);
1172        assert!(n.take().is_empty());
1173    }
1174
1175    #[test]
1176    fn close_all_closes_every_app() {
1177        let mut n = notifier(true, false);
1178        n.record(&freeze("/a"), ok());
1179        n.end_tick(0, Memory::Low, &mut names);
1180        n.take();
1181        n.close_all();
1182        assert_eq!(n.take(), vec![Call::Close("firefox".into())]);
1183    }
1184
1185    fn warning() -> Call {
1186        Call::Show(
1187            PRESSURE_KEY.into(),
1188            PRESSURE_TITLE.into(),
1189            PRESSURE_BODY.into(),
1190        )
1191    }
1192
1193    #[test]
1194    fn pressure_warning_is_off_by_default() {
1195        let mut n = Notifier::new(Fake::default(), &GuardConfig::default());
1196        n.end_tick(0, Memory::Low, &mut names);
1197        assert!(n.take().is_empty());
1198    }
1199
1200    #[test]
1201    fn pressure_warning_is_rate_limited_and_only_at_high() {
1202        let mut n = notifier(true, true);
1203        n.end_tick(0, Memory::Fine, &mut names);
1204        assert!(n.take().is_empty(), "Warn is not enough");
1205        n.end_tick(1_000, Memory::Low, &mut names);
1206        assert_eq!(n.take(), vec![warning()]);
1207        n.end_tick(30_000, Memory::Low, &mut names);
1208        assert!(n.take().is_empty(), "at most once a minute");
1209        n.end_tick(61_000, Memory::Low, &mut names);
1210        assert_eq!(n.take(), vec![warning()]);
1211    }
1212
1213    #[test]
1214    fn memory_is_low_only_when_pressure_is_high_and_memory_short() {
1215        let t = common::GuardTrigger::default();
1216        let sample = |avail| Sample {
1217            some_avg10: 40.0,
1218            full_avg10: 0.0,
1219            mem_available_mb: avail,
1220            mem_total_mb: 16_000,
1221            source: super::super::types::PsiSource::AppSlice,
1222        };
1223        assert_eq!(memory_state(Level::High, &sample(1_000), &t), Memory::Low);
1224        assert_eq!(
1225            memory_state(Level::High, &sample(8_000), &t),
1226            Memory::Fine,
1227            "a stall with half the RAM free is not low memory"
1228        );
1229        assert_eq!(memory_state(Level::Warn, &sample(1_000), &t), Memory::Fine);
1230        assert_eq!(memory_state(Level::Critical, &sample(300), &t), Memory::Low);
1231    }
1232
1233    #[test]
1234    fn pressure_warning_gives_way_to_an_intervention() {
1235        let mut n = notifier(true, true);
1236        n.end_tick(0, Memory::Low, &mut names);
1237        n.take();
1238        n.record(&freeze("/a"), ok());
1239        n.end_tick(1_000, Memory::Low, &mut names);
1240        assert_eq!(n.take(), vec![paused(), Call::Close(PRESSURE_KEY.into())]);
1241        n.end_tick(70_000, Memory::Low, &mut names);
1242        assert!(n.take().is_empty(), "no warning while an app is held");
1243    }
1244
1245    #[test]
1246    fn pressure_warning_needs_notify() {
1247        let mut n = notifier(false, true);
1248        n.end_tick(0, Memory::Low, &mut names);
1249        assert!(n.take().is_empty());
1250    }
1251
1252    #[test]
1253    fn reload_takes_only_valid_notification_flags() {
1254        let current = NotifyFlags {
1255            notify: true,
1256            pressure: false,
1257        };
1258        let mut changed = GuardConfig {
1259            notify_pressure: true,
1260            enabled: false,
1261            ..GuardConfig::default()
1262        };
1263        changed.timing.freeze_hold_secs = 30;
1264        assert_eq!(
1265            flags_after_reload(current, Ok(&changed)),
1266            NotifyFlags {
1267                notify: true,
1268                pressure: true
1269            }
1270        );
1271        let err = common::Error::Config("bad".into());
1272        assert_eq!(flags_after_reload(current, Err(&err)), current);
1273    }
1274
1275    use std::cell::{Cell, RefCell};
1276    use std::rc::Rc;
1277
1278    /// What the fake bus and connector saw, shared with the test.
1279    #[derive(Default)]
1280    struct Log {
1281        calls: RefCell<Vec<String>>,
1282        /// Answers for the next connect attempts; `true` connects. Empty
1283        /// means connect.
1284        connects: RefCell<Vec<bool>>,
1285        /// When set, the next bus call fails with this.
1286        fail_next: RefCell<Option<BusError>>,
1287        next_id: Cell<u32>,
1288    }
1289
1290    struct FakeBus(Rc<Log>, u32);
1291
1292    impl Bus for FakeBus {
1293        fn notify(&self, replaces_id: u32, title: &str, _body: &str) -> Result<u32, BusError> {
1294            if let Some(e) = self.0.fail_next.borrow_mut().take() {
1295                return Err(e);
1296            }
1297            let id = if replaces_id == 0 {
1298                self.0.next_id.set(self.0.next_id.get() + 1);
1299                self.0.next_id.get()
1300            } else {
1301                replaces_id
1302            };
1303            self.0
1304                .calls
1305                .borrow_mut()
1306                .push(format!("conn{} notify {replaces_id}->{id} {title}", self.1));
1307            Ok(id)
1308        }
1309
1310        fn close(&self, id: u32) -> Result<(), BusError> {
1311            if let Some(e) = self.0.fail_next.borrow_mut().take() {
1312                return Err(e);
1313            }
1314            self.0
1315                .calls
1316                .borrow_mut()
1317                .push(format!("conn{} close {id}", self.1));
1318            Ok(())
1319        }
1320    }
1321
1322    struct FakeConnector(Rc<Log>, u32);
1323
1324    impl Connector for FakeConnector {
1325        type Bus = FakeBus;
1326
1327        fn connect(&mut self) -> Option<FakeBus> {
1328            let ok = {
1329                let mut answers = self.0.connects.borrow_mut();
1330                if answers.is_empty() {
1331                    true
1332                } else {
1333                    answers.remove(0)
1334                }
1335            };
1336            self.0
1337                .calls
1338                .borrow_mut()
1339                .push(format!("connect {}", if ok { "ok" } else { "failed" }));
1340            ok.then(|| {
1341                self.1 += 1;
1342                FakeBus(Rc::clone(&self.0), self.1)
1343            })
1344        }
1345
1346        fn fallback(&mut self, title: &str, _body: &str) {
1347            self.0
1348                .calls
1349                .borrow_mut()
1350                .push(format!("notify-send {title}"));
1351        }
1352    }
1353
1354    fn show(key: &str, title: &str) -> Cmd {
1355        Cmd::Show {
1356            key: key.into(),
1357            title: title.into(),
1358            body: String::new(),
1359        }
1360    }
1361
1362    fn sender_with(log: &Rc<Log>) -> Sender<FakeConnector> {
1363        Sender::new(FakeConnector(Rc::clone(log), 0))
1364    }
1365
1366    fn take(log: &Log) -> Vec<String> {
1367        std::mem::take(&mut *log.calls.borrow_mut())
1368    }
1369
1370    #[test]
1371    fn a_failed_connect_is_retried_after_the_backoff() {
1372        let log = Rc::new(Log::default());
1373        log.connects.borrow_mut().push(false);
1374        let mut s = sender_with(&log);
1375        let t0 = Instant::now();
1376        s.handle(show("a", "A paused"), t0);
1377        s.handle(show("a", "A slowed"), t0 + Duration::from_secs(10));
1378        assert_eq!(
1379            take(&log),
1380            [
1381                "connect failed",
1382                "notify-send A paused",
1383                "notify-send A slowed"
1384            ],
1385            "no new attempt during the backoff"
1386        );
1387        s.handle(show("a", "A slowed"), t0 + RECONNECT_AFTER);
1388        s.handle(show("a", "A slowed again"), t0 + RECONNECT_AFTER);
1389        s.handle(
1390            Cmd::Close { key: "a".into() },
1391            t0 + RECONNECT_AFTER + Duration::from_secs(1),
1392        );
1393        assert_eq!(
1394            take(&log),
1395            [
1396                "connect ok",
1397                "conn1 notify 0->1 A slowed",
1398                "conn1 notify 1->1 A slowed again",
1399                "conn1 close 1"
1400            ]
1401        );
1402    }
1403
1404    #[test]
1405    fn a_broken_connection_is_dropped_and_rebuilt() {
1406        let log = Rc::new(Log::default());
1407        let mut s = sender_with(&log);
1408        let t0 = Instant::now();
1409        s.handle(show("a", "A paused"), t0);
1410        take(&log);
1411        *log.fail_next.borrow_mut() = Some(BusError::Gone("broken pipe".into()));
1412        s.handle(show("a", "A slowed"), t0);
1413        assert_eq!(
1414            take(&log),
1415            ["connect ok", "conn2 notify 0->2 A slowed"],
1416            "reconnects and shows it as a new notification"
1417        );
1418        *log.fail_next.borrow_mut() = Some(BusError::Gone("gone".into()));
1419        s.handle(Cmd::Close { key: "a".into() }, t0);
1420        s.handle(show("b", "B paused"), t0);
1421        assert_eq!(
1422            take(&log),
1423            ["connect ok", "conn3 notify 0->3 B paused"],
1424            "a failed close also resets"
1425        );
1426    }
1427
1428    #[test]
1429    fn a_timed_out_call_keeps_the_connection() {
1430        let log = Rc::new(Log::default());
1431        let mut s = sender_with(&log);
1432        let t0 = Instant::now();
1433        s.handle(show("a", "A paused"), t0);
1434        *log.fail_next.borrow_mut() = Some(BusError::Call("timed out".into()));
1435        s.handle(show("a", "A slowed"), t0);
1436        s.handle(show("a", "A slowed"), t0);
1437        assert_eq!(
1438            take(&log),
1439            [
1440                "connect ok",
1441                "conn1 notify 0->1 A paused",
1442                "conn1 notify 1->1 A slowed"
1443            ]
1444        );
1445    }
1446
1447    #[test]
1448    fn io_errors_other_than_timeouts_mean_the_connection_is_gone() {
1449        let io = |kind| zbus::Error::InputOutput(std::sync::Arc::new(std::io::Error::from(kind)));
1450        assert!(matches!(
1451            classify(io(std::io::ErrorKind::BrokenPipe)),
1452            BusError::Gone(_)
1453        ));
1454        assert!(matches!(
1455            classify(io(std::io::ErrorKind::TimedOut)),
1456            BusError::Call(_)
1457        ));
1458        assert!(matches!(
1459            classify(zbus::Error::InvalidReply),
1460            BusError::Call(_)
1461        ));
1462    }
1463
1464    #[test]
1465    fn a_slow_connect_gives_up() {
1466        let start = std::time::Instant::now();
1467        let got = within(Duration::from_millis(50), || {
1468            std::thread::sleep(Duration::from_secs(5));
1469            Some(1)
1470        });
1471        assert_eq!(got, None);
1472        assert!(start.elapsed() < Duration::from_secs(2));
1473        assert_eq!(within(Duration::from_secs(2), || Some(7)), Some(7));
1474    }
1475
1476    #[test]
1477    fn notifications_use_the_server_default_lifetime() {
1478        assert_eq!(EXPIRE_TIMEOUT, -1);
1479    }
1480}