main zmx / src / ipc.zig
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}