Valentine Silvansky
·
2026-07-27
1const std = @import("std");
2const cross = @import("cross.zig");
3const socket = @import("socket.zig");
4const lib_posix = @import("posix.zig");
5
6pub const Tag = enum(u8) {
7 Input = 0,
8 Output = 1,
9 Resize = 2,
10 Detach = 3,
11 DetachAll = 4,
12 Kill = 5,
13 Info = 6,
14 Init = 7,
15 History = 8,
16 Run = 9,
17 Ack = 10,
18 Switch = 11,
19 Write = 12,
20 TaskComplete = 13,
21 LabelGet = 14,
22 LabelSet = 15,
23 LabelClear = 16,
24 LabelData = 17,
25 Send = 18,
26 EnvGet = 19,
27 EnvSet = 20,
28 EnvData = 21,
29 // Non-exhaustive: this enum comes off the wire via bytesToValue and
30 // @enumFromInt, so out-of-range values are representable
31 // rather than UB. Switches must handle `_` (unknown tag).
32 _,
33};
34
35comptime {
36 if (@typeInfo(Tag).@"enum".is_exhaustive) @compileError(
37 "ipc.Tag must stay non-exhaustive -- old daemons rely on `_` to ignore unknown tags",
38 );
39}
40
41pub const Header = packed struct {
42 tag: Tag,
43 len: u32,
44};
45
46pub const Resize = packed struct {
47 rows: u16,
48 cols: u16,
49 xpixel: u16 = 0,
50 ypixel: u16 = 0,
51
52 pub fn winsize(self: Resize) cross.c.struct_winsize {
53 return .{ .ws_row = self.rows, .ws_col = self.cols, .ws_xpixel = self.xpixel, .ws_ypixel = self.ypixel };
54 }
55};
56
57pub fn getTerminalSize(fd: i32) Resize {
58 var ws: cross.c.struct_winsize = undefined;
59 if (cross.c.ioctl(fd, cross.c.TIOCGWINSZ, &ws) == 0 and ws.ws_row > 0 and ws.ws_col > 0) {
60 return .{ .rows = ws.ws_row, .cols = ws.ws_col, .xpixel = ws.ws_xpixel, .ypixel = ws.ws_ypixel };
61 }
62 inline for (.{ lib_posix.STDOUT_FILENO, lib_posix.STDIN_FILENO, lib_posix.STDERR_FILENO }) |fallback_fd| {
63 if (fallback_fd != fd) {
64 if (cross.c.ioctl(fallback_fd, cross.c.TIOCGWINSZ, &ws) == 0 and ws.ws_row > 0 and ws.ws_col > 0) {
65 return .{ .rows = ws.ws_row, .cols = ws.ws_col, .xpixel = ws.ws_xpixel, .ypixel = ws.ws_ypixel };
66 }
67 }
68 }
69 if (lib_posix.open("/dev/tty", .{ .ACCMODE = .RDWR }, 0)) |tty_fd| {
70 defer lib_posix.close(tty_fd);
71 if (cross.c.ioctl(tty_fd, cross.c.TIOCGWINSZ, &ws) == 0 and ws.ws_row > 0 and ws.ws_col > 0) {
72 return .{ .rows = ws.ws_row, .cols = ws.ws_col, .xpixel = ws.ws_xpixel, .ypixel = ws.ws_ypixel };
73 }
74 } else |_| {}
75 return .{ .rows = 24, .cols = 120 };
76}
77
78pub const MAX_CMD_LEN = 256;
79pub const MAX_CWD_LEN = 256;
80
81/// Frozen wire shape. Do NOT add fields! New stats go in new `Tag` values
82/// so old daemons (whose `_` arm ignores unknown tags) stay reachable.
83/// Changing `@sizeOf(Info)` breaks `zmx list` against running daemons.
84pub const Info = extern struct {
85 clients_len: u64,
86 pid: i32,
87 cmd_len: u16,
88 cwd_len: u16,
89 cmd: [MAX_CMD_LEN]u8,
90 cwd: [MAX_CWD_LEN]u8,
91 created_at: u64,
92 task_ended_at: u64,
93 task_exit_code: u8,
94};
95
96pub fn expectedLength(data: []const u8) ?usize {
97 if (data.len < @sizeOf(Header)) return null;
98 const header = std.mem.bytesToValue(Header, data[0..@sizeOf(Header)]);
99 // header.len comes off the wire; widen to usize before adding so a
100 // near-u32-max value can't wrap (panic in safe mode, UB in release).
101 return @as(usize, @sizeOf(Header)) + @as(usize, header.len);
102}
103
104pub fn send(fd: i32, tag: Tag, data: []const u8) !void {
105 const header = Header{
106 .tag = tag,
107 .len = @intCast(data.len),
108 };
109 const header_bytes = std.mem.asBytes(&header);
110 try writeAll(fd, header_bytes);
111 if (data.len > 0) {
112 try writeAll(fd, data);
113 }
114}
115
116pub fn appendMessage(
117 gpa: std.mem.Allocator,
118 list: *std.ArrayList(u8),
119 tag: Tag,
120 data: []const u8,
121) !void {
122 const header = Header{
123 .tag = tag,
124 .len = @intCast(data.len),
125 };
126 // Guarantee capacity for header + payload in one check to avoid
127 // intermediate realloc between the two appends on the hot path.
128 try list.ensureTotalCapacity(gpa, list.items.len + @sizeOf(Header) + data.len);
129 list.appendSliceAssumeCapacity(std.mem.asBytes(&header));
130 if (data.len > 0) {
131 list.appendSliceAssumeCapacity(data);
132 }
133}
134
135/// Pre-0.7.0 daemons expect a 4-byte Init/Resize payload (rows+cols, no
136/// pixel size) and silently drop the 8-byte form, hanging `zmx attach`
137/// against a running old daemon (#211). Append both encodings: every daemon
138/// drops the length it doesn't expect and processes the other exactly once.
139pub const LEGACY_RESIZE_LEN = 4;
140
141pub fn appendSizeMessage(
142 alloc: std.mem.Allocator,
143 list: *std.ArrayList(u8),
144 tag: Tag,
145 size: Resize,
146) !void {
147 const bytes = std.mem.asBytes(&size);
148 try appendMessage(alloc, list, tag, bytes);
149 try appendMessage(alloc, list, tag, bytes[0..LEGACY_RESIZE_LEN]);
150}
151
152fn writeAll(fd: i32, data: []const u8) !void {
153 var index: usize = 0;
154 while (index < data.len) {
155 const n = try lib_posix.write(fd, data[index..]);
156 if (n == 0) return error.DiskQuota;
157 index += n;
158 }
159}
160
161pub const Message = struct {
162 tag: Tag,
163 data: []u8,
164
165 pub fn deinit(self: Message, alloc: std.mem.Allocator) void {
166 if (self.data.len > 0) {
167 alloc.free(self.data);
168 }
169 }
170};
171
172pub const SocketMsg = struct {
173 header: Header,
174 payload: []const u8,
175};
176
177pub const SocketBuffer = struct {
178 buf: std.ArrayList(u8),
179 alloc: std.mem.Allocator,
180 head: usize,
181
182 pub fn init(alloc: std.mem.Allocator) !SocketBuffer {
183 return .{
184 .buf = try std.ArrayList(u8).initCapacity(alloc, 4096),
185 .alloc = alloc,
186 .head = 0,
187 };
188 }
189
190 pub fn deinit(self: *SocketBuffer) void {
191 self.buf.deinit(self.alloc);
192 }
193
194 /// Reads from fd into buffer.
195 /// Returns number of bytes read.
196 /// Propagates error.WouldBlock and other errors to caller.
197 /// Returns 0 on EOF.
198 pub fn read(self: *SocketBuffer, fd: i32) !usize {
199 if (self.head > 0) {
200 const remaining = self.buf.items.len - self.head;
201 if (remaining > 0) {
202 std.mem.copyForwards(u8, self.buf.items[0..remaining], self.buf.items[self.head..]);
203 self.buf.items.len = remaining;
204 } else {
205 self.buf.clearRetainingCapacity();
206 }
207 self.head = 0;
208 }
209
210 var tmp: [4096]u8 = undefined;
211 const n = try lib_posix.read(fd, &tmp);
212 if (n > 0) {
213 try self.buf.appendSlice(self.alloc, tmp[0..n]);
214 }
215 return n;
216 }
217
218 /// Returns the next complete message or `null` when none available.
219 /// `buf` is advanced automatically; caller keeps the returned slices
220 /// valid until the following `next()` (or `deinit`).
221 pub fn next(self: *SocketBuffer) ?SocketMsg {
222 const available = self.buf.items[self.head..];
223 const total = expectedLength(available) orelse return null;
224 if (available.len < total) return null;
225
226 const hdr = std.mem.bytesToValue(Header, available[0..@sizeOf(Header)]);
227 const pay = available[@sizeOf(Header)..total];
228
229 self.head += total;
230 return .{ .header = hdr, .payload = pay };
231 }
232};
233
234const ConnectError = error{
235 ConnectionRefused,
236 Unexpected,
237};
238
239/// Connect-only liveness check. Callers that don't read `Info` should use
240/// this (not `probeSession`) so they survive `Info` shape changes.
241pub fn connectSession(socket_path: []const u8) ConnectError!i32 {
242 return socket.sessionConnect(socket_path) catch |err| switch (err) {
243 error.ConnectionRefused => return error.ConnectionRefused,
244 else => return error.Unexpected,
245 };
246}
247
248const SessionProbeError = error{
249 Timeout,
250 ConnectionRefused,
251 Unexpected,
252 InfoSizeMismatch,
253};
254
255const SessionProbeResult = struct {
256 fd: i32,
257 info: Info,
258 labels: ?[]const u8,
259 alloc: std.mem.Allocator,
260
261 pub fn deinit(self: *const SessionProbeResult) void {
262 if (self.labels) |lbl| self.alloc.free(lbl);
263 lib_posix.close(self.fd);
264 }
265};
266
267pub fn probeSession(
268 alloc: std.mem.Allocator,
269 socket_path: []const u8,
270) SessionProbeError!SessionProbeResult {
271 const timeout_ms = 1000;
272 const fd = try connectSession(socket_path);
273 errdefer lib_posix.close(fd);
274
275 send(fd, .Info, "") catch return error.Unexpected;
276 send(fd, .LabelGet, "") catch {};
277
278 var poll_fds = [_]lib_posix.pollfd{.{ .fd = fd, .events = lib_posix.POLL.IN, .revents = 0 }};
279 const poll_result = lib_posix.poll(&poll_fds, timeout_ms) catch return error.Unexpected;
280 if (poll_result == 0) {
281 return error.Timeout;
282 }
283
284 var sb = SocketBuffer.init(alloc) catch return error.Unexpected;
285 defer sb.deinit();
286
287 const n = sb.read(fd) catch return error.Unexpected;
288 if (n == 0) return error.Unexpected;
289
290 var info_result: ?Info = null;
291 var labels: ?[]const u8 = null;
292 errdefer if (labels) |lbl| alloc.free(lbl);
293
294 while (true) {
295 if (sb.next()) |msg| {
296 if (msg.header.tag == .Info) {
297 if (msg.payload.len != @sizeOf(Info)) return error.InfoSizeMismatch;
298 info_result = std.mem.bytesToValue(Info, msg.payload[0..@sizeOf(Info)]);
299 }
300 if (msg.header.tag == .LabelData) {
301 labels = alloc.dupe(u8, msg.payload) catch null;
302 }
303
304 if (info_result != null and labels != null) break;
305 continue;
306 }
307
308 // No complete message available, wait for more data
309 const more = lib_posix.poll(&poll_fds, 50) catch break;
310 if (more == 0) break;
311 const n_read = sb.read(fd) catch break;
312 if (n_read == 0) break;
313 }
314
315 if (info_result) |info| {
316 return .{
317 .fd = fd,
318 .info = info,
319 .labels = labels,
320 .alloc = alloc,
321 };
322 }
323 return error.Unexpected;
324}
325
326// WIRE PROTOCOL FREEZE: read before "fixing" any test below.
327//
328// Changing these constants does not fix the test; it breaks every
329// running daemon for every user until they `pkill -f zmx`.
330//
331// Need a new field? → add a new `Tag` value (next free integer).
332// Need to remove one? → don't. Reserve the integer, stop sending it.
333test "Info wire size is frozen" {
334 try std.testing.expectEqual(@as(usize, 552), @sizeOf(Info));
335 // packed struct{u8,u32} backs to u40 → @sizeOf rounds to 8, not 5.
336 try std.testing.expectEqual(@as(usize, 8), @sizeOf(Header));
337}
338
339test "Tag wire values are frozen" {
340 inline for (.{
341 .{ Tag.Input, 0 }, .{ Tag.Output, 1 }, .{ Tag.Resize, 2 },
342 .{ Tag.Detach, 3 }, .{ Tag.DetachAll, 4 }, .{ Tag.Kill, 5 },
343 .{ Tag.Info, 6 }, .{ Tag.Init, 7 }, .{ Tag.History, 8 },
344 .{ Tag.Run, 9 }, .{ Tag.Ack, 10 }, .{ Tag.Switch, 11 },
345 .{ Tag.Write, 12 }, .{ Tag.TaskComplete, 13 }, .{ Tag.LabelGet, 14 },
346 .{ Tag.LabelSet, 15 }, .{ Tag.LabelClear, 16 }, .{ Tag.LabelData, 17 },
347 .{ Tag.Send, 18 },
348 }) |p| try std.testing.expectEqual(@as(u8, p[1]), @intFromEnum(p[0]));
349}
350
351pub fn roundTripForTag(
352 alloc: std.mem.Allocator,
353 socket_path: []const u8,
354 request_tag: Tag,
355 payload: []const u8,
356 expected_tag: Tag,
357) SessionProbeError![]u8 {
358 const timeout_ms = 1000;
359 const fd = try connectSession(socket_path);
360 defer lib_posix.close(fd);
361
362 send(fd, request_tag, payload) catch return error.Unexpected;
363
364 var poll_fds = [_]lib_posix.pollfd{.{ .fd = fd, .events = lib_posix.POLL.IN, .revents = 0 }};
365 const poll_result = lib_posix.poll(&poll_fds, timeout_ms) catch return error.Unexpected;
366 if (poll_result == 0) return error.Timeout;
367
368 var sb = SocketBuffer.init(alloc) catch return error.Unexpected;
369 defer sb.deinit();
370
371 const n = sb.read(fd) catch return error.Unexpected;
372 if (n == 0) return error.Unexpected;
373
374 while (sb.next()) |msg| {
375 if (msg.header.tag == expected_tag) {
376 return alloc.dupe(u8, msg.payload) catch return error.Unexpected;
377 }
378 }
379 return error.Unexpected;
380}
381
382test "appendSizeMessage emits current and legacy encodings" {
383 const alloc = std.testing.allocator;
384 var list = try std.ArrayList(u8).initCapacity(alloc, 64);
385 defer list.deinit(alloc);
386
387 const size = Resize{ .rows = 45, .cols = 170, .xpixel = 900, .ypixel = 1800 };
388 try appendSizeMessage(alloc, &list, .Init, size);
389
390 const h1 = std.mem.bytesToValue(Header, list.items[0..@sizeOf(Header)]);
391 try std.testing.expectEqual(Tag.Init, h1.tag);
392 try std.testing.expectEqual(@as(u32, @sizeOf(Resize)), h1.len);
393 const p1 = list.items[@sizeOf(Header)..][0..@sizeOf(Resize)];
394 try std.testing.expectEqual(size, std.mem.bytesToValue(Resize, p1));
395
396 const off2 = @sizeOf(Header) + @sizeOf(Resize);
397 const h2 = std.mem.bytesToValue(Header, list.items[off2..][0..@sizeOf(Header)]);
398 try std.testing.expectEqual(Tag.Init, h2.tag);
399 try std.testing.expectEqual(@as(u32, LEGACY_RESIZE_LEN), h2.len);
400 // Legacy payload is the rows+cols prefix of the current encoding.
401 const p2 = list.items[off2 + @sizeOf(Header) ..][0..LEGACY_RESIZE_LEN];
402 try std.testing.expectEqualSlices(u8, p1[0..LEGACY_RESIZE_LEN], p2);
403 try std.testing.expectEqual(off2 + @sizeOf(Header) + LEGACY_RESIZE_LEN, list.items.len);
404}
405
406test "zeroed Info has no stack garbage in wire bytes" {
407 var info = std.mem.zeroes(Info);
408 info.clients_len = 3;
409 info.pid = 999;
410 info.task_exit_code = 7;
411 const bytes = std.mem.asBytes(&info);
412 // Tail padding after task_exit_code must be zero (asBytes ships it).
413 const last_field_end = @offsetOf(Info, "task_exit_code") + @sizeOf(u8);
414 for (bytes[last_field_end..]) |b| try std.testing.expectEqual(@as(u8, 0), b);
415}