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 }