Commit e889a18
Ian Tay
·
2026-04-10 12:06:49 -0400 EDT
parent cd82c84
fix(daemon): self-pipe signal wakeup + version-tolerant IPC liveness
Self-pipe (idle daemon was deaf to SIGTERM/SIGWINCH):
std.posix.poll loops on .INTR internally (PollError has no Interrupted
member), so the prior `catch error.Interrupted` was unreachable since
8640be51/690487b — signals only "worked" when unrelated I/O woke poll().
Replace with the self-pipe trick: posix.pipe2(.{CLOEXEC,NONBLOCK}), one
wakeSignalPipe handler with errno save/restore, pipe read-end at fixed
poll_fds[2] in both loops; resize/shutdown act in the drain branch.
Removes the now-redundant sigwinch/sigterm atomic flags and collapses
setupSig*Handler into installWakeHandler(sig).
IPC versioning (Option E — decouple liveness from Info shape):
connectSession() does connect-only liveness; 5 of 6 callers (kill,
detach, history, run, attach) switch to it so they survive Info struct
changes. probeSession keeps the Info round-trip for `list`, which now
reports InfoSizeMismatch instead of Unexpected on version skew.
Hygiene: clients_len usize->u64 (extern struct); handleInfo zero-inits
Info so asBytes() doesn't ship 7B tail padding + cmd/cwd stack tails.
Tests freeze @sizeOf(Info)=552 / @sizeOf(Header)=8 and assert zeroed
Info ships zero padding; wired via test{_=ipc;}.
2 files changed,
+144,
-108
+45,
-12
| ... | ... | @@ -19,7 +19,7 @@ pub const Tag = enum(u8) { | |
| 19 | 19 | Write = 12, | |
| 20 | 20 | TaskComplete = 13, | |
| 21 | 21 | // Non-exhaustive: this enum comes off the wire via bytesToValue and | |
| 22 | - | // @enumFromInt, so out-of-range values (11-255) are representable | |
| 22 | + | // @enumFromInt, so out-of-range values (14-255) are representable | |
| 23 | 23 | // rather than UB. Switches must handle `_` (unknown tag). | |
| 24 | 24 | _, | |
| 25 | 25 | }; |
| ... | ... | @@ -45,8 +45,11 @@ pub fn getTerminalSize(fd: i32) Resize { | |
| 45 | 45 | pub const MAX_CMD_LEN = 256; | |
| 46 | 46 | pub const MAX_CWD_LEN = 256; | |
| 47 | 47 | ||
| 48 | + | /// Frozen wire shape. Do NOT add fields — new stats go in new `Tag` values | |
| 49 | + | /// so old daemons (whose `_` arm ignores unknown tags) stay reachable. | |
| 50 | + | /// Changing `@sizeOf(Info)` breaks `zmx list` against running daemons. | |
| 48 | 51 | pub const Info = extern struct { | |
| 49 | - | clients_len: usize, | |
| 52 | + | clients_len: u64, | |
| 50 | 53 | pid: i32, | |
| 51 | 54 | cmd_len: u16, | |
| 52 | 55 | cwd_len: u16, |
| ... | ... | @@ -176,10 +179,25 @@ pub const SocketBuffer = struct { | |
| 176 | 179 | } | |
| 177 | 180 | }; | |
| 178 | 181 | ||
| 182 | + | const ConnectError = error{ | |
| 183 | + | ConnectionRefused, | |
| 184 | + | Unexpected, | |
| 185 | + | }; | |
| 186 | + | ||
| 187 | + | /// Unlike `probeSession`, does not round-trip `Info` — kill/detach/history/run | |
| 188 | + | /// stay usable against version-skewed daemons. | |
| 189 | + | pub fn connectSession(socket_path: []const u8) ConnectError!i32 { | |
| 190 | + | return socket.sessionConnect(socket_path) catch |err| switch (err) { | |
| 191 | + | error.ConnectionRefused => return error.ConnectionRefused, | |
| 192 | + | else => return error.Unexpected, | |
| 193 | + | }; | |
| 194 | + | } | |
| 195 | + | ||
| 179 | 196 | const SessionProbeError = error{ | |
| 180 | 197 | Timeout, | |
| 181 | 198 | ConnectionRefused, | |
| 182 | 199 | Unexpected, | |
| 200 | + | InfoSizeMismatch, | |
| 183 | 201 | }; | |
| 184 | 202 | ||
| 185 | 203 | const SessionProbeResult = struct { |
| ... | ... | @@ -192,10 +210,7 @@ pub fn probeSession( | |
| 192 | 210 | socket_path: []const u8, | |
| 193 | 211 | ) SessionProbeError!SessionProbeResult { | |
| 194 | 212 | const timeout_ms = 1000; | |
| 195 | - | const fd = socket.sessionConnect(socket_path) catch |err| switch (err) { | |
| 196 | - | error.ConnectionRefused => return error.ConnectionRefused, | |
| 197 | - | else => return error.Unexpected, | |
| 198 | - | }; | |
| 213 | + | const fd = try connectSession(socket_path); | |
| 199 | 214 | errdefer posix.close(fd); | |
| 200 | 215 | ||
| 201 | 216 | send(fd, .Info, "") catch return error.Unexpected; |
| ... | ... | @@ -214,13 +229,31 @@ pub fn probeSession( | |
| 214 | 229 | ||
| 215 | 230 | while (sb.next()) |msg| { | |
| 216 | 231 | if (msg.header.tag == .Info) { | |
| 217 | - | if (msg.payload.len == @sizeOf(Info)) { | |
| 218 | - | return .{ | |
| 219 | - | .fd = fd, | |
| 220 | - | .info = std.mem.bytesToValue(Info, msg.payload[0..@sizeOf(Info)]), | |
| 221 | - | }; | |
| 222 | - | } | |
| 232 | + | if (msg.payload.len != @sizeOf(Info)) return error.InfoSizeMismatch; | |
| 233 | + | return .{ | |
| 234 | + | .fd = fd, | |
| 235 | + | .info = std.mem.bytesToValue(Info, msg.payload[0..@sizeOf(Info)]), | |
| 236 | + | }; | |
| 223 | 237 | } | |
| 224 | 238 | } | |
| 225 | 239 | return error.Unexpected; | |
| 226 | 240 | } | |
| 241 | + | ||
| 242 | + | test "Info wire size is frozen" { | |
| 243 | + | // Bumping this means version-skewed `zmx list` breaks. See doc comment | |
| 244 | + | // on `Info` — add a new `Tag` instead of growing this struct. | |
| 245 | + | try std.testing.expectEqual(@as(usize, 552), @sizeOf(Info)); | |
| 246 | + | // packed struct{u8,u32} backs to u40 → @sizeOf rounds to 8, not 5. | |
| 247 | + | try std.testing.expectEqual(@as(usize, 8), @sizeOf(Header)); | |
| 248 | + | } | |
| 249 | + | ||
| 250 | + | test "zeroed Info has no stack garbage in wire bytes" { | |
| 251 | + | var info = std.mem.zeroes(Info); | |
| 252 | + | info.clients_len = 3; | |
| 253 | + | info.pid = 999; | |
| 254 | + | info.task_exit_code = 7; | |
| 255 | + | const bytes = std.mem.asBytes(&info); | |
| 256 | + | // Tail padding after task_exit_code must be zero (asBytes ships it). | |
| 257 | + | const last_field_end = @offsetOf(Info, "task_exit_code") + @sizeOf(u8); | |
| 258 | + | for (bytes[last_field_end..]) |b| try std.testing.expectEqual(@as(u8, 0), b); | |
| 259 | + | } |
+99,
-96
| ... | ... | @@ -30,8 +30,10 @@ fn zmxLogFn( | |
| 30 | 30 | log_system.log(level, scope, format, args); | |
| 31 | 31 | } | |
| 32 | 32 | ||
| 33 | - | var sigwinch_received: std.atomic.Value(bool) = std.atomic.Value(bool).init(false); | |
| 34 | - | var sigterm_received: std.atomic.Value(bool) = std.atomic.Value(bool).init(false); | |
| 33 | + | /// Self-pipe woken by signal handlers. std.posix.poll loops on .INTR internally | |
| 34 | + | /// (PollError has no Interrupted member), so a signal that lands during poll() | |
| 35 | + | /// never surfaces; the handler writes a byte here and poll() wakes on POLLIN. | |
| 36 | + | var sig_pipe: [2]posix.fd_t = .{ -1, -1 }; | |
| 35 | 37 | ||
| 36 | 38 | // https://github.com/ziglang/zig/blob/738d2be9d6b6ef3ff3559130c05159ef53336224/lib/std/posix.zig#L3505 | |
| 37 | 39 | const O_NONBLOCK: usize = 1 << @bitOffsetOf(posix.O, "NONBLOCK"); |
| ... | ... | @@ -55,6 +57,18 @@ fn parseSessionArg(alloc: std.mem.Allocator, raw: []const u8) !SessionMatch { | |
| 55 | 57 | return .{ .name = name, .is_prefix = false }; | |
| 56 | 58 | } | |
| 57 | 59 | ||
| 60 | + | fn openSignalPipe() !void { | |
| 61 | + | sig_pipe = try posix.pipe2(.{ .CLOEXEC = true, .NONBLOCK = true }); | |
| 62 | + | } | |
| 63 | + | ||
| 64 | + | fn drainSignalPipe() void { | |
| 65 | + | var b: [16]u8 = undefined; | |
| 66 | + | while (true) { | |
| 67 | + | const n = posix.read(sig_pipe[0], &b) catch return; | |
| 68 | + | if (n == 0) return; | |
| 69 | + | } | |
| 70 | + | } | |
| 71 | + | ||
| 58 | 72 | pub fn main() !void { | |
| 59 | 73 | // use c_allocator to avoid "reached unreachable code" panic in DebugAllocator when forking | |
| 60 | 74 | const alloc = std.heap.c_allocator; |
| ... | ... | @@ -694,8 +708,8 @@ const Daemon = struct { | |
| 694 | 708 | var should_create = !exists; | |
| 695 | 709 | ||
| 696 | 710 | if (exists) { | |
| 697 | - | if (ipc.probeSession(self.alloc, self.socket_path)) |result| { | |
| 698 | - | posix.close(result.fd); | |
| 711 | + | if (ipc.connectSession(self.socket_path)) |fd| { | |
| 712 | + | posix.close(fd); | |
| 699 | 713 | if (self.command != null) { | |
| 700 | 714 | std.log.warn( | |
| 701 | 715 | "session already exists, ignoring command session={s}", |
| ... | ... | @@ -708,12 +722,12 @@ const Daemon = struct { | |
| 708 | 722 | socket.cleanupStaleSocket(dir, self.session_name); | |
| 709 | 723 | should_create = true; | |
| 710 | 724 | }, | |
| 711 | - | // Probe didn't respond in time -- daemon may just be busy. | |
| 712 | - | // The probe is only to decide create-vs-attach; the session | |
| 713 | - | // exists, so proceed to attach rather than fail or orphan. | |
| 725 | + | // Connect failed for an unusual reason. The check is only to | |
| 726 | + | // decide create-vs-attach; the socket file exists, so proceed | |
| 727 | + | // to attach rather than fail or orphan. | |
| 714 | 728 | else => { | |
| 715 | 729 | std.log.warn( | |
| 716 | - | "probe slow ({s}), proceeding to attach session={s}", | |
| 730 | + | "connect failed ({s}), proceeding to attach session={s}", | |
| 717 | 731 | .{ @errorName(err), self.session_name }, | |
| 718 | 732 | ); | |
| 719 | 733 | }, |
| ... | ... | @@ -1012,12 +1026,17 @@ const Daemon = struct { | |
| 1012 | 1026 | } | |
| 1013 | 1027 | ||
| 1014 | 1028 | pub fn handleInfo(self: *Daemon, client: *Client) !void { | |
| 1015 | - | const clients_len = self.clients.items.len - 1; | |
| 1029 | + | // zeroes() so asBytes() doesn't ship struct padding + unused cmd/cwd | |
| 1030 | + | // tail bytes (daemon stack contents) to clients. | |
| 1031 | + | var info = std.mem.zeroes(ipc.Info); | |
| 1032 | + | info.clients_len = self.clients.items.len - 1; | |
| 1033 | + | info.pid = self.pid; | |
| 1034 | + | info.created_at = self.created_at; | |
| 1035 | + | info.task_ended_at = self.task_ended_at orelse 0; | |
| 1036 | + | info.task_exit_code = self.task_exit_code orelse 0; | |
| 1016 | 1037 | ||
| 1017 | 1038 | // Build command string from args, re-quoting args that contain | |
| 1018 | 1039 | // shell-special characters so the displayed command is copy-pasteable. | |
| 1019 | - | var cmd_buf: [ipc.MAX_CMD_LEN]u8 = undefined; | |
| 1020 | - | var cmd_len: u16 = 0; | |
| 1021 | 1040 | const cur_cmd = self.command; | |
| 1022 | 1041 | if (cur_cmd) |args| { | |
| 1023 | 1042 | for (args, 0..) |arg, i| { |
| ... | ... | @@ -1029,40 +1048,27 @@ const Daemon = struct { | |
| 1029 | 1048 | const src = quoted orelse arg; | |
| 1030 | 1049 | ||
| 1031 | 1050 | const need = src.len + @as(usize, if (i > 0) 1 else 0); | |
| 1032 | - | if (cmd_len + need > ipc.MAX_CMD_LEN) { | |
| 1051 | + | if (info.cmd_len + need > ipc.MAX_CMD_LEN) { | |
| 1033 | 1052 | const ellipsis = "..."; | |
| 1034 | - | if (cmd_len + ellipsis.len <= ipc.MAX_CMD_LEN) { | |
| 1035 | - | @memcpy(cmd_buf[cmd_len..][0..ellipsis.len], ellipsis); | |
| 1036 | - | cmd_len += ellipsis.len; | |
| 1053 | + | if (info.cmd_len + ellipsis.len <= ipc.MAX_CMD_LEN) { | |
| 1054 | + | @memcpy(info.cmd[info.cmd_len..][0..ellipsis.len], ellipsis); | |
| 1055 | + | info.cmd_len += ellipsis.len; | |
| 1037 | 1056 | } | |
| 1038 | 1057 | break; | |
| 1039 | 1058 | } | |
| 1040 | 1059 | ||
| 1041 | 1060 | if (i > 0) { | |
| 1042 | - | cmd_buf[cmd_len] = ' '; | |
| 1043 | - | cmd_len += 1; | |
| 1061 | + | info.cmd[info.cmd_len] = ' '; | |
| 1062 | + | info.cmd_len += 1; | |
| 1044 | 1063 | } | |
| 1045 | - | @memcpy(cmd_buf[cmd_len..][0..src.len], src); | |
| 1046 | - | cmd_len += @intCast(src.len); | |
| 1064 | + | @memcpy(info.cmd[info.cmd_len..][0..src.len], src); | |
| 1065 | + | info.cmd_len += @intCast(src.len); | |
| 1047 | 1066 | } | |
| 1048 | 1067 | } | |
| 1049 | 1068 | ||
| 1050 | - | // Copy cwd | |
| 1051 | - | var cwd_buf: [ipc.MAX_CWD_LEN]u8 = undefined; | |
| 1052 | - | const cwd_len: u16 = @intCast(@min(self.cwd.len, ipc.MAX_CWD_LEN)); | |
| 1053 | - | @memcpy(cwd_buf[0..cwd_len], self.cwd[0..cwd_len]); | |
| 1054 | - | ||
| 1055 | - | const info = ipc.Info{ | |
| 1056 | - | .clients_len = clients_len, | |
| 1057 | - | .pid = self.pid, | |
| 1058 | - | .cmd_len = cmd_len, | |
| 1059 | - | .cwd_len = cwd_len, | |
| 1060 | - | .cmd = cmd_buf, | |
| 1061 | - | .cwd = cwd_buf, | |
| 1062 | - | .created_at = self.created_at, | |
| 1063 | - | .task_ended_at = self.task_ended_at orelse 0, | |
| 1064 | - | .task_exit_code = self.task_exit_code orelse 0, | |
| 1065 | - | }; | |
| 1069 | + | info.cwd_len = @intCast(@min(self.cwd.len, ipc.MAX_CWD_LEN)); | |
| 1070 | + | @memcpy(info.cwd[0..info.cwd_len], self.cwd[0..info.cwd_len]); | |
| 1071 | + | ||
| 1066 | 1072 | try ipc.appendMessage(self.alloc, &client.write_buf, .Info, std.mem.asBytes(&info)); | |
| 1067 | 1073 | client.has_pending_output = true; | |
| 1068 | 1074 | } |
| ... | ... | @@ -1617,13 +1623,13 @@ fn detachAll(cfg: *Cfg) !void { | |
| 1617 | 1623 | error.OutOfMemory => return err, | |
| 1618 | 1624 | }; | |
| 1619 | 1625 | defer alloc.free(socket_path); | |
| 1620 | - | const result = ipc.probeSession(alloc, socket_path) catch |err| { | |
| 1626 | + | const fd = ipc.connectSession(socket_path) catch |err| { | |
| 1621 | 1627 | std.log.err("session unresponsive: {s}", .{@errorName(err)}); | |
| 1622 | 1628 | if (err == error.ConnectionRefused) socket.cleanupStaleSocket(dir, session_name); | |
| 1623 | 1629 | return; | |
| 1624 | 1630 | }; | |
| 1625 | - | defer posix.close(result.fd); | |
| 1626 | - | ipc.send(result.fd, .DetachAll, "") catch |err| switch (err) { | |
| 1631 | + | defer posix.close(fd); | |
| 1632 | + | ipc.send(fd, .DetachAll, "") catch |err| switch (err) { | |
| 1627 | 1633 | error.BrokenPipe, error.ConnectionResetByPeer => return, | |
| 1628 | 1634 | else => return err, | |
| 1629 | 1635 | }; |
| ... | ... | @@ -1651,7 +1657,7 @@ fn kill(cfg: *Cfg, session_name: []const u8, force: bool) !void { | |
| 1651 | 1657 | w.interface.flush() catch {}; | |
| 1652 | 1658 | return error.SessionNotFound; | |
| 1653 | 1659 | } | |
| 1654 | - | const result = ipc.probeSession(alloc, socket_path) catch |err| { | |
| 1660 | + | const fd = ipc.connectSession(socket_path) catch |err| { | |
| 1655 | 1661 | std.log.err("session unresponsive: {s}", .{@errorName(err)}); | |
| 1656 | 1662 | var buf: [4096]u8 = undefined; | |
| 1657 | 1663 | var w = std.fs.File.stdout().writer(&buf); |
| ... | ... | @@ -1668,8 +1674,8 @@ fn kill(cfg: *Cfg, session_name: []const u8, force: bool) !void { | |
| 1668 | 1674 | return; | |
| 1669 | 1675 | }; | |
| 1670 | 1676 | ||
| 1671 | - | defer posix.close(result.fd); | |
| 1672 | - | ipc.send(result.fd, .Kill, "") catch |err| switch (err) { | |
| 1677 | + | defer posix.close(fd); | |
| 1678 | + | ipc.send(fd, .Kill, "") catch |err| switch (err) { | |
| 1673 | 1679 | error.BrokenPipe, error.ConnectionResetByPeer => return, | |
| 1674 | 1680 | else => return err, | |
| 1675 | 1681 | }; |
| ... | ... | @@ -1702,15 +1708,15 @@ fn history(cfg: *Cfg, session_name: []const u8, format: util.HistoryFormat) !voi | |
| 1702 | 1708 | w.interface.flush() catch {}; | |
| 1703 | 1709 | return error.SessionNotFound; | |
| 1704 | 1710 | } | |
| 1705 | - | const result = ipc.probeSession(alloc, socket_path) catch |err| { | |
| 1711 | + | const fd = ipc.connectSession(socket_path) catch |err| { | |
| 1706 | 1712 | std.log.err("session unresponsive: {s}", .{@errorName(err)}); | |
| 1707 | 1713 | if (err == error.ConnectionRefused) socket.cleanupStaleSocket(dir, session_name); | |
| 1708 | 1714 | return; | |
| 1709 | 1715 | }; | |
| 1710 | - | defer posix.close(result.fd); | |
| 1716 | + | defer posix.close(fd); | |
| 1711 | 1717 | ||
| 1712 | 1718 | const format_byte = [_]u8{@intFromEnum(format)}; | |
| 1713 | - | ipc.send(result.fd, .History, &format_byte) catch |err| switch (err) { | |
| 1719 | + | ipc.send(fd, .History, &format_byte) catch |err| switch (err) { | |
| 1714 | 1720 | error.BrokenPipe, error.ConnectionResetByPeer => return, | |
| 1715 | 1721 | else => return err, | |
| 1716 | 1722 | }; |
| ... | ... | @@ -1719,14 +1725,14 @@ fn history(cfg: *Cfg, session_name: []const u8, format: util.HistoryFormat) !voi | |
| 1719 | 1725 | defer sb.deinit(); | |
| 1720 | 1726 | ||
| 1721 | 1727 | while (true) { | |
| 1722 | - | var poll_fds = [_]posix.pollfd{.{ .fd = result.fd, .events = posix.POLL.IN, .revents = 0 }}; | |
| 1728 | + | var poll_fds = [_]posix.pollfd{.{ .fd = fd, .events = posix.POLL.IN, .revents = 0 }}; | |
| 1723 | 1729 | const poll_result = posix.poll(&poll_fds, 5000) catch return; | |
| 1724 | 1730 | if (poll_result == 0) { | |
| 1725 | 1731 | std.log.err("timeout waiting for history response", .{}); | |
| 1726 | 1732 | return; | |
| 1727 | 1733 | } | |
| 1728 | 1734 | ||
| 1729 | - | const n = sb.read(result.fd) catch return; | |
| 1735 | + | const n = sb.read(fd) catch return; | |
| 1730 | 1736 | if (n == 0) return; | |
| 1731 | 1737 | ||
| 1732 | 1738 | while (sb.next()) |msg| { |
| ... | ... | @@ -2103,7 +2109,10 @@ fn run(daemon: *Daemon, detached: bool, command_args: [][]const u8) !void { | |
| 2103 | 2109 | return error.CommandRequired; | |
| 2104 | 2110 | } | |
| 2105 | 2111 | ||
| 2106 | - | const client_sock = try socket.sessionConnect(daemon.socket_path); | |
| 2112 | + | const client_sock = ipc.connectSession(daemon.socket_path) catch |err| { | |
| 2113 | + | std.log.err("session not ready: {s}", .{@errorName(err)}); | |
| 2114 | + | return error.SessionNotReady; | |
| 2115 | + | }; | |
| 2107 | 2116 | defer posix.close(client_sock); | |
| 2108 | 2117 | ||
| 2109 | 2118 | var fds = try std.ArrayList(i32).initCapacity(alloc, 1); |
| ... | ... | @@ -2134,7 +2143,8 @@ fn clientLoop(client_sock_fd: i32) !ClientResult { | |
| 2134 | 2143 | const alloc = std.heap.c_allocator; | |
| 2135 | 2144 | defer posix.close(client_sock_fd); | |
| 2136 | 2145 | ||
| 2137 | - | setupSigwinchHandler(); | |
| 2146 | + | try openSignalPipe(); | |
| 2147 | + | installWakeHandler(posix.SIG.WINCH); | |
| 2138 | 2148 | ||
| 2139 | 2149 | // Make socket non-blocking to avoid blocking on writes | |
| 2140 | 2150 | var sock_flags = try posix.fcntl(client_sock_fd, posix.F.GETFL, 0); |
| ... | ... | @@ -2168,12 +2178,6 @@ fn clientLoop(client_sock_fd: i32) !ClientResult { | |
| 2168 | 2178 | defer _ = posix.fcntl(stdin_fd, posix.F.SETFL, stdin_orig_flags) catch {}; | |
| 2169 | 2179 | ||
| 2170 | 2180 | while (true) { | |
| 2171 | - | // Check for pending SIGWINCH | |
| 2172 | - | if (sigwinch_received.swap(false, .acq_rel)) { | |
| 2173 | - | const next_size = ipc.getTerminalSize(posix.STDOUT_FILENO); | |
| 2174 | - | try ipc.appendMessage(alloc, &sock_write_buf, .Resize, std.mem.asBytes(&next_size)); | |
| 2175 | - | } | |
| 2176 | - | ||
| 2177 | 2181 | poll_fds.clearRetainingCapacity(); | |
| 2178 | 2182 | ||
| 2179 | 2183 | try poll_fds.append(alloc, .{ |
| ... | ... | @@ -2193,6 +2197,8 @@ fn clientLoop(client_sock_fd: i32) !ClientResult { | |
| 2193 | 2197 | .revents = 0, | |
| 2194 | 2198 | }); | |
| 2195 | 2199 | ||
| 2200 | + | try poll_fds.append(alloc, .{ .fd = sig_pipe[0], .events = posix.POLL.IN, .revents = 0 }); | |
| 2201 | + | ||
| 2196 | 2202 | if (stdout_buf.items.len > 0) { | |
| 2197 | 2203 | try poll_fds.append(alloc, .{ | |
| 2198 | 2204 | .fd = posix.STDOUT_FILENO, |
| ... | ... | @@ -2201,10 +2207,13 @@ fn clientLoop(client_sock_fd: i32) !ClientResult { | |
| 2201 | 2207 | }); | |
| 2202 | 2208 | } | |
| 2203 | 2209 | ||
| 2204 | - | _ = posix.poll(poll_fds.items, -1) catch |err| { | |
| 2205 | - | if (err == error.Interrupted) continue; // EINTR from signal, loop again | |
| 2206 | - | return err; | |
| 2207 | - | }; | |
| 2210 | + | _ = try posix.poll(poll_fds.items, -1); | |
| 2211 | + | ||
| 2212 | + | if (poll_fds.items[2].revents & posix.POLL.IN != 0) { | |
| 2213 | + | drainSignalPipe(); | |
| 2214 | + | const next_size = ipc.getTerminalSize(posix.STDOUT_FILENO); | |
| 2215 | + | try ipc.appendMessage(alloc, &sock_write_buf, .Resize, std.mem.asBytes(&next_size)); | |
| 2216 | + | } | |
| 2208 | 2217 | ||
| 2209 | 2218 | // Handle stdin -> socket (Input) | |
| 2210 | 2219 | const inp_flags = (posix.POLL.IN | posix.POLL.HUP | posix.POLL.ERR | posix.POLL.NVAL); |
| ... | ... | @@ -2308,7 +2317,8 @@ fn clientLoop(client_sock_fd: i32) !ClientResult { | |
| 2308 | 2317 | fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 2309 | 2318 | std.log.info("daemon started session={s} pty_fd={d}", .{ daemon.session_name, pty_fd }); | |
| 2310 | 2319 | daemon.pty_fd = pty_fd; | |
| 2311 | - | setupSigtermHandler(); | |
| 2320 | + | try openSignalPipe(); | |
| 2321 | + | installWakeHandler(posix.SIG.TERM); | |
| 2312 | 2322 | var poll_fds = try std.ArrayList(posix.pollfd).initCapacity(daemon.alloc, 8); | |
| 2313 | 2323 | defer poll_fds.deinit(daemon.alloc); | |
| 2314 | 2324 |
| ... | ... | @@ -2323,14 +2333,6 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 2323 | 2333 | defer vt_stream.deinit(); | |
| 2324 | 2334 | ||
| 2325 | 2335 | daemon_loop: while (daemon.running) { | |
| 2326 | - | if (sigterm_received.swap(false, .acq_rel)) { | |
| 2327 | - | std.log.info( | |
| 2328 | - | "SIGTERM received, shutting down gracefully session={s}", | |
| 2329 | - | .{daemon.session_name}, | |
| 2330 | - | ); | |
| 2331 | - | break :daemon_loop; | |
| 2332 | - | } | |
| 2333 | - | ||
| 2334 | 2336 | poll_fds.clearRetainingCapacity(); | |
| 2335 | 2337 | ||
| 2336 | 2338 | try poll_fds.append(daemon.alloc, .{ |
| ... | ... | @@ -2349,6 +2351,8 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 2349 | 2351 | .revents = 0, | |
| 2350 | 2352 | }); | |
| 2351 | 2353 | ||
| 2354 | + | try poll_fds.append(daemon.alloc, .{ .fd = sig_pipe[0], .events = posix.POLL.IN, .revents = 0 }); | |
| 2355 | + | ||
| 2352 | 2356 | for (daemon.clients.items) |client| { | |
| 2353 | 2357 | var events: i16 = posix.POLL.IN; | |
| 2354 | 2358 | if (client.has_pending_output) { |
| ... | ... | @@ -2361,10 +2365,16 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 2361 | 2365 | }); | |
| 2362 | 2366 | } | |
| 2363 | 2367 | ||
| 2364 | - | _ = posix.poll(poll_fds.items, -1) catch |err| { | |
| 2365 | - | if (err == error.Interrupted) continue; | |
| 2366 | - | return err; | |
| 2367 | - | }; | |
| 2368 | + | _ = try posix.poll(poll_fds.items, -1); | |
| 2369 | + | ||
| 2370 | + | if (poll_fds.items[2].revents & posix.POLL.IN != 0) { | |
| 2371 | + | drainSignalPipe(); | |
| 2372 | + | std.log.info( | |
| 2373 | + | "SIGTERM received, shutting down gracefully session={s}", | |
| 2374 | + | .{daemon.session_name}, | |
| 2375 | + | ); | |
| 2376 | + | break :daemon_loop; | |
| 2377 | + | } | |
| 2368 | 2378 | ||
| 2369 | 2379 | if (poll_fds.items[0].revents & (posix.POLL.ERR | posix.POLL.HUP | posix.POLL.NVAL) != 0) { | |
| 2370 | 2380 | std.log.err("server socket error revents={d}", .{poll_fds.items[0].revents}); |
| ... | ... | @@ -2473,9 +2483,9 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 2473 | 2483 | ||
| 2474 | 2484 | var i: usize = daemon.clients.items.len; | |
| 2475 | 2485 | // Only iterate over clients that were present when poll_fds was constructed | |
| 2476 | - | // poll_fds contains [server, pty, client0, client1, ...] | |
| 2477 | - | // So number of clients in poll_fds is poll_fds.items.len - 2 | |
| 2478 | - | const num_polled_clients = poll_fds.items.len - 2; | |
| 2486 | + | // poll_fds contains [server, pty, sig_pipe, client0, client1, ...] | |
| 2487 | + | // So number of clients in poll_fds is poll_fds.items.len - 3 | |
| 2488 | + | const num_polled_clients = poll_fds.items.len - 3; | |
| 2479 | 2489 | if (i > num_polled_clients) { | |
| 2480 | 2490 | // If we have more clients than polled (i.e. we just accepted one), start from the | |
| 2481 | 2491 | // polled ones |
| ... | ... | @@ -2485,7 +2495,7 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 2485 | 2495 | clients_loop: while (i > 0) { | |
| 2486 | 2496 | i -= 1; | |
| 2487 | 2497 | const client = daemon.clients.items[i]; | |
| 2488 | - | const revents = poll_fds.items[i + 2].revents; | |
| 2498 | + | const revents = poll_fds.items[i + 3].revents; | |
| 2489 | 2499 | ||
| 2490 | 2500 | if (revents & posix.POLL.IN != 0) { | |
| 2491 | 2501 | const n = client.read_buf.read(client.socket_fd) catch |err| { |
| ... | ... | @@ -2564,33 +2574,22 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 2564 | 2574 | } | |
| 2565 | 2575 | } | |
| 2566 | 2576 | ||
| 2567 | - | fn handleSigwinch(_: i32, _: *const posix.siginfo_t, _: ?*anyopaque) callconv(.c) void { | |
| 2568 | - | sigwinch_received.store(true, .release); | |
| 2577 | + | fn wakeSignalPipe(_: i32, _: *const posix.siginfo_t, _: ?*anyopaque) callconv(.c) void { | |
| 2578 | + | const saved = std.c._errno().*; | |
| 2579 | + | _ = std.c.write(sig_pipe[1], "x", 1); | |
| 2580 | + | std.c._errno().* = saved; | |
| 2569 | 2581 | } | |
| 2570 | 2582 | ||
| 2571 | - | fn handleSigterm(_: i32, _: *const posix.siginfo_t, _: ?*anyopaque) callconv(.c) void { | |
| 2572 | - | sigterm_received.store(true, .release); | |
| 2573 | - | } | |
| 2574 | - | ||
| 2575 | - | // No SA_RESTART: we want the signal to interrupt poll() so the | |
| 2576 | - | // loop can check the flag. On BSD/macOS, SA_RESTART makes poll restartable, | |
| 2577 | - | // which would leave an idle daemon deaf to SIGTERM until other I/O wakes it. | |
| 2578 | - | fn setupSigwinchHandler() void { | |
| 2583 | + | // std.posix.poll retries EINTR internally, so SA_RESTART is moot — neither | |
| 2584 | + | // setting wakes the loop. The handler writes to sig_pipe instead; poll() | |
| 2585 | + | // wakes on its read end. | |
| 2586 | + | fn installWakeHandler(sig: u6) void { | |
| 2579 | 2587 | const act: posix.Sigaction = .{ | |
| 2580 | - | .handler = .{ .sigaction = handleSigwinch }, | |
| 2588 | + | .handler = .{ .sigaction = wakeSignalPipe }, | |
| 2581 | 2589 | .mask = posix.sigemptyset(), | |
| 2582 | 2590 | .flags = posix.SA.SIGINFO, | |
| 2583 | 2591 | }; | |
| 2584 | - | posix.sigaction(posix.SIG.WINCH, &act, null); | |
| 2585 | - | } | |
| 2586 | - | ||
| 2587 | - | fn setupSigtermHandler() void { | |
| 2588 | - | const act: posix.Sigaction = .{ | |
| 2589 | - | .handler = .{ .sigaction = handleSigterm }, | |
| 2590 | - | .mask = posix.sigemptyset(), | |
| 2591 | - | .flags = posix.SA.SIGINFO, | |
| 2592 | - | }; | |
| 2593 | - | posix.sigaction(posix.SIG.TERM, &act, null); | |
| 2592 | + | posix.sigaction(sig, &act, null); | |
| 2594 | 2593 | } | |
| 2595 | 2594 | ||
| 2596 | 2595 | fn ignoreSigpipe() void { |