Commit 3477fb2
Ian Tay
·
2026-03-18 15:54:12 -0400 EDT
parent 98fa966
feat(daemon): buffer PTY stdin writes and flush via POLLOUT Replaces the best-effort ptyWrite (drop on EAGAIN) with a buffered queue flushed by the daemon's poll loop. Mirrors the pattern already used for client-socket writes. - pty_write_buf on Daemon, capped at 256KB (drop new payload on overflow — same failure mode as before, 64x higher threshold) - POLLOUT registered on pty_fd when buffer non-empty; flush handler loops until EAGAIN - handleInput/handleRun queue instead of writing directly - respondToDeviceAttributes routes through the same buffer so DA responses can't interleave with a draining .Run payload after client disconnect Follow-up to #82.
2 files changed,
+60,
-43
+52,
-40
| ... | ... | @@ -341,9 +341,11 @@ const Daemon = struct { | |
| 341 | 341 | task_exit_code: ?u8 = null, // null = running or n/a, set when task completes | |
| 342 | 342 | task_ended_at: ?u64 = null, // timestamp when task exited | |
| 343 | 343 | task_command: ?[]const []const u8 = null, | |
| 344 | + | pty_write_buf: std.ArrayList(u8) = .empty, | |
| 344 | 345 | ||
| 345 | 346 | pub fn deinit(self: *Daemon) void { | |
| 346 | 347 | self.clients.deinit(self.alloc); | |
| 348 | + | self.pty_write_buf.deinit(self.alloc); | |
| 347 | 349 | self.alloc.free(self.socket_path); | |
| 348 | 350 | } | |
| 349 | 351 |
| ... | ... | @@ -547,40 +549,32 @@ const Daemon = struct { | |
| 547 | 549 | return .{ .created = false, .is_daemon = false }; | |
| 548 | 550 | } | |
| 549 | 551 | ||
| 550 | - | /// Best-effort write to the (non-blocking) PTY fd. Retries short writes | |
| 551 | - | /// until complete, but on WouldBlock (kernel buffer full) gives up and | |
| 552 | - | /// drops the remainder — the daemon is single-threaded, so blocking here | |
| 553 | - | /// to wait for POLLOUT would deadlock against a shell that's itself | |
| 554 | - | /// blocked writing echo to a full PTY output buffer that we're not | |
| 555 | - | /// draining. Dropping is the same trade-off the old code made implicitly | |
| 556 | - | /// (short writes were silently truncated), just without the crash. | |
| 557 | - | fn ptyWrite(pty_fd: i32, data: []const u8) void { | |
| 558 | - | var remaining = data; | |
| 559 | - | while (remaining.len > 0) { | |
| 560 | - | const n = posix.write(pty_fd, remaining) catch |err| { | |
| 561 | - | if (err == error.WouldBlock) { | |
| 562 | - | std.log.warn( | |
| 563 | - | "pty write dropped {d}/{d} bytes (buffer full)", | |
| 564 | - | .{ remaining.len, data.len }, | |
| 565 | - | ); | |
| 566 | - | } else { | |
| 567 | - | std.log.warn( | |
| 568 | - | "pty write failed, {d} bytes lost: {s}", | |
| 569 | - | .{ remaining.len, @errorName(err) }, | |
| 570 | - | ); | |
| 571 | - | } | |
| 572 | - | return; | |
| 573 | - | }; | |
| 574 | - | if (n == 0) return; | |
| 575 | - | remaining = remaining[n..]; | |
| 552 | + | const PTY_WRITE_BUF_MAX = 256 * 1024; | |
| 553 | + | ||
| 554 | + | /// Queue bytes for the PTY's stdin. Flushed by daemonLoop on POLLOUT. | |
| 555 | + | /// Drops the payload if the buffer is over cap -- same failure mode as | |
| 556 | + | /// the old direct-write ptyWrite (drop on EAGAIN), just at a 64x higher | |
| 557 | + | /// threshold. Capping avoids OOM when the shell stops reading; dropping | |
| 558 | + | /// new (not old) bytes avoids tearing a partially-accepted sequence. | |
| 559 | + | fn queuePtyInput(self: *Daemon, data: []const u8) void { | |
| 560 | + | if (data.len == 0) return; | |
| 561 | + | if (self.pty_write_buf.items.len + data.len > PTY_WRITE_BUF_MAX) { | |
| 562 | + | std.log.warn( | |
| 563 | + | "pty input dropped {d} bytes (buffer full, shell not reading)", | |
| 564 | + | .{data.len}, | |
| 565 | + | ); | |
| 566 | + | return; | |
| 576 | 567 | } | |
| 568 | + | self.pty_write_buf.appendSlice(self.alloc, data) catch |err| { | |
| 569 | + | std.log.warn( | |
| 570 | + | "pty input dropped {d} bytes: {s}", | |
| 571 | + | .{ data.len, @errorName(err) }, | |
| 572 | + | ); | |
| 573 | + | }; | |
| 577 | 574 | } | |
| 578 | 575 | ||
| 579 | - | pub fn handleInput(self: *Daemon, pty_fd: i32, payload: []const u8) void { | |
| 580 | - | _ = self; | |
| 581 | - | if (payload.len > 0) { | |
| 582 | - | ptyWrite(pty_fd, payload); | |
| 583 | - | } | |
| 576 | + | pub fn handleInput(self: *Daemon, payload: []const u8) void { | |
| 577 | + | self.queuePtyInput(payload); | |
| 584 | 578 | } | |
| 585 | 579 | ||
| 586 | 580 | pub fn handleInit( |
| ... | ... | @@ -760,7 +754,7 @@ const Daemon = struct { | |
| 760 | 754 | } | |
| 761 | 755 | } | |
| 762 | 756 | ||
| 763 | - | pub fn handleRun(self: *Daemon, client: *Client, pty_fd: i32, payload: []const u8) !void { | |
| 757 | + | pub fn handleRun(self: *Daemon, client: *Client, payload: []const u8) !void { | |
| 764 | 758 | // Reset task tracking so the new command's exit marker is detected. | |
| 765 | 759 | // Without this, a second `zmx run` on the same session is ignored | |
| 766 | 760 | // because task_exit_code is still set from the first run. |
| ... | ... | @@ -768,9 +762,7 @@ const Daemon = struct { | |
| 768 | 762 | self.task_ended_at = null; | |
| 769 | 763 | self.is_task_mode = true; | |
| 770 | 764 | ||
| 771 | - | if (payload.len > 0) { | |
| 772 | - | ptyWrite(pty_fd, payload); | |
| 773 | - | } | |
| 765 | + | self.queuePtyInput(payload); | |
| 774 | 766 | try ipc.appendMessage(self.alloc, &client.write_buf, .Ack, ""); | |
| 775 | 767 | client.has_pending_output = true; | |
| 776 | 768 | self.has_had_client = true; |
| ... | ... | @@ -1523,9 +1515,13 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 1523 | 1515 | .revents = 0, | |
| 1524 | 1516 | }); | |
| 1525 | 1517 | ||
| 1518 | + | var pty_events: i16 = posix.POLL.IN; | |
| 1519 | + | if (daemon.pty_write_buf.items.len > 0) { | |
| 1520 | + | pty_events |= posix.POLL.OUT; | |
| 1521 | + | } | |
| 1526 | 1522 | try poll_fds.append(daemon.alloc, .{ | |
| 1527 | 1523 | .fd = pty_fd, | |
| 1528 | - | .events = posix.POLL.IN, | |
| 1524 | + | .events = pty_events, | |
| 1529 | 1525 | .revents = 0, | |
| 1530 | 1526 | }); | |
| 1531 | 1527 |
| ... | ... | @@ -1595,8 +1591,10 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 1595 | 1591 | // This prevents shells like from fish from waiting 2s | |
| 1596 | 1592 | // and then sending a no DA query response warning because | |
| 1597 | 1593 | // there's no client terminal to respond to the query. | |
| 1598 | - | if (daemon.clients.items.len == 0) { | |
| 1599 | - | util.respondToDeviceAttributes(pty_fd, buf[0..n]); | |
| 1594 | + | if (daemon.clients.items.len == 0 and | |
| 1595 | + | daemon.pty_write_buf.items.len < Daemon.PTY_WRITE_BUF_MAX) | |
| 1596 | + | { | |
| 1597 | + | util.respondToDeviceAttributes(daemon.alloc, &daemon.pty_write_buf, buf[0..n]); | |
| 1600 | 1598 | } | |
| 1601 | 1599 | ||
| 1602 | 1600 | // In run mode, scan output for exit code marker |
| ... | ... | @@ -1625,6 +1623,20 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 1625 | 1623 | } | |
| 1626 | 1624 | } | |
| 1627 | 1625 | ||
| 1626 | + | if (poll_fds.items[1].revents & posix.POLL.OUT != 0) { | |
| 1627 | + | while (daemon.pty_write_buf.items.len > 0) { | |
| 1628 | + | const n = posix.write(pty_fd, daemon.pty_write_buf.items) catch |err| { | |
| 1629 | + | if (err != error.WouldBlock) { | |
| 1630 | + | std.log.warn("pty write failed: {s}", .{@errorName(err)}); | |
| 1631 | + | daemon.pty_write_buf.clearRetainingCapacity(); | |
| 1632 | + | } | |
| 1633 | + | break; | |
| 1634 | + | }; | |
| 1635 | + | if (n == 0) break; | |
| 1636 | + | daemon.pty_write_buf.replaceRange(daemon.alloc, 0, n, &[_]u8{}) catch unreachable; | |
| 1637 | + | } | |
| 1638 | + | } | |
| 1639 | + | ||
| 1628 | 1640 | var i: usize = daemon.clients.items.len; | |
| 1629 | 1641 | // Only iterate over clients that were present when poll_fds was constructed | |
| 1630 | 1642 | // poll_fds contains [server, pty, client0, client1, ...] |
| ... | ... | @@ -1662,7 +1674,7 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 1662 | 1674 | ||
| 1663 | 1675 | while (client.read_buf.next()) |msg| { | |
| 1664 | 1676 | switch (msg.header.tag) { | |
| 1665 | - | .Input => daemon.handleInput(pty_fd, msg.payload), | |
| 1677 | + | .Input => daemon.handleInput(msg.payload), | |
| 1666 | 1678 | .Init => try daemon.handleInit(client, pty_fd, &term, msg.payload), | |
| 1667 | 1679 | .Resize => try daemon.handleResize(pty_fd, &term, msg.payload), | |
| 1668 | 1680 | .Detach => { |
| ... | ... | @@ -1678,7 +1690,7 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 1678 | 1690 | }, | |
| 1679 | 1691 | .Info => try daemon.handleInfo(client), | |
| 1680 | 1692 | .History => try daemon.handleHistory(client, &term, msg.payload), | |
| 1681 | - | .Run => try daemon.handleRun(client, pty_fd, msg.payload), | |
| 1693 | + | .Run => try daemon.handleRun(client, msg.payload), | |
| 1682 | 1694 | .Output, .Ack => {}, | |
| 1683 | 1695 | _ => std.log.warn( | |
| 1684 | 1696 | "ignoring unknown IPC tag={d}", |
+8,
-3
| ... | ... | @@ -145,11 +145,16 @@ const DA2_QUERY_EXPLICIT = "\x1b[>0c"; | |
| 145 | 145 | const DA1_RESPONSE = "\x1b[?62;22c"; | |
| 146 | 146 | const DA2_RESPONSE = "\x1b[>1;10;0c"; | |
| 147 | 147 | ||
| 148 | - | pub fn respondToDeviceAttributes(pty_fd: i32, data: []const u8) void { | |
| 148 | + | pub fn respondToDeviceAttributes(alloc: std.mem.Allocator, buf: *std.ArrayList(u8), data: []const u8) void { | |
| 149 | 149 | // Scan for DA queries in PTY output and respond on behalf of the terminal. | |
| 150 | 150 | // This handles the case where no client is attached (e.g. zmx run) | |
| 151 | 151 | // and the shell (e.g. fish) sends a DA query that would otherwise go unanswered. | |
| 152 | 152 | // | |
| 153 | + | // Responses are queued into the daemon's pty_write_buf (not written | |
| 154 | + | // directly) so they don't interleave with any already-buffered input — | |
| 155 | + | // e.g. a large `zmx run` payload still draining after the client | |
| 156 | + | // disconnected. | |
| 157 | + | // | |
| 153 | 158 | // DA1 query: ESC [ c or ESC [ 0 c | |
| 154 | 159 | // DA2 query: ESC [ > c or ESC [ > 0 c | |
| 155 | 160 | // DA1 response (from terminal): ESC [ ? ... c (has '?' after '[') |
| ... | ... | @@ -164,9 +169,9 @@ pub fn respondToDeviceAttributes(pty_fd: i32, data: []const u8) void { | |
| 164 | 169 | continue; | |
| 165 | 170 | } | |
| 166 | 171 | if (matchSeq(data[i..], DA2_QUERY) or matchSeq(data[i..], DA2_QUERY_EXPLICIT)) { | |
| 167 | - | _ = posix.write(pty_fd, DA2_RESPONSE) catch {}; | |
| 172 | + | buf.appendSlice(alloc, DA2_RESPONSE) catch {}; | |
| 168 | 173 | } else if (matchSeq(data[i..], DA1_QUERY) or matchSeq(data[i..], DA1_QUERY_EXPLICIT)) { | |
| 169 | - | _ = posix.write(pty_fd, DA1_RESPONSE) catch {}; | |
| 174 | + | buf.appendSlice(alloc, DA1_RESPONSE) catch {}; | |
| 170 | 175 | } | |
| 171 | 176 | } | |
| 172 | 177 | i += 1; |