git.lucas.co / cce-status-interface
status bar
git clone https://git.lucas.co/cce-status-interface.git

src/listeners.rs (4.8K)

 1 //! Compositor push-update listeners: the status-feed subscriptions and the
 2 //! switcher trigger socket.
 3 
 4 use crate::CustomEvent;
 5 
 6 /// Subscribe to one compositor status topic and forward its pushes as
 7 /// [`CustomEvent`]s, reconnecting every second until the socket is there.
 8 ///
 9 /// `sub` is the whole subscription line; the topic is its first word, so a
10 /// topic that takes an argument can still be subscribed to.
11 pub(crate) async fn spawn_status_listener(sub: String, sender: calloop::channel::Sender<CustomEvent>) {
12     use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
13     use tokio::net::UnixStream;
14     // Retry delay, doubling while connections keep ending without ever
15     // delivering a line and reset the moment one does. A compositor that
16     // does not know this topic drops the subscription on sight, so the flat
17     // 1s retry turned an unrecognized topic into a permanent once-a-second
18     // reconnect from every module process — which is exactly what a bar
19     // running ahead of its compositor does with a new topic (the two halves
20     // deploy separately, and cce-fx only restarts at login).
21     let mut retry_s = 1u64;
22     loop {
23         let socket_path = match std::env::var("WAYLAND_DISPLAY") {
24             Ok(display) => {
25                 let primary = format!("/tmp/cce-status-interface-{}.sock", display);
26                 if std::path::Path::new(&primary).exists() {
27                     primary
28                 } else {
29                     format!("/tmp/cce-status-{}.sock", display)
30                 }
31             }
32             Err(_) => {
33                 let primary = "/tmp/cce-status-interface.sock".to_string();
34                 if std::path::Path::new(&primary).exists() {
35                     primary
36                 } else {
37                     "/tmp/cce-status.sock".to_string()
38                 }
39             }
40         };
41         if let Ok(mut stream) = UnixStream::connect(&socket_path).await {
42             log::info!("[status-listener] connected to {} for sub '{}'", socket_path, sub);
43             if stream.write_all(format!("{}\n", sub).as_bytes()).await.is_ok() {
44                 let mut reader = BufReader::new(stream);
45                 let mut line = String::new();
46                 while reader.read_line(&mut line).await.unwrap_or(0) > 0 {
47                     // Anything at all means the topic is understood.
48                     retry_s = 1;
49                     let val = line.trim().to_string();
50                     log::debug!("[status-listener] received '{}' update: '{}'", sub, val);
51                     if !val.is_empty() {
52                         let topic = sub.split_whitespace().next().unwrap_or("");
53                         let ev = match topic {
54                             "layout" => CustomEvent::LayoutUpdated(val.clone()),
55                             "title" => CustomEvent::TitleUpdated(val.clone()),
56                             // Click-away-close: the payload is the app_id of
57                             // the segment the press landed on ("-" for none).
58                             "dismiss" => CustomEvent::MenuDismiss(val.clone()),
59                             _ => unreachable!(),
60                         };
61                         let _ = sender.send(ev);
62                     }
63                     line.clear();
64                 }
65             }
66         }
67         tokio::time::sleep(std::time::Duration::from_secs(retry_s)).await;
68         retry_s = (retry_s * 2).min(30);
69     }
70 }
71 
72 pub(crate) async fn spawn_switcher_listener(sender: calloop::channel::Sender<CustomEvent>) {
73     use tokio::io::AsyncBufReadExt;
74     use tokio::net::UnixListener;
75     let display = std::env::var("WAYLAND_DISPLAY").unwrap_or_else(|_| "wayland-0".to_string());
76     let socket_path = format!("/tmp/cce-status-interface-switcher-{}.sock", display);
77     let _ = std::fs::remove_file(&socket_path);
78 
79     if let Ok(listener) = UnixListener::bind(&socket_path) {
80         log::info!("[switcher-listener] Listening on {}", socket_path);
81         loop {
82             if let Ok((stream, _)) = listener.accept().await {
83                 // Bounded in size and in TOTAL time: connections are served
84                 // one at a time here, so one that never finished its line
85                 // (until 2026-10-02 nothing timed it out) left the Super+Tab
86                 // switcher dead for the rest of the session.
87                 use tokio::io::AsyncReadExt;
88                 let mut reader = tokio::io::BufReader::new(stream.take(4096));
89                 let mut line = String::new();
90                 let read = tokio::time::timeout(std::time::Duration::from_secs(2), reader.read_line(&mut line)).await;
91                 if matches!(read, Ok(Ok(n)) if n > 0) {
92                     let _ = sender.send(CustomEvent::SwitcherTriggered);
93                 }
94             }
95         }
96     } else {
97         log::warn!("[switcher-listener] Failed to bind to {}", socket_path);
98     }
99 }