Wayland compositor (wlroots)
git clone https://git.lucas.co/cce-compositor.git
src/server/status_server.rs (27.4K)
1 // Status socket server for monolithic cce server
2 //
3 // Runs in a dedicated thread. cce-status connects to
4 // /tmp/cce-status-{WAYLAND_DISPLAY}.sock, sends a subscription line
5 // ("layout", "title", "modifiers", "adjust", "dismiss", "shortcuts",
6 // "clickaway" or "selection") and receives lines whenever the status changes.
7 //
8 // The main loop sends updates through an mpsc channel. The server thread
9 // owns the socket and handles all I/O independently of the Wayland event loop.
10
11 use std::io::Write;
12 use std::os::fd::{AsRawFd, OwnedFd};
13 use std::os::unix::net::{UnixListener, UnixStream};
14 use std::sync::mpsc;
15 use std::sync::Arc;
16
17 use crate::ipc_server::{drain_wake_fd, new_wake_fd, wake_fd};
18
19 /// A status update sent from the main loop to the server thread.
20 #[derive(Debug, Clone, PartialEq, Eq)]
21 pub struct StatusUpdate {
22 /// Plain text for layout module subscribers
23 pub layout_text: String,
24 /// Plain text for title module subscribers
25 pub title_text: String,
26 /// Plain text for modifiers subscriber
27 pub modifiers_text: String,
28 /// "on" while window-adjust mode is active (overview, or Super held —
29 /// `WindowManager::window_adjust_active`), else "off". The desktop grid
30 /// subscribes to show its own resize handles on the pinned images in
31 /// step with the windows' handles; it never holds keyboard focus, so
32 /// it cannot read the modifier state for itself.
33 pub adjust_text: String,
34 }
35
36 /// A message from the main loop to the server thread: either a new state
37 /// snapshot for the state topics, or a one-shot menu-dismiss event.
38 #[derive(Debug, Clone)]
39 pub enum StatusMsg {
40 State(StatusUpdate),
41 /// Click-away-close for in-surface status menus: every `dismiss`
42 /// subscriber EXCEPT the segment whose app_id is carried here should
43 /// close its open menu (the exempt segment saw the press itself).
44 MenuDismiss { except_app_id: String },
45 /// A portal-bound chord went down or up (`global_shortcuts`): one line,
46 /// `activated|deactivated <session> <id> <time_msec>`, for every
47 /// `shortcuts` subscriber — in practice the one portal backend.
48 Shortcut(String),
49 /// A button press that landed on no X11 surface while an X11
50 /// override-redirect window was showing: one line, `press`, for every
51 /// `clickaway` subscriber. An X11 popup only hears clicks on its own
52 /// client's X windows (Xwayland's pointer never reaches it over a
53 /// Wayland surface), so the tray bridge closes the popup it opened on
54 /// this cue — see `cce-status-interface`'s `cce-xembed-tray`.
55 ClickAway,
56 /// The overview drag-selection carrying the grid's desktop images: one
57 /// line for every `selection` subscriber (in practice the grid) —
58 /// `move <id>:<x>:<y> ...` with the images' new virtual positions as a
59 /// group move steps, then `drop` at its release. See
60 /// [`crate::selection`].
61 Selection(String),
62 }
63
64 /// Subscription types that the status bar script can request.
65 #[derive(Debug, Clone, PartialEq, Eq)]
66 enum Subscription {
67 Layout,
68 Title,
69 Modifiers,
70 /// `adjust` — "on"/"off" as window-adjust mode comes and goes.
71 Adjust,
72 /// One-shot menu-dismiss events only — never receives state pushes.
73 Dismiss,
74 /// `shortcuts` — one-shot portal shortcut press/release lines only
75 /// (see `StatusMsg::Shortcut`); never receives state pushes.
76 Shortcuts,
77 /// `clickaway` — one-shot `press` lines only (see
78 /// `StatusMsg::ClickAway`); never receives state pushes.
79 ClickAway,
80 /// `selection` — one-shot group-move lines for the desktop images
81 /// (see `StatusMsg::Selection`); never receives state pushes.
82 Selection,
83 Unknown,
84 }
85
86 impl Subscription {
87 fn from_str(s: &str) -> Self {
88 let s = s.trim();
89 match s {
90 "layout" => Subscription::Layout,
91 "title" => Subscription::Title,
92 "modifiers" => Subscription::Modifiers,
93 "adjust" => Subscription::Adjust,
94 "dismiss" => Subscription::Dismiss,
95 "shortcuts" => Subscription::Shortcuts,
96 "clickaway" => Subscription::ClickAway,
97 "selection" => Subscription::Selection,
98 _ => Subscription::Unknown,
99 }
100 }
101 }
102
103 /// A connected client with a known subscription.
104 struct Client {
105 subscription: Subscription,
106 stream: UnixStream,
107 /// The last line actually written to this client. A state push only
108 /// re-sends a topic whose formatted line CHANGED — a StatusUpdate is
109 /// one struct, so e.g. a camera animation (viewport text embeds pan/
110 /// zoom) used to re-broadcast identical layout/title lines at frame
111 /// rate, and every subscriber rebuilt its segment per frame.
112 last_line: Option<String>,
113 /// Bytes owed to this client that its socket has not taken yet, oldest
114 /// first. Client sockets are non-blocking, and a `write_all` that hit
115 /// WouldBlock part-way through a line used to be skipped: the part it
116 /// had written stayed on the wire and the whole line was sent again
117 /// later, so a slow reader got spliced or doubled lines — on
118 /// `shortcuts`, `selection` and `dismiss` a garbled line is a wrong
119 /// action. Now every line is queued here whole and sent in order.
120 pending: Vec<u8>,
121 }
122
123 /// Most bytes a subscriber may fall behind before it is dropped. A status
124 /// line is tens of bytes; a client this far behind is not reading at all.
125 const MAX_PENDING: usize = 256 * 1024;
126
127 impl Client {
128 fn new(subscription: Subscription, stream: UnixStream) -> Self {
129 Self { subscription, stream, last_line: None, pending: Vec::new() }
130 }
131
132 /// Queue `line` (newline added) behind anything still owed, and send
133 /// what the socket takes now. Err: drop this client.
134 fn send_line(&mut self, line: &str) -> Result<(), ()> {
135 self.pending.extend_from_slice(line.as_bytes());
136 self.pending.push(b'\n');
137 if self.pending.len() > MAX_PENDING {
138 log::warn!(
139 "[status] client {:?} is {} bytes behind; dropping it",
140 self.subscription,
141 self.pending.len()
142 );
143 return Err(());
144 }
145 self.flush()
146 }
147
148 /// Send as much of `pending` as the socket takes without blocking.
149 fn flush(&mut self) -> Result<(), ()> {
150 while !self.pending.is_empty() {
151 match self.stream.write(&self.pending) {
152 Ok(0) => return Err(()),
153 Ok(n) => {
154 self.pending.drain(..n);
155 }
156 Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => return Ok(()),
157 Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => {}
158 Err(_) => return Err(()),
159 }
160 }
161 Ok(())
162 }
163 }
164
165 /// Handle to the status server for sending updates from the main loop.
166 ///
167 /// Every send bumps `wake`, the eventfd the server thread `poll()`s on
168 /// alongside its sockets. The thread used to spin on `try_recv` with a 20 ms
169 /// sleep — 50 wakeups/s forever, subscribers or not; now it blocks until a
170 /// socket or the main loop has something for it.
171 #[derive(Debug, Clone)]
172 pub struct StatusSender {
173 tx: mpsc::Sender<StatusMsg>,
174 wake: Arc<OwnedFd>,
175 }
176
177 impl StatusSender {
178 pub fn send(&self, update: StatusUpdate) {
179 // If the channel is full or the receiver is gone, just drop it.
180 if self.tx.send(StatusMsg::State(update)).is_ok() {
181 wake_fd(&self.wake);
182 }
183 }
184
185 /// Fire a one-shot menu-dismiss at every `dismiss` subscriber except the
186 /// segment with this app_id (pass "-" to exempt nobody).
187 pub fn send_menu_dismiss(&self, except_app_id: &str) {
188 if self.tx.send(StatusMsg::MenuDismiss { except_app_id: except_app_id.to_string() }).is_ok() {
189 wake_fd(&self.wake);
190 }
191 }
192 }
193
194 impl StatusSender {
195 /// Report a portal shortcut edge to every `shortcuts` subscriber.
196 pub fn send_shortcut_event(&self, line: &str) {
197 if self.tx.send(StatusMsg::Shortcut(line.to_string())).is_ok() {
198 wake_fd(&self.wake);
199 }
200 }
201
202 /// Report a press that missed every X11 surface (`StatusMsg::ClickAway`).
203 pub fn send_click_away(&self) {
204 if self.tx.send(StatusMsg::ClickAway).is_ok() {
205 wake_fd(&self.wake);
206 }
207 }
208
209 /// Tell the grid where a group move put its images
210 /// (`StatusMsg::Selection`).
211 pub fn send_selection_line(&self, line: &str) {
212 if self.tx.send(StatusMsg::Selection(line.to_string())).is_ok() {
213 wake_fd(&self.wake);
214 }
215 }
216 }
217
218 /// Write each of `lines` to every client subscribed to `topic`, dropping
219 /// the clients whose socket has gone away. One-shot topics only: nothing is
220 /// remembered for a client that subscribes later.
221 fn push_one_shot(clients: &mut Vec<Client>, topic: &Subscription, lines: &[String]) {
222 if lines.is_empty() {
223 return;
224 }
225 let mut dead_clients = Vec::new();
226 for (i, client) in clients.iter_mut().enumerate() {
227 if &client.subscription != topic {
228 continue;
229 }
230 for line in lines {
231 if client.send_line(line).is_err() {
232 dead_clients.push(i);
233 break;
234 }
235 }
236 }
237 dead_clients.dedup();
238 for i in dead_clients.into_iter().rev() {
239 clients.remove(i);
240 }
241 }
242
243 impl Drop for StatusSender {
244 /// Dropping the last handle disconnects the channel; the thread only
245 /// notices when it next wakes, so give it one.
246 fn drop(&mut self) {
247 wake_fd(&self.wake);
248 }
249 }
250
251 pub fn get_status_socket_path(display_socket: Option<&str>) -> String {
252 if let Some(display) = display_socket {
253 format!("/tmp/cce-status-interface-{}.sock", display)
254 } else {
255 "/tmp/cce-status-interface.sock".to_string()
256 }
257 }
258
259 /// Spawn the status server thread. Returns a StatusSender for the main loop.
260 pub fn spawn_status_server(display_socket: Option<String>) -> StatusSender {
261 let (tx, rx) = mpsc::channel::<StatusMsg>();
262 let wake = new_wake_fd().expect("failed to create status wake eventfd");
263 let thread_wake = wake.clone();
264
265 std::thread::Builder::new()
266 .name("cce-status-server".into())
267 .spawn(move || {
268 status_server_main(rx, thread_wake, display_socket);
269 })
270 .expect("failed to spawn status server thread");
271
272 StatusSender { tx, wake }
273 }
274
275 /// Block until the wake eventfd, the listener, or any subscriber socket is
276 /// readable. Returns `(wake, accept, per-client readiness)`; a client is
277 /// "ready" on data, hangup or error alike, since all three are handled by
278 /// reading it.
279 fn wait_for_activity(wake: &OwnedFd, listener: &UnixListener, clients: &[Client]) -> Option<(bool, bool, Vec<bool>)> {
280 let mut fds: Vec<libc::pollfd> = Vec::with_capacity(2 + clients.len());
281 for fd in [wake.as_raw_fd(), listener.as_raw_fd()] {
282 fds.push(libc::pollfd { fd, events: libc::POLLIN, revents: 0 });
283 }
284 for client in clients {
285 // POLLOUT only while something is owed, or poll() would return at
286 // once forever on a socket that is simply writable.
287 let events = if client.pending.is_empty() { libc::POLLIN } else { libc::POLLIN | libc::POLLOUT };
288 fds.push(libc::pollfd { fd: client.stream.as_raw_fd(), events, revents: 0 });
289 }
290 let n = unsafe { libc::poll(fds.as_mut_ptr(), fds.len() as libc::nfds_t, -1) };
291 if n < 0 {
292 let err = std::io::Error::last_os_error();
293 if err.kind() == std::io::ErrorKind::Interrupted {
294 return Some((false, false, vec![false; clients.len()]));
295 }
296 log::error!("[status] poll failed: {}", err);
297 return None;
298 }
299 let ready = |f: &libc::pollfd| f.revents != 0;
300 Some((ready(&fds[0]), ready(&fds[1]), fds[2..].iter().map(ready).collect()))
301 }
302
303 fn status_server_main(rx: mpsc::Receiver<StatusMsg>, wake: Arc<OwnedFd>, display_socket: Option<String>) {
304 let socket_path = get_status_socket_path(display_socket.as_deref());
305 // Remove stale socket
306 let _ = std::fs::remove_file(&socket_path);
307
308 let listener = match UnixListener::bind(&socket_path) {
309 Ok(l) => l,
310 Err(e) => {
311 log::error!("[status] failed to bind {}: {}", socket_path, e);
312 return;
313 }
314 };
315
316 // Set non-blocking so accept() doesn't hang the thread
317 if let Err(e) = listener.set_nonblocking(true) {
318 log::error!("[status] failed to set non-blocking: {}", e);
319 return;
320 }
321
322 log::info!("[status] listening on {}", socket_path);
323
324 let mut clients: Vec<Client> = Vec::new();
325 let mut latest: Option<StatusUpdate> = None;
326
327 loop {
328 let Some((wake_ready, accept_ready, client_ready)) = wait_for_activity(&wake, &listener, &clients) else {
329 // poll() itself failing is not something a retry fixes fast;
330 // back off so the error line cannot flood the log.
331 std::thread::sleep(std::time::Duration::from_millis(100));
332 continue;
333 };
334 if wake_ready {
335 drain_wake_fd(wake.as_raw_fd());
336 }
337
338 let mut has_new_update = false;
339
340 // Accept new connections (the listener is non-blocking)
341 for _ in 0..5 {
342 if !accept_ready {
343 break;
344 }
345 match listener.accept() {
346 Ok((stream, _addr)) => {
347 // Read the subscription line while the socket is still
348 // blocking (bounded by a read timeout): poll() hands us
349 // the connection the instant it lands, which can be
350 // before the client's first line is in the buffer.
351 let sub = read_subscription(&stream);
352 if let Err(e) = stream.set_nonblocking(true) {
353 log::error!("[status] failed to set non-blocking on client: {}", e);
354 continue;
355 }
356 if sub != Subscription::Unknown {
357 log::info!("[status] new subscriber for {:?}", sub);
358 clients.push(Client::new(sub, stream));
359 has_new_update = true; // push the latest status to the new client
360 }
361 }
362 Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
363 break;
364 }
365 Err(e) => {
366 // EMFILE and friends: the socket stays readable, so
367 // without a pause this loop (and its log line) spins the
368 // thread at 100% — observed as a 167GB log once dead
369 // subscribers had exhausted the fd table.
370 log::error!("[status] accept error: {}", e);
371 std::thread::sleep(std::time::Duration::from_millis(100));
372 break;
373 }
374 }
375 }
376
377 // Reap dead subscribers by reading: a subscriber never sends after
378 // its subscription line, so a successful zero-byte read is EOF (the
379 // client vanished). Waiting for a WRITE to fail leaked them instead
380 // — last_line dedup means a quiet topic may never write again, and
381 // every bar restart stranded its whole subscriber set. Enough
382 // restarts exhausted the fd table and took the session down.
383 {
384 let mut buf = [0u8; 64];
385 let mut dead_clients = Vec::new();
386 for (i, client) in clients.iter_mut().enumerate() {
387 // Only sockets poll() flagged; the rest are quiet, not dead.
388 if !client_ready.get(i).copied().unwrap_or(false) {
389 continue;
390 }
391 loop {
392 use std::io::Read;
393 match client.stream.read(&mut buf) {
394 Ok(0) => {
395 log::info!(
396 "[status] client {:?} disconnected (eof)",
397 client.subscription
398 );
399 dead_clients.push(i);
400 break;
401 }
402 // Unexpected chatter: drain and keep the client.
403 Ok(_) => continue,
404 Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => break,
405 Err(_) => {
406 dead_clients.push(i);
407 break;
408 }
409 }
410 }
411 }
412 for i in dead_clients.into_iter().rev() {
413 clients.remove(i);
414 }
415 }
416
417 // Send what slow clients are still owed (poll() woke on POLLOUT).
418 {
419 let mut dead_clients = Vec::new();
420 for (i, client) in clients.iter_mut().enumerate() {
421 if !client.pending.is_empty() && client.flush().is_err() {
422 dead_clients.push(i);
423 }
424 }
425 for i in dead_clients.into_iter().rev() {
426 clients.remove(i);
427 }
428 }
429
430 // Process incoming updates from the main loop
431 let mut dismiss_events: Vec<String> = Vec::new();
432 let mut shortcut_events: Vec<String> = Vec::new();
433 let mut clickaway_events: Vec<String> = Vec::new();
434 let mut selection_events: Vec<String> = Vec::new();
435 loop {
436 match rx.try_recv() {
437 Ok(StatusMsg::State(update)) => {
438 latest = Some(update);
439 has_new_update = true;
440 }
441 Ok(StatusMsg::MenuDismiss { except_app_id }) => {
442 dismiss_events.push(except_app_id);
443 }
444 Ok(StatusMsg::Shortcut(line)) => {
445 shortcut_events.push(line);
446 }
447 Ok(StatusMsg::ClickAway) => {
448 clickaway_events.push("press".to_string());
449 }
450 Ok(StatusMsg::Selection(line)) => {
451 selection_events.push(line);
452 }
453 Err(mpsc::TryRecvError::Empty) => break,
454 Err(mpsc::TryRecvError::Disconnected) => {
455 log::info!("[status] channel disconnected, exiting");
456 let _ = std::fs::remove_file(&socket_path);
457 return;
458 }
459 }
460 }
461
462 // One-shot dismiss lines go only to `dismiss` subscribers; the line
463 // payload is the exempt app_id.
464 if !dismiss_events.is_empty() {
465 let mut dead_clients = Vec::new();
466 for (i, client) in clients.iter_mut().enumerate() {
467 if client.subscription != Subscription::Dismiss {
468 continue;
469 }
470 for except in &dismiss_events {
471 if client.send_line(except).is_err() {
472 dead_clients.push(i);
473 break;
474 }
475 }
476 }
477 dead_clients.dedup();
478 for i in dead_clients.into_iter().rev() {
479 clients.remove(i);
480 }
481 }
482
483 // Portal shortcut edges go only to `shortcuts` subscribers, in
484 // order; click-aways only to `clickaway` subscribers.
485 push_one_shot(&mut clients, &Subscription::Shortcuts, &shortcut_events);
486 push_one_shot(&mut clients, &Subscription::ClickAway, &clickaway_events);
487 push_one_shot(&mut clients, &Subscription::Selection, &selection_events);
488
489 // If we got a new update, push it to all clients
490 if has_new_update {
491 if let Some(ref update) = latest {
492 let mut dead_clients = Vec::new();
493
494 for (i, client) in clients.iter_mut().enumerate() {
495 // Dismiss and shortcuts subscribers get one-shot events
496 // only, never state pushes.
497 if matches!(
498 client.subscription,
499 Subscription::Dismiss
500 | Subscription::Shortcuts
501 | Subscription::ClickAway
502 | Subscription::Selection
503 ) {
504 continue;
505 }
506 let msg = format_for_subscription(&client.subscription, update);
507 // Only lines that changed for THIS topic go out (see
508 // Client::last_line); a fresh client always gets one.
509 if client.last_line.as_deref() == Some(msg.as_str()) {
510 continue;
511 }
512 // Queued means delivered: what a slow client has not
513 // taken yet goes out, in order, as its socket drains.
514 if client.send_line(&msg).is_ok() {
515 client.last_line = Some(msg);
516 } else {
517 log::info!("[status] client {:?} dropped (write failed)", client.subscription);
518 dead_clients.push(i);
519 }
520 }
521
522 // Remove dead clients (iterate in reverse to preserve indices)
523 for i in dead_clients.into_iter().rev() {
524 clients.remove(i);
525 }
526 }
527 }
528 }
529 }
530
531 fn read_subscription(stream: &UnixStream) -> Subscription {
532 // A subscriber writes its one line right after connecting, so this
533 // returns at once in practice. The bound is for one that does not: this
534 // thread pushes every status update, so the wait is capped in total time
535 // and size, not per read (`read_line_bounded`).
536 match crate::ipc_server::read_line_bounded(stream, 256, std::time::Duration::from_millis(200)) {
537 Some(line) => Subscription::from_str(&line),
538 None => {
539 log::error!("[status] no subscription line within 200ms / 256 bytes");
540 Subscription::Unknown
541 }
542 }
543 }
544
545 fn format_for_subscription(sub: &Subscription, update: &StatusUpdate) -> String {
546 match sub {
547 Subscription::Layout => update.layout_text.clone(),
548 Subscription::Title => update.title_text.clone(),
549 Subscription::Modifiers => update.modifiers_text.clone(),
550 Subscription::Adjust => update.adjust_text.clone(),
551 Subscription::Dismiss
552 | Subscription::Shortcuts
553 | Subscription::ClickAway
554 | Subscription::Selection
555 | Subscription::Unknown => String::new(),
556 }
557 }
558
559 pub unsafe fn build_status_update(wm: &crate::window_manager::WindowManager) -> StatusUpdate {
560 // `focused_window()` falls back to the most recent real window so the
561 // bar doesn't flash while overlay UI (the launcher) briefly holds
562 // focus. But an explicit Focus::None (desktop click) is a real,
563 // user-visible state — keystrokes go nowhere — and the status feed
564 // must report it honestly instead of showing the last window as if it
565 // still had focus.
566 let seat_focus_is_none = wm
567 .first_seat()
568 .map(|s| matches!((*s).focused, crate::seat::Focus::None))
569 .unwrap_or(false);
570 let focused_window = if seat_focus_is_none {
571 std::ptr::null_mut()
572 } else {
573 wm.focused_window()
574 };
575
576 let layout_text = if !focused_window.is_null() {
577 (*focused_window).tiling_mode.as_str().to_string()
578 } else {
579 let focused_layer = wm.focused_layer_surface();
580 let mut is_cce_cloud = false;
581 if !focused_layer.is_null() {
582 let wlr_layer_surface = crate::ffi::wlr_layer_surface_v1_try_from_wlr_surface(focused_layer);
583 if !wlr_layer_surface.is_null() && !(*wlr_layer_surface).namespace.is_null() {
584 let ns = std::ffi::CStr::from_ptr((*wlr_layer_surface).namespace).to_string_lossy();
585 if ns.starts_with("cce-cloud") {
586 is_cce_cloud = true;
587 }
588 }
589 }
590 if is_cce_cloud {
591 "Overlay".to_string()
592 } else {
593 "---".to_string()
594 }
595 };
596
597 let title_text = if !focused_window.is_null() {
598 let title_ptr = (*focused_window).get_title();
599 if !title_ptr.is_null() {
600 std::ffi::CStr::from_ptr(title_ptr).to_string_lossy().into_owned()
601 } else {
602 "(none)".to_string()
603 }
604 } else {
605 let focused_layer = wm.focused_layer_surface();
606 if !focused_layer.is_null() {
607 let wlr_layer_surface = crate::ffi::wlr_layer_surface_v1_try_from_wlr_surface(focused_layer);
608 if !wlr_layer_surface.is_null() && !(*wlr_layer_surface).namespace.is_null() {
609 std::ffi::CStr::from_ptr((*wlr_layer_surface).namespace).to_string_lossy().into_owned()
610 } else {
611 "(none)".to_string()
612 }
613 } else {
614 "(none)".to_string()
615 }
616 };
617
618 let seat_ptr = wm.first_seat().unwrap_or(std::ptr::null_mut());
619 let mut super_pressed = false;
620 if !seat_ptr.is_null() {
621 let wlr_keyboard = crate::ffi::river_wlr_seat_get_keyboard((*seat_ptr).wlr_seat);
622 if !wlr_keyboard.is_null() {
623 let modifiers = crate::ffi::wlr_keyboard_get_modifiers(wlr_keyboard);
624 super_pressed = modifiers & 0x40 != 0;
625 }
626 }
627 let modifiers_text = if super_pressed { "super" } else { "none" }.to_string();
628
629 StatusUpdate {
630 layout_text,
631 title_text,
632 modifiers_text,
633 adjust_text: if wm.window_adjust_active() { "on" } else { "off" }.to_string(),
634 }
635 }
636
637 #[cfg(test)]
638 mod send_queue_tests {
639 use super::*;
640 use std::io::Read;
641
642 #[test]
643 fn a_slow_reader_gets_every_line_whole_and_in_order() {
644 let (server, mut reader) = UnixStream::pair().unwrap();
645 server.set_nonblocking(true).unwrap();
646 let mut client = Client::new(Subscription::Shortcuts, server);
647 // ~240 KB: more than a Unix socket buffers, less than MAX_PENDING.
648 let lines: Vec<String> = (0..2400).map(|i| format!("activated /s/1 id{i:05} {}", "x".repeat(70))).collect();
649 for line in &lines {
650 client.send_line(line).unwrap();
651 }
652 assert!(!client.pending.is_empty(), "the socket must actually have backed up");
653
654 // The reader catches up a chunk at a time; the server loop flushes
655 // between chunks, as poll()'s POLLOUT has it do.
656 let expected: usize = lines.iter().map(|l| l.len() + 1).sum();
657 let mut got = Vec::new();
658 let mut buf = [0u8; 16 * 1024];
659 reader.set_read_timeout(Some(std::time::Duration::from_millis(200))).unwrap();
660 while got.len() < expected {
661 let n = reader.read(&mut buf).unwrap();
662 assert!(n > 0);
663 got.extend_from_slice(&buf[..n]);
664 client.flush().unwrap();
665 }
666 let text = String::from_utf8(got).unwrap();
667 let received: Vec<&str> = text.lines().collect();
668 assert_eq!(received.len(), lines.len());
669 for (want, have) in lines.iter().zip(&received) {
670 assert_eq!(want, have, "a line arrived spliced or out of order");
671 }
672 }
673
674 #[test]
675 fn a_client_that_never_reads_is_dropped_at_the_cap() {
676 let (server, _reader) = UnixStream::pair().unwrap();
677 server.set_nonblocking(true).unwrap();
678 let mut client = Client::new(Subscription::Layout, server);
679 let line = "y".repeat(1000);
680 let mut dropped = false;
681 for _ in 0..2000 {
682 if client.send_line(&line).is_err() {
683 dropped = true;
684 break;
685 }
686 }
687 assert!(dropped, "2 MB to a reader that never reads must hit MAX_PENDING");
688 assert!(client.pending.len() <= MAX_PENDING + line.len() + 1);
689 }
690 }