Wayland compositor (wlroots)
git clone https://git.lucas.co/cce-compositor.git
status socket: queue each subscriber's lines instead of resending them
Subscriber sockets are non-blocking. A write_all that hit WouldBlock
part-way through a line was treated as skip-for-now: the part it had
written stayed on the wire, and the whole line went out again on the next
push. A reader that fell behind (the status bar busy, a stalled grid)
therefore got spliced or doubled lines. On the one-shot topics
(shortcuts, selection, dismiss) a garbled line is a wrong action, and a
state topic's last_line could record a line that never fully arrived.
Each client now has a send queue. A line is appended whole and the
socket takes what it can, poll() watches for POLLOUT while anything is
owed, and the rest goes out in order. A client more than 256 KB behind
is dropped, as one that stopped reading always was eventually.
Tests: 240 KB of lines to a reader that catches up a chunk at a time
arrive whole and in order, and a client that never reads is dropped at
the cap. In a shadow, a title subscriber gets the update when a window
opens, the same as on the installed build.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
src/server/status_server.rs | 190 +++++++++++++++++++++++++++++++-------------
1 file changed, 136 insertions(+), 54 deletions(-)
diff --git a/src/server/status_server.rs b/src/server/status_server.rs
index f5b411a5..d8ff05e9 100644
--- a/src/server/status_server.rs
+++ b/src/server/status_server.rs
@@ -110,6 +110,56 @@ struct Client {
/// zoom) used to re-broadcast identical layout/title lines at frame
/// rate, and every subscriber rebuilt its segment per frame.
last_line: Option<String>,
+ /// Bytes owed to this client that its socket has not taken yet, oldest
+ /// first. Client sockets are non-blocking, and a `write_all` that hit
+ /// WouldBlock part-way through a line used to be skipped: the part it
+ /// had written stayed on the wire and the whole line was sent again
+ /// later, so a slow reader got spliced or doubled lines — on
+ /// `shortcuts`, `selection` and `dismiss` a garbled line is a wrong
+ /// action. Now every line is queued here whole and sent in order.
+ pending: Vec<u8>,
+}
+
+/// Most bytes a subscriber may fall behind before it is dropped. A status
+/// line is tens of bytes; a client this far behind is not reading at all.
+const MAX_PENDING: usize = 256 * 1024;
+
+impl Client {
+ fn new(subscription: Subscription, stream: UnixStream) -> Self {
+ Self { subscription, stream, last_line: None, pending: Vec::new() }
+ }
+
+ /// Queue `line` (newline added) behind anything still owed, and send
+ /// what the socket takes now. Err: drop this client.
+ fn send_line(&mut self, line: &str) -> Result<(), ()> {
+ self.pending.extend_from_slice(line.as_bytes());
+ self.pending.push(b'\n');
+ if self.pending.len() > MAX_PENDING {
+ log::warn!(
+ "[status] client {:?} is {} bytes behind; dropping it",
+ self.subscription,
+ self.pending.len()
+ );
+ return Err(());
+ }
+ self.flush()
+ }
+
+ /// Send as much of `pending` as the socket takes without blocking.
+ fn flush(&mut self) -> Result<(), ()> {
+ while !self.pending.is_empty() {
+ match self.stream.write(&self.pending) {
+ Ok(0) => return Err(()),
+ Ok(n) => {
+ self.pending.drain(..n);
+ }
+ Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => return Ok(()),
+ Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => {}
+ Err(_) => return Err(()),
+ }
+ }
+ Ok(())
+ }
}
/// Handle to the status server for sending updates from the main loop.
@@ -178,17 +228,9 @@ fn push_one_shot(clients: &mut Vec<Client>, topic: &Subscription, lines: &[Strin
continue;
}
for line in lines {
- match client
- .stream
- .write_all(line.as_bytes())
- .and_then(|_| client.stream.write_all(b"\n"))
- {
- Ok(_) => {}
- Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {}
- Err(_) => {
- dead_clients.push(i);
- break;
- }
+ if client.send_line(line).is_err() {
+ dead_clients.push(i);
+ break;
}
}
}
@@ -240,7 +282,10 @@ fn wait_for_activity(wake: &OwnedFd, listener: &UnixListener, clients: &[Client]
fds.push(libc::pollfd { fd, events: libc::POLLIN, revents: 0 });
}
for client in clients {
- fds.push(libc::pollfd { fd: client.stream.as_raw_fd(), events: libc::POLLIN, revents: 0 });
+ // POLLOUT only while something is owed, or poll() would return at
+ // once forever on a socket that is simply writable.
+ let events = if client.pending.is_empty() { libc::POLLIN } else { libc::POLLIN | libc::POLLOUT };
+ fds.push(libc::pollfd { fd: client.stream.as_raw_fd(), events, revents: 0 });
}
let n = unsafe { libc::poll(fds.as_mut_ptr(), fds.len() as libc::nfds_t, -1) };
if n < 0 {
@@ -310,12 +355,7 @@ fn status_server_main(rx: mpsc::Receiver<StatusMsg>, wake: Arc<OwnedFd>, display
}
if sub != Subscription::Unknown {
log::info!("[status] new subscriber for {:?}", sub);
- let client = Client {
- subscription: sub,
- stream,
- last_line: None,
- };
- clients.push(client);
+ clients.push(Client::new(sub, stream));
has_new_update = true; // push the latest status to the new client
}
}
@@ -374,6 +414,19 @@ fn status_server_main(rx: mpsc::Receiver<StatusMsg>, wake: Arc<OwnedFd>, display
}
}
+ // Send what slow clients are still owed (poll() woke on POLLOUT).
+ {
+ let mut dead_clients = Vec::new();
+ for (i, client) in clients.iter_mut().enumerate() {
+ if !client.pending.is_empty() && client.flush().is_err() {
+ dead_clients.push(i);
+ }
+ }
+ for i in dead_clients.into_iter().rev() {
+ clients.remove(i);
+ }
+ }
+
// Process incoming updates from the main loop
let mut dismiss_events: Vec<String> = Vec::new();
let mut shortcut_events: Vec<String> = Vec::new();
@@ -415,17 +468,9 @@ fn status_server_main(rx: mpsc::Receiver<StatusMsg>, wake: Arc<OwnedFd>, display
continue;
}
for except in &dismiss_events {
- match client
- .stream
- .write_all(except.as_bytes())
- .and_then(|_| client.stream.write_all(b"\n"))
- {
- Ok(_) => {}
- Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {}
- Err(_) => {
- dead_clients.push(i);
- break;
- }
+ if client.send_line(except).is_err() {
+ dead_clients.push(i);
+ break;
}
}
}
@@ -464,31 +509,13 @@ fn status_server_main(rx: mpsc::Receiver<StatusMsg>, wake: Arc<OwnedFd>, display
if client.last_line.as_deref() == Some(msg.as_str()) {
continue;
}
- match client
- .stream
- .write_all(msg.as_bytes())
- .and_then(|_| client.stream.write_all(b"\n"))
- {
- Ok(_) => {
- client.last_line = Some(msg);
- }
- Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
- // Client not ready to receive — skip for now
- }
- Err(ref e) if e.kind() == std::io::ErrorKind::BrokenPipe => {
- log::info!(
- "[status] client {:?} disconnected (broken pipe)",
- client.subscription
- );
- dead_clients.push(i);
- }
- Err(e) => {
- log::error!(
- "[status] write error to client {:?}: {}",
- client.subscription, e
- );
- dead_clients.push(i);
- }
+ // Queued means delivered: what a slow client has not
+ // taken yet goes out, in order, as its socket drains.
+ if client.send_line(&msg).is_ok() {
+ client.last_line = Some(msg);
+ } else {
+ log::info!("[status] client {:?} dropped (write failed)", client.subscription);
+ dead_clients.push(i);
}
}
@@ -606,3 +633,58 @@ pub unsafe fn build_status_update(wm: &crate::window_manager::WindowManager) ->
adjust_text: if wm.window_adjust_active() { "on" } else { "off" }.to_string(),
}
}
+
+#[cfg(test)]
+mod send_queue_tests {
+ use super::*;
+ use std::io::Read;
+
+ #[test]
+ fn a_slow_reader_gets_every_line_whole_and_in_order() {
+ let (server, mut reader) = UnixStream::pair().unwrap();
+ server.set_nonblocking(true).unwrap();
+ let mut client = Client::new(Subscription::Shortcuts, server);
+ // ~240 KB: more than a Unix socket buffers, less than MAX_PENDING.
+ let lines: Vec<String> = (0..2400).map(|i| format!("activated /s/1 id{i:05} {}", "x".repeat(70))).collect();
+ for line in &lines {
+ client.send_line(line).unwrap();
+ }
+ assert!(!client.pending.is_empty(), "the socket must actually have backed up");
+
+ // The reader catches up a chunk at a time; the server loop flushes
+ // between chunks, as poll()'s POLLOUT has it do.
+ let expected: usize = lines.iter().map(|l| l.len() + 1).sum();
+ let mut got = Vec::new();
+ let mut buf = [0u8; 16 * 1024];
+ reader.set_read_timeout(Some(std::time::Duration::from_millis(200))).unwrap();
+ while got.len() < expected {
+ let n = reader.read(&mut buf).unwrap();
+ assert!(n > 0);
+ got.extend_from_slice(&buf[..n]);
+ client.flush().unwrap();
+ }
+ let text = String::from_utf8(got).unwrap();
+ let received: Vec<&str> = text.lines().collect();
+ assert_eq!(received.len(), lines.len());
+ for (want, have) in lines.iter().zip(&received) {
+ assert_eq!(want, have, "a line arrived spliced or out of order");
+ }
+ }
+
+ #[test]
+ fn a_client_that_never_reads_is_dropped_at_the_cap() {
+ let (server, _reader) = UnixStream::pair().unwrap();
+ server.set_nonblocking(true).unwrap();
+ let mut client = Client::new(Subscription::Layout, server);
+ let line = "y".repeat(1000);
+ let mut dropped = false;
+ for _ in 0..2000 {
+ if client.send_line(&line).is_err() {
+ dropped = true;
+ break;
+ }
+ }
+ assert!(dropped, "2 MB to a reader that never reads must hit MAX_PENDING");
+ assert!(client.pending.len() <= MAX_PENDING + line.len() + 1);
+ }
+}