main zmx / src / loop.zig
Eric Bower  ·  2026-09-10
   1const std = @import("std");
   2const ghostty_vt = @import("ghostty-vt");
   3const ipc = @import("ipc.zig");
   4const log = @import("log.zig");
   5const util = @import("util.zig");
   6const cross = @import("cross.zig");
   7const socket = @import("socket.zig");
   8const label = @import("label.zig");
   9const lib_posix = @import("posix.zig");
  10const Cfg = @import("cfg.zig");
  11const signal = @import("signal.zig");
  12const assert = std.debug.assert;
  13const daemonize = @import("daemonize.zig");
  14const builtin = @import("builtin");
  15
  16/// clientLoop sends ipc commands to its corresponding daemon.  It uses poll() as its non-blocking
  17/// mechanism. It will send stdin to the daemon and receive stdout from the daemon.
  18pub fn clientLoop(client_sock_fd: i32, env_str: []const u8) !ClientResult {
  19    std.log.info("client loop fd={d}", .{client_sock_fd});
  20    const gpa: std.mem.Allocator = blk: {
  21        if (builtin.mode == .Debug) {
  22            const GPA = std.heap.DebugAllocator(.{});
  23            const Static = struct {
  24                var gpa: GPA = .{};
  25            };
  26            break :blk Static.gpa.allocator();
  27        }
  28        break :blk std.heap.c_allocator;
  29    };
  30    defer lib_posix.close(client_sock_fd);
  31
  32    try signal.openSignalPipe();
  33    signal.installWakeHandler(@intFromEnum(lib_posix.SIG.WINCH));
  34
  35    // Make socket non-blocking to avoid blocking on writes
  36    var sock_flags = try lib_posix.fcntl(client_sock_fd, lib_posix.F.GETFL, 0);
  37    sock_flags |= lib_posix.O_NONBLOCK;
  38    _ = try lib_posix.fcntl(client_sock_fd, lib_posix.F.SETFL, sock_flags);
  39
  40    // Buffer for outgoing socket writes
  41    var sock_write_buf = try std.ArrayList(u8).initCapacity(gpa, 4096);
  42    defer sock_write_buf.deinit(gpa);
  43
  44    if (env_str.len > 0) {
  45        try ipc.appendMessage(gpa, &sock_write_buf, .EnvSet, env_str);
  46    }
  47
  48    // Send init message with terminal size (buffered)
  49    const size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO);
  50    try ipc.appendSizeMessage(gpa, &sock_write_buf, .Init, size);
  51
  52    var poll_fds = try std.ArrayList(lib_posix.pollfd).initCapacity(gpa, 4);
  53    defer poll_fds.deinit(gpa);
  54
  55    var read_buf = try ipc.SocketBuffer.init(gpa);
  56    defer read_buf.deinit();
  57
  58    var stdout_buf = try std.ArrayList(u8).initCapacity(gpa, 4096);
  59    defer stdout_buf.deinit(gpa);
  60
  61    const stdin_fd = lib_posix.STDIN_FILENO;
  62
  63    // Make stdin non-blocking. O_NONBLOCK is set on the open file description,
  64    // which is shared with the parent shell; restore on exit to avoid
  65    // corrupting the parent's stdin.
  66    const stdin_orig_flags = try lib_posix.fcntl(stdin_fd, lib_posix.F.GETFL, 0);
  67    _ = try lib_posix.fcntl(stdin_fd, lib_posix.F.SETFL, stdin_orig_flags | lib_posix.O_NONBLOCK);
  68    defer _ = lib_posix.fcntl(stdin_fd, lib_posix.F.SETFL, stdin_orig_flags) catch {};
  69
  70    const detach_key_disabled = util.isDetachKeyDisabled();
  71
  72    while (true) {
  73        poll_fds.clearRetainingCapacity();
  74
  75        try poll_fds.append(gpa, .{
  76            .fd = stdin_fd,
  77            .events = lib_posix.POLL.IN,
  78            .revents = 0,
  79        });
  80
  81        // Poll socket for read, and also for write if we have pending data
  82        var sock_events: i16 = lib_posix.POLL.IN;
  83        if (sock_write_buf.items.len > 0) {
  84            sock_events |= lib_posix.POLL.OUT;
  85        }
  86        try poll_fds.append(gpa, .{
  87            .fd = client_sock_fd,
  88            .events = sock_events,
  89            .revents = 0,
  90        });
  91
  92        try poll_fds.append(gpa, .{ .fd = signal.sig_pipe[0], .events = lib_posix.POLL.IN, .revents = 0 });
  93
  94        if (stdout_buf.items.len > 0) {
  95            try poll_fds.append(gpa, .{
  96                .fd = lib_posix.STDOUT_FILENO,
  97                .events = lib_posix.POLL.OUT,
  98                .revents = 0,
  99            });
 100        }
 101
 102        _ = try lib_posix.poll(poll_fds.items, -1);
 103
 104        if (poll_fds.items[2].revents & lib_posix.POLL.IN != 0) {
 105            signal.drainSignalPipe();
 106            const next_size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO);
 107            try ipc.appendSizeMessage(gpa, &sock_write_buf, .Resize, next_size);
 108        }
 109
 110        // Handle stdin -> socket (Input)
 111        const inp_flags = (lib_posix.POLL.IN | lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL);
 112        if (poll_fds.items[0].revents & inp_flags != 0) {
 113            var buf: [4096]u8 = undefined;
 114            const n_opt: ?usize = lib_posix.read(stdin_fd, &buf) catch |err| blk: {
 115                if (err == error.WouldBlock) break :blk null;
 116                return err;
 117            };
 118
 119            if (n_opt) |n| {
 120                if (n > 0) {
 121                    // Check for detach sequences (ctrl+\ as first byte or Kitty escape sequence)
 122                    if (!detach_key_disabled and util.isCtrlBackslash(buf[0..n])) {
 123                        std.log.info("detach key detected", .{});
 124                        try ipc.appendMessage(gpa, &sock_write_buf, .Detach, "");
 125                    } else {
 126                        try ipc.appendMessage(gpa, &sock_write_buf, .Input, buf[0..n]);
 127                    }
 128                } else {
 129                    std.log.info("eof stdin", .{});
 130                    // EOF on stdin
 131                    return ClientResult{ .kind = .detach, .session_name = null };
 132                }
 133            }
 134        }
 135
 136        // Handle socket read (incoming Output messages from daemon)
 137        if (poll_fds.items[1].revents & lib_posix.POLL.IN != 0) {
 138            const n = read_buf.read(client_sock_fd) catch |err| {
 139                if (err == error.WouldBlock) continue;
 140                if (err == error.ConnectionResetByPeer or err == error.BrokenPipe) {
 141                    return ClientResult{ .kind = .detach, .session_name = null };
 142                }
 143                std.log.err("daemon read err={s}", .{@errorName(err)});
 144                return err;
 145            };
 146            if (n == 0) {
 147                std.log.info("server closed connection", .{});
 148                // Server closed connection
 149                return ClientResult{ .kind = .detach, .session_name = null };
 150            }
 151
 152            while (read_buf.next()) |msg| {
 153                switch (msg.header.tag) {
 154                    .Output => {
 155                        if (msg.payload.len > 0) {
 156                            try stdout_buf.appendSlice(gpa, msg.payload);
 157                        }
 158                    },
 159                    .Resize => {
 160                        // daemon is asking for the client's window size usually in response
 161                        // to this client being set as leader.
 162                        const next_size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO);
 163                        try ipc.appendSizeMessage(gpa, &sock_write_buf, .Resize, next_size);
 164                    },
 165                    .Switch => {
 166                        std.log.info("switch session", .{});
 167                        // Payload format: "session_name\ncwd" from the daemon
 168                        const newline_idx = std.mem.indexOfScalar(u8, msg.payload, '\n') orelse {
 169                            // No cwd provided (backward compat or old daemon)
 170                            return ClientResult{ .kind = .switch_session, .session_name = try gpa.dupe(u8, msg.payload) };
 171                        };
 172                        return ClientResult{
 173                            .kind = .switch_session,
 174                            .session_name = try gpa.dupe(u8, msg.payload[0..newline_idx]),
 175                            .cwd = if (newline_idx + 1 < msg.payload.len) try gpa.dupe(u8, msg.payload[newline_idx + 1 ..]) else null,
 176                        };
 177                    },
 178                    else => {},
 179                }
 180            }
 181        }
 182
 183        // Handle socket write (flush buffered messages to daemon)
 184        if (poll_fds.items[1].revents & lib_posix.POLL.OUT != 0) {
 185            if (sock_write_buf.items.len > 0) {
 186                const n = lib_posix.write(client_sock_fd, sock_write_buf.items) catch |err| blk: {
 187                    if (err == error.WouldBlock) break :blk 0;
 188                    if (err == error.ConnectionResetByPeer or err == error.BrokenPipe) {
 189                        std.log.info("connection reset or broken pipe", .{});
 190                        return ClientResult{ .kind = .detach, .session_name = null };
 191                    }
 192                    return err;
 193                };
 194                if (n > 0) {
 195                    try sock_write_buf.replaceRange(gpa, 0, n, &[_]u8{});
 196                }
 197            }
 198        }
 199
 200        if (stdout_buf.items.len > 0) {
 201            const n = lib_posix.write(lib_posix.STDOUT_FILENO, stdout_buf.items) catch |err| blk: {
 202                if (err == error.WouldBlock) break :blk 0;
 203                return err;
 204            };
 205            if (n > 0) {
 206                try stdout_buf.replaceRange(gpa, 0, n, &[_]u8{});
 207            }
 208        }
 209
 210        if (poll_fds.items[1].revents & (lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL) != 0) {
 211            std.log.info("poll hup|err|nval", .{});
 212            return ClientResult{ .kind = .detach, .session_name = null };
 213        }
 214    }
 215}
 216
 217fn initTerminal(gpa: std.mem.Allocator, io: std.Io, size: ipc.Resize, cfg: *const Cfg) !ghostty_vt.Terminal {
 218    return ghostty_vt.Terminal.init(io, gpa, .{
 219        .cols = size.cols,
 220        .rows = size.rows,
 221        .max_scrollback_lines = cfg.max_scrollback_lines,
 222        .max_scrollback_bytes = null, // Let the line limit control scrollback.
 223    });
 224}
 225
 226/// dameonLoop is what the daemon runs to send and receive ipc commands from its corresponding
 227/// clients.  It uses poll() as its non-blocking mechanism.
 228fn daemonLoop(daemon: *Daemon, gpa: std.mem.Allocator, io: std.Io, server_sock_fd: lib_posix.socket_t, pty_fd: i32) !void {
 229    std.log.info("daemon started session={s} pty_fd={d}", .{ daemon.session_name, pty_fd });
 230
 231    try signal.openSignalPipe();
 232    signal.installWakeHandler(@intFromEnum(lib_posix.SIG.TERM));
 233    var poll_fds = try std.ArrayList(lib_posix.pollfd).initCapacity(gpa, 8);
 234    defer poll_fds.deinit(gpa);
 235
 236    const init_size = ipc.getTerminalSize(pty_fd);
 237    var term = try initTerminal(gpa, io, init_size, daemon.cfg);
 238    defer term.deinit(gpa);
 239    var vt_stream = term.vtStream();
 240    defer vt_stream.deinit();
 241
 242    // Carries the tail of the previous PTY read so the task-exit marker
 243    // search below can see across a read() boundary. Sized to comfortably
 244    // hold "ZMX_TASK_COMPLETED:" (19 bytes) plus a u8 exit code and CRLF.
 245    var marker_carry: [32]u8 = undefined;
 246    var marker_carry_len: usize = 0;
 247
 248    var had_terminal_client = daemon.hasTerminalClient();
 249
 250    daemon_loop: while (daemon.running) {
 251        // If the program asked for focus reports (DECSET 1004), send focus-out
 252        // when the last attached client leaves and focus-in when one returns,
 253        // as a terminal would when its window loses/gains focus.
 254        const has_terminal_client = daemon.hasTerminalClient();
 255        if (has_terminal_client != had_terminal_client and term.modes.get(.focus_event)) {
 256            daemon.queuePtyInput(gpa, if (has_terminal_client) "\x1b[I" else "\x1b[O");
 257        }
 258        had_terminal_client = has_terminal_client;
 259
 260        poll_fds.clearRetainingCapacity();
 261
 262        try poll_fds.append(gpa, .{
 263            .fd = server_sock_fd,
 264            .events = lib_posix.POLL.IN,
 265            .revents = 0,
 266        });
 267
 268        var pty_events: i16 = lib_posix.POLL.IN;
 269        if (daemon.pty_write_buf.items.len > 0) {
 270            pty_events |= lib_posix.POLL.OUT;
 271        }
 272        try poll_fds.append(gpa, .{
 273            .fd = pty_fd,
 274            .events = pty_events,
 275            .revents = 0,
 276        });
 277
 278        try poll_fds.append(gpa, .{ .fd = signal.sig_pipe[0], .events = lib_posix.POLL.IN, .revents = 0 });
 279
 280        for (daemon.clients.items) |client| {
 281            var events: i16 = lib_posix.POLL.IN;
 282            if (client.has_pending_output) {
 283                events |= lib_posix.POLL.OUT;
 284            }
 285            try poll_fds.append(gpa, .{
 286                .fd = client.socket_fd,
 287                .events = events,
 288                .revents = 0,
 289            });
 290        }
 291
 292        _ = try lib_posix.poll(poll_fds.items, -1);
 293
 294        if (poll_fds.items[2].revents & lib_posix.POLL.IN != 0) {
 295            signal.drainSignalPipe();
 296            std.log.info(
 297                "SIGTERM received, shutting down gracefully session={s}",
 298                .{daemon.session_name},
 299            );
 300            break :daemon_loop;
 301        }
 302
 303        if (poll_fds.items[0].revents & (lib_posix.POLL.ERR | lib_posix.POLL.HUP | lib_posix.POLL.NVAL) != 0) {
 304            std.log.err("server socket error revents={d}", .{poll_fds.items[0].revents});
 305            break :daemon_loop;
 306        } else if (poll_fds.items[0].revents & lib_posix.POLL.IN != 0) {
 307            const client_fd = try lib_posix.accept(
 308                server_sock_fd,
 309                null,
 310                null,
 311                lib_posix.SOCK.NONBLOCK | lib_posix.SOCK.CLOEXEC,
 312            );
 313            const client = try gpa.create(Client);
 314            client.* = Client{
 315                .alloc = gpa,
 316                .socket_fd = client_fd,
 317                .read_buf = try ipc.SocketBuffer.init(gpa),
 318                .write_buf = undefined,
 319            };
 320            // 64KB initial capacity lets ~15 broadcast cycles (N_TTY_BUF_SIZE reads
 321            // * header) accumulate before the first ArrayList growth. The write
 322            // buffer is userspace-only: it drains via POLLOUT to the client socket,
 323            // which has no corresponding kernel-imposed per-write limit.
 324            client.write_buf = try std.ArrayList(u8).initCapacity(client.alloc, 65536);
 325            try daemon.clients.append(gpa, client);
 326            std.log.info(
 327                "client connected fd={d} total={d}",
 328                .{ client_fd, daemon.clients.items.len },
 329            );
 330        }
 331
 332        const inp_flags = lib_posix.POLL.IN | lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL;
 333        if (poll_fds.items[1].revents & inp_flags != 0) {
 334            // Read from PTY. Buffer is sized to N_TTY_BUF_SIZE (4096): the hard
 335            // kernel limit for the N_TTY line discipline. A larger buffer doesn't
 336            // help: each read() from a PTY master returns at most 4096 bytes
 337            // regardless of the userspace buffer size.
 338            var buf: [4096]u8 = undefined;
 339            const n_opt: ?usize = lib_posix.read(pty_fd, &buf) catch |err| blk: {
 340                if (err == error.WouldBlock) break :blk null;
 341                break :blk 0;
 342            };
 343
 344            if (n_opt) |n| {
 345                if (n == 0) {
 346                    // EOF: Shell exited
 347                    std.log.info("shell exited pty_fd={d}", .{pty_fd});
 348                    // Let the rest of this poll iteration complete so client
 349                    // write buffers are flushed via the normal POLLOUT path.
 350                    // On the next iteration, daemon.running will be false.
 351                    daemon.running = false;
 352                } else {
 353                    // Feed PTY output to terminal emulator for state tracking
 354                    vt_stream.nextSlice(buf[0..n]);
 355                    daemon.setPwd(&term);
 356                    daemon.has_pty_output = true;
 357
 358                    // When no real terminal client has attached yet, respond to
 359                    // terminal queries (e.g. DA1/DA2) on behalf of the terminal.
 360                    // This prevents fish from waiting 10s for unanswered queries.
 361                    // Only clients that sent .Init (a real zmx attach) count,
 362                    // not a `zmx run` tail-only client.
 363                    if (!daemon.hasTerminalClient() and
 364                        daemon.pty_write_buf.items.len < Daemon.PTY_WRITE_BUF_MAX)
 365                    {
 366                        util.respondToDeviceAttributes(gpa, &daemon.pty_write_buf, buf[0..n]);
 367                    }
 368
 369                    // In run mode, scan output for exit code marker. The marker
 370                    // can straddle two PTY reads (more likely under a throttled
 371                    // scheduler, e.g. containers), so prepend the tail carried
 372                    // over from the previous read before searching.
 373                    if (daemon.is_task_mode and daemon.task_exit_code == null) {
 374                        var scan_buf: [marker_carry.len + buf.len]u8 = undefined;
 375                        @memcpy(scan_buf[0..marker_carry_len], marker_carry[0..marker_carry_len]);
 376                        @memcpy(scan_buf[marker_carry_len..][0..n], buf[0..n]);
 377                        const scan_len = marker_carry_len + n;
 378
 379                        if (try util.findTaskExitMarker(scan_buf[0..scan_len], daemon.task_id)) |exit_code| {
 380                            daemon.task_exit_code = exit_code;
 381                            daemon.task_ended_at = @intCast(std.Io.Timestamp.now(io, .real).toSeconds());
 382
 383                            std.log.info("task completed exit_code={d}", .{exit_code});
 384
 385                            // Notify connected clients
 386                            for (daemon.clients.items) |c| {
 387                                ipc.appendMessage(gpa, &c.write_buf, .TaskComplete, &[_]u8{exit_code}) catch {};
 388                                c.has_pending_output = true;
 389                            }
 390                        }
 391
 392                        marker_carry_len = @min(marker_carry.len, scan_len);
 393                        @memcpy(
 394                            marker_carry[0..marker_carry_len],
 395                            scan_buf[scan_len - marker_carry_len .. scan_len],
 396                        );
 397                    }
 398
 399                    // Broadcast data to all clients.
 400                    // Rewrite OSC 133;A to include redraw=0 so the outer terminal
 401                    // does not clear prompt lines on resize (issue #111).
 402                    const broadcast_data = util.rewritePromptRedraw(gpa, buf[0..n]) orelse buf[0..n];
 403                    defer if (broadcast_data.ptr != buf[0..n].ptr) gpa.free(broadcast_data);
 404                    for (daemon.clients.items) |client| {
 405                        ipc.appendMessage(gpa, &client.write_buf, .Output, broadcast_data) catch |err| {
 406                            std.log.warn(
 407                                "failed to buffer output for client err={s}",
 408                                .{@errorName(err)},
 409                            );
 410                            continue;
 411                        };
 412                        client.has_pending_output = true;
 413                    }
 414                }
 415            }
 416        }
 417
 418        if (poll_fds.items[1].revents & lib_posix.POLL.OUT != 0) {
 419            while (daemon.pty_write_buf.items.len > 0) {
 420                const n = lib_posix.write(pty_fd, daemon.pty_write_buf.items) catch |err| {
 421                    if (err != error.WouldBlock) {
 422                        std.log.warn("pty write failed: {s}", .{@errorName(err)});
 423                        daemon.pty_write_buf.clearRetainingCapacity();
 424                    }
 425                    break;
 426                };
 427                if (n == 0) break;
 428                daemon.pty_write_buf.replaceRange(gpa, 0, n, &[_]u8{}) catch unreachable;
 429            }
 430        }
 431
 432        var i: usize = daemon.clients.items.len;
 433        // Only iterate over clients that were present when poll_fds was constructed
 434        // poll_fds contains [server, pty, sig_pipe, client0, client1, ...]
 435        // So number of clients in poll_fds is poll_fds.items.len - 3
 436        const num_polled_clients = poll_fds.items.len - 3;
 437        if (i > num_polled_clients) {
 438            // If we have more clients than polled (i.e. we just accepted one), start from the
 439            // polled ones
 440            i = num_polled_clients;
 441        }
 442
 443        clients_loop: while (i > 0) {
 444            i -= 1;
 445            const client = daemon.clients.items[i];
 446            const revents = poll_fds.items[i + 3].revents;
 447
 448            if (revents & lib_posix.POLL.IN != 0) {
 449                const n = client.read_buf.read(client.socket_fd) catch |err| {
 450                    if (err == error.WouldBlock) continue;
 451                    std.log.debug(
 452                        "client read err={s} fd={d}",
 453                        .{ @errorName(err), client.socket_fd },
 454                    );
 455                    const last = daemon.closeClient(gpa, client, i, false);
 456                    if (last) break :daemon_loop;
 457                    continue;
 458                };
 459
 460                if (n == 0) {
 461                    // Client closed connection
 462                    const last = daemon.closeClient(gpa, client, i, false);
 463                    if (last) break :daemon_loop;
 464                    continue;
 465                }
 466
 467                while (client.read_buf.next()) |msg| {
 468                    switch (msg.header.tag) {
 469                        .Input => try daemon.handleInput(gpa, client, msg.payload),
 470                        .Send => daemon.handleSend(gpa, msg.payload),
 471                        .Output => try daemon.handleOutput(gpa, msg.payload, &term, &vt_stream),
 472                        .Init => try daemon.handleInit(gpa, client, pty_fd, &term, msg.payload),
 473                        .Switch => try daemon.handleSwitch(gpa, msg.payload),
 474                        .Resize => try daemon.handleResize(gpa, client, pty_fd, &term, msg.payload),
 475                        .Detach => {
 476                            daemon.handleDetach(gpa, client, i);
 477                            break :clients_loop;
 478                        },
 479                        .DetachAll => {
 480                            daemon.handleDetachAll(gpa);
 481                            break :clients_loop;
 482                        },
 483                        .Kill => {
 484                            break :daemon_loop;
 485                        },
 486                        .Info => try daemon.handleInfo(gpa, client, &term),
 487                        .LabelGet => try daemon.handleLabelGet(gpa, client),
 488                        .LabelSet => try daemon.handleLabelSet(gpa, client, msg.payload),
 489                        .LabelClear => try daemon.handleLabelClear(gpa, client),
 490                        .EnvGet => try daemon.handleEnvGet(gpa, client),
 491                        .EnvSet => try daemon.handleEnvSet(gpa, client, msg.payload),
 492                        .History => try daemon.handleHistory(gpa, client, &term, msg.payload),
 493                        .Run => try daemon.handleRun(gpa, io, client, msg.payload),
 494                        .Ack, .TaskComplete, .LabelData, .EnvData => {},
 495                        .Write => try daemon.handleWrite(gpa, client, msg.payload),
 496                        _ => std.log.warn(
 497                            "ignoring unknown IPC tag={d}",
 498                            .{@intFromEnum(msg.header.tag)},
 499                        ),
 500                    }
 501                }
 502            }
 503
 504            if (revents & lib_posix.POLL.OUT != 0) {
 505                // Flush pending output buffers
 506                const n = lib_posix.write(client.socket_fd, client.write_buf.items) catch |err| blk: {
 507                    if (err == error.WouldBlock) break :blk 0;
 508                    // Error on write, close client
 509                    const last = daemon.closeClient(gpa, client, i, false);
 510                    if (last) break :daemon_loop;
 511                    continue;
 512                };
 513
 514                if (n > 0) {
 515                    client.write_buf.replaceRange(gpa, 0, n, &[_]u8{}) catch unreachable;
 516                }
 517
 518                if (client.write_buf.items.len == 0) {
 519                    client.has_pending_output = false;
 520                }
 521            }
 522
 523            if (revents & (lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL) != 0) {
 524                const last = daemon.closeClient(gpa, client, i, false);
 525                if (last) break :daemon_loop;
 526            }
 527        }
 528    }
 529}
 530
 531const ClientResult = struct {
 532    kind: enum {
 533        detach,
 534        switch_session,
 535    },
 536    session_name: ?[]const u8,
 537    cwd: ?[]const u8 = null,
 538};
 539
 540/// Client represents each terminal that has connected to a session.
 541///
 542/// Multiple Clients can connect to a single session.
 543pub const Client = struct {
 544    alloc: std.mem.Allocator,
 545    socket_fd: i32,
 546    has_pending_output: bool = false,
 547    is_terminal: bool = false, // sent .Init (a `zmx attach`), not a run/send/tail client
 548    read_buf: ipc.SocketBuffer,
 549    write_buf: std.ArrayList(u8),
 550    env_str: ?[]u8 = null,
 551
 552    pub fn deinit(self: *Client, gpa: std.mem.Allocator) void {
 553        lib_posix.close(self.socket_fd);
 554        self.read_buf.deinit();
 555        self.write_buf.deinit(self.alloc);
 556        if (self.env_str) |s| gpa.free(s);
 557    }
 558
 559    fn setEnv(self: *Client, gpa: std.mem.Allocator, env_str: []const u8) !void {
 560        if (self.env_str) |s| gpa.free(s);
 561        self.env_str = if (env_str.len > 0) try gpa.dupe(u8, env_str) else null;
 562    }
 563};
 564
 565/// Daemon is responsible for managing a zmx session.
 566///
 567/// It holds all the state for a running session.  Instead of a single daemon for all sessions, we
 568/// create a daemon for every session.  This has some benefits. The ipc communication between
 569/// session clients and the daemon doesn't need to be tagged with the session name.  If a daemon
 570/// crashes for one session won't crash all the other sessions.
 571///
 572/// Conceptually it's also much simpler to reason about.
 573pub const Daemon = struct {
 574    cfg: *Cfg,
 575    session_name: []const u8,
 576    socket_path: []const u8,
 577    // === opt ===
 578    pty_write_buf: std.ArrayList(u8) = .empty,
 579    clients: std.ArrayList(*Client) = .empty,
 580    labels: std.StringHashMapUnmanaged([]u8) = .empty,
 581    // This control which client is the leader.  The leader controls terminal state and
 582    // cols/rows of session.
 583    leader_client_fd: ?i32 = null,
 584    running: bool = true,
 585    pid: i32 = undefined,
 586    command: ?[]const []const u8 = null,
 587    /// The session's working directory in OSC 7 form, `file://<host><path>`.
 588    /// Kept as a URI rather than a path so `zmx list` shows the host, which is
 589    /// what tells you a session is inside SSH. Points into `cwd_buf` once set,
 590    /// so a Daemon must not be copied by value after that.
 591    cwd: []const u8 = "",
 592    /// The same directory as a path that can be opened: percent-decoding
 593    /// applied, scheme and host stripped. Empty when the cwd is on another
 594    /// host, since then it names no directory here and nothing should chdir
 595    /// into it. Points into `cwd_path_buf`.
 596    cwd_path: []const u8 = "",
 597    cwd_buf: [std.fs.max_path_bytes]u8 = undefined,
 598    cwd_path_buf: [std.fs.max_path_bytes]u8 = undefined,
 599    has_pty_output: bool = false,
 600    has_had_client: bool = false,
 601    created_at: u64, // unix timestamp (ns)
 602    is_task_mode: bool = false, // flag for when session is run as a task
 603    task_id: [4]u8 = undefined,
 604    task_exit_code: ?u8 = null, // null = running or n/a, set when task completes
 605    task_ended_at: ?u64 = null, // timestamp when task exited
 606    pty_fd: i32 = -1, // set by daemonLoop so handleRun can probe the foreground process
 607    shell: []const u8 = "/bin/sh",
 608
 609    /// Create a Daemon. Caller is responsible for freeing all variables passed
 610    /// into the init fn.
 611    pub fn init(io: std.Io, cfg: *Cfg, sesh_name: []const u8, socket_path: []const u8) Daemon {
 612        return .{
 613            .cfg = cfg,
 614            .session_name = sesh_name,
 615            .socket_path = socket_path,
 616            .created_at = @intCast(std.Io.Timestamp.now(io, .real).toSeconds()),
 617        };
 618    }
 619
 620    pub fn deinit(self: *Daemon, gpa: std.mem.Allocator) void {
 621        self.clients.deinit(gpa);
 622        var it = self.labels.iterator();
 623        while (it.next()) |entry| {
 624            gpa.free(entry.key_ptr.*);
 625            gpa.free(entry.value_ptr.*);
 626        }
 627        self.labels.deinit(gpa);
 628        self.pty_write_buf.deinit(gpa);
 629        gpa.free(self.socket_path);
 630    }
 631
 632    pub fn shutdown(self: *Daemon, gpa: std.mem.Allocator) void {
 633        std.log.info("shutting down daemon session={s}", .{self.session_name});
 634        self.running = false;
 635
 636        for (self.clients.items) |client| {
 637            client.deinit(gpa);
 638            gpa.destroy(client);
 639        }
 640        self.clients.clearRetainingCapacity();
 641    }
 642
 643    /// Resize the daemon's terminal with prompt_redraw disabled. On resize the
 644    /// terminal would clear prompt lines expecting the shell to redraw them,
 645    /// but the shell's redraw goes to the PTY (forwarded to clients), not to
 646    /// this terminal, so the clearing only corrupts our snapshot state.
 647    fn resizeTerm(gpa: std.mem.Allocator, term: *ghostty_vt.Terminal, cols: u16, rows: u16) !void {
 648        const saved = term.flags.shell_redraws_prompt;
 649        term.flags.shell_redraws_prompt = .false;
 650        defer term.flags.shell_redraws_prompt = saved;
 651        try term.resize(gpa, .{ .cols = cols, .rows = rows });
 652    }
 653
 654    /// True while a client that sent .Init (a real `zmx attach`) is connected.
 655    fn hasTerminalClient(self: *const Daemon) bool {
 656        for (self.clients.items) |c| {
 657            if (c.is_terminal) return true;
 658        }
 659        return false;
 660    }
 661
 662    pub fn closeClient(self: *Daemon, gpa: std.mem.Allocator, client: *Client, i: usize, shutdown_on_last: bool) bool {
 663        const fd = client.socket_fd;
 664        // leader is disconnected, remove ref and let another client claim leader on input
 665        if (self.leader_client_fd == client.socket_fd) {
 666            std.log.info(
 667                "unsetting leader session={s} fd={d}",
 668                .{ self.session_name, client.socket_fd },
 669            );
 670            self.leader_client_fd = null;
 671        }
 672        client.deinit(gpa);
 673        gpa.destroy(client);
 674        _ = self.clients.orderedRemove(i);
 675        std.log.info("client disconnected fd={d} remaining={d}", .{ fd, self.clients.items.len });
 676        if (shutdown_on_last and self.clients.items.len == 0) {
 677            self.shutdown(gpa);
 678            return true;
 679        }
 680        return false;
 681    }
 682
 683    /// ensureSession will either create or re-use the daemon used for a session.
 684    /// It will spin up a unix socket, double-fork the process (so it survives
 685    /// the terminal dying), and automatically attach the client to the ipc unix
 686    /// socket.
 687    ///
 688    /// The return bool value indicates if the current process is the daemon
 689    /// or the client since they have different behaviors post-fork.
 690    ///
 691    /// E.g. If it's the client process then we need to connect to the unix socket
 692    /// and run the clientLoop.  If it's the daemon then we need to bail since
 693    /// the daemonLoop is created inside this fn and when it returns that means
 694    /// the daemon stopped and needs to exit.
 695    pub fn ensureSession(self: *Daemon, io: std.Io) !bool {
 696        const sesh_name = self.session_name;
 697        std.log.info("ensure session session={s}", .{sesh_name});
 698        var dir = try std.Io.Dir.openDirAbsolute(io, self.cfg.socket_dir, .{});
 699        defer dir.close(io);
 700
 701        const exists = try socket.sessionExists(io, dir, sesh_name);
 702        // if daemon is gone then we flip this to true
 703        var should_create = !exists;
 704
 705        if (exists) {
 706            if (ipc.connectSession(self.socket_path)) |fd| {
 707                lib_posix.close(fd);
 708                if (self.command != null) {
 709                    std.log.warn(
 710                        "session already exists, ignoring command session={s}",
 711                        .{sesh_name},
 712                    );
 713                }
 714            } else |err| switch (err) {
 715                // Daemon is definitively gone: safe to replace.
 716                error.ConnectionRefused => {
 717                    socket.cleanupStaleSocket(io, dir, sesh_name);
 718                    should_create = true;
 719                },
 720                // Connect failed for an unusual reason. The check is only to
 721                // decide create-vs-attach; the socket file exists, so proceed
 722                // to attach rather than fail or orphan.
 723                else => {
 724                    std.log.warn(
 725                        "connect failed ({s}), proceeding to attach session={s}",
 726                        .{ @errorName(err), sesh_name },
 727                    );
 728                },
 729            }
 730        }
 731
 732        if (!should_create) {
 733            return false;
 734        }
 735
 736        return self.run(io, dir, sesh_name);
 737    }
 738
 739    fn run(self: *Daemon, io: std.Io, dir: std.Io.Dir, sesh_name: []const u8) !bool {
 740        std.log.info("creating session={s}", .{sesh_name});
 741        const server_sock_fd: lib_posix.socket_t = try socket.createSocket(self.socket_path);
 742        const log_fd = log.log_system.file.?.handle;
 743
 744        var keep_fds_open = [_]i32{ server_sock_fd, dir.handle, log_fd };
 745        const cmd = try daemonize.createCmdZ(self.shell, self.is_task_mode, self.command);
 746
 747        // `cwd_path` is the decoded path, and is empty when the cwd is on
 748        // another host: OSC 7 crosses SSH boundaries, so a session that ssh'd
 749        // elsewhere reports a directory that does not exist on this machine.
 750        std.log.info("checking pwd={s} path={s}", .{ self.cwd, self.cwd_path });
 751        if (self.cwd_path.len > 0) {
 752            const pwd_dir = std.Io.Dir.openDirAbsolute(io, self.cwd_path, .{}) catch |err| blk: {
 753                std.log.warn("failed to open dir={s} err={s}", .{ self.cwd_path, @errorName(err) });
 754                break :blk null;
 755            };
 756            if (pwd_dir) |pdir| {
 757                defer std.Io.Dir.close(pdir, io);
 758                std.log.info("set directory dir={s}", .{self.cwd_path});
 759                try std.process.setCurrentDir(io, pdir);
 760            }
 761        }
 762
 763        const pty_info = daemonize.daemonize(
 764            sesh_name,
 765            cmd,
 766            &keep_fds_open,
 767        ) catch |err| {
 768            switch (err) {
 769                error.IsClientProc => {
 770                    // send a msg to the client that the session was created.
 771                    var w_buf: [2048]u8 = undefined;
 772                    var w = std.Io.File.stdout().writer(io, &w_buf);
 773                    try w.interface.print("session \"{s}\" created\n", .{sesh_name});
 774                    try w.interface.flush();
 775                    lib_posix.close(server_sock_fd);
 776                    return false;
 777                },
 778                else => {
 779                    lib_posix.close(server_sock_fd);
 780                    dir.deleteFile(io, self.session_name) catch {};
 781                    return err;
 782                },
 783            }
 784        };
 785        // =======
 786        // WARNING: cannot use upstream allocator or io after this point since
 787        // we forked the process and there's a risk of a mutex (e.g. thread-safe
 788        // allocator) being locked by a thread prior to fork which can cause a
 789        // deadlock.
 790        // =======
 791
 792        self.pid = pty_info.pid;
 793
 794        var threaded: std.Io.Threaded = .init_single_threaded;
 795        defer threaded.deinit();
 796        const new_io = threaded.io();
 797
 798        { // re-initialize logs with the session name as the filename
 799            log.log_system.deinit();
 800            var log_buf: [4096]u8 = undefined;
 801            const session_log_name = try std.fmt.bufPrint(
 802                &log_buf,
 803                "{s}.log",
 804                .{sesh_name},
 805            );
 806            var fba_buf: [4096]u8 = undefined;
 807            var fba = std.heap.FixedBufferAllocator.init(&fba_buf);
 808            const session_log_path = try std.fs.path.join(
 809                fba.allocator(),
 810                &.{ self.cfg.log_dir, session_log_name },
 811            );
 812            const log_mode = std.Io.File.Permissions.fromMode(@intCast(self.cfg.log_mode));
 813            log.log_system.init(new_io, session_log_path, log_mode) catch {};
 814        }
 815
 816        const gpa: std.mem.Allocator = blk: {
 817            if (builtin.mode == .Debug) {
 818                const GPA = std.heap.DebugAllocator(.{});
 819                const Static = struct {
 820                    var gpa: GPA = .{};
 821                };
 822                break :blk Static.gpa.allocator();
 823            }
 824            break :blk std.heap.c_allocator;
 825        };
 826
 827        defer {
 828            // Close and unlink the listen socket BEFORE handleKill()'s
 829            // 500ms SIGHUP->SIGKILL grace sleep. Otherwise a `zmx run`
 830            // for the same name issued in that window will hang waiting
 831            // for a connect.
 832            lib_posix.close(server_sock_fd);
 833            std.log.info("deleting socket file session={s}", .{sesh_name});
 834            dir.deleteFile(new_io, sesh_name) catch |err| {
 835                std.log.warn("failed to delete socket file err={s}", .{@errorName(err)});
 836            };
 837            self.handleKill(gpa, new_io);
 838            self.deinit(gpa);
 839            lib_posix.close(pty_info.master_fd);
 840            _ = lib_posix.waitpid(self.pid, 0);
 841        }
 842
 843        try daemonLoop(self, gpa, new_io, server_sock_fd, pty_info.master_fd);
 844        std.log.info("daemon loop shutdown", .{});
 845        return true;
 846    }
 847
 848    fn setLeader(self: *Daemon, gpa: std.mem.Allocator, client: *Client) !void {
 849        std.log.info("setting new leader client_fd={d}", .{client.socket_fd});
 850        self.leader_client_fd = client.socket_fd;
 851        // Send a resize message to the client so it can send us back their window size
 852        // so we can resize the pty and ghostty state.
 853        try ipc.appendMessage(gpa, &client.write_buf, .Resize, "");
 854        client.has_pending_output = true;
 855    }
 856
 857    fn getLeaderClient(self: *Daemon) ?*Client {
 858        for (self.clients.items) |client| {
 859            if (self.leader_client_fd == client.socket_fd) {
 860                return client;
 861            }
 862        }
 863        return null;
 864    }
 865
 866    const PTY_WRITE_BUF_MAX = 256 * 1024;
 867
 868    /// Queue bytes for the PTY's stdin. Flushed by daemonLoop on POLLOUT.
 869    /// Drops the payload if the buffer is over cap -- same failure mode as
 870    /// the old direct-write ptyWrite (drop on EAGAIN), just at a 64x higher
 871    /// threshold. Capping avoids OOM when the shell stops reading; dropping
 872    /// new (not old) bytes avoids tearing a partially-accepted sequence.
 873    fn queuePtyInput(self: *Daemon, gpa: std.mem.Allocator, data: []const u8) void {
 874        if (data.len == 0) return;
 875        if (self.pty_write_buf.items.len + data.len > PTY_WRITE_BUF_MAX) {
 876            std.log.warn(
 877                "pty input dropped {d} bytes (buffer full, shell not reading)",
 878                .{data.len},
 879            );
 880            return;
 881        }
 882
 883        // NOTE: for local dev only
 884        // std.log.debug("buffering pty input data={x}", .{data});
 885
 886        self.pty_write_buf.appendSlice(gpa, data) catch |err| {
 887            std.log.warn(
 888                "pty input dropped {d} bytes: {s}",
 889                .{ data.len, @errorName(err) },
 890            );
 891        };
 892    }
 893
 894    pub fn handleInput(self: *Daemon, gpa: std.mem.Allocator, client: *Client, payload: []const u8) !void {
 895        // NOTE: for local dev only
 896        // std.log.debug("buffering pty input data={x}", .{payload});
 897
 898        // client is leader, send entire payload (ansi escape codes + text)
 899        if (self.leader_client_fd == client.socket_fd) {
 900            self.queuePtyInput(gpa, payload);
 901            return;
 902        }
 903
 904        // check if leader needs to be updated by detecting any user input
 905        if (util.isUserInput(payload)) {
 906            try self.setLeader(gpa, client);
 907            self.queuePtyInput(gpa, payload);
 908        }
 909    }
 910
 911    /// Queue input from `zmx send` without changing interactive client leadership.
 912    pub fn handleSend(self: *Daemon, gpa: std.mem.Allocator, payload: []const u8) void {
 913        self.queuePtyInput(gpa, payload);
 914    }
 915
 916    pub fn handleSwitch(self: *Daemon, gpa: std.mem.Allocator, session_name: []const u8) !void {
 917        for (self.clients.items) |client| {
 918            if (self.leader_client_fd == client.socket_fd) {
 919                // Include the daemon's current cwd so the new session can start
 920                // in the right directory. A remote cwd is left out: it names no
 921                // directory here, so the new session is better off with the
 922                // attaching client's own cwd than with a path it cannot enter.
 923                if (self.cwd.len > 0 and self.cwd_path.len > 0) {
 924                    var payload = gpa.alloc(u8, session_name.len + 1 + self.cwd.len) catch return;
 925                    defer gpa.free(payload);
 926                    @memcpy(payload[0..session_name.len], session_name);
 927                    payload[session_name.len] = '\n';
 928                    @memcpy(payload[session_name.len + 1 ..], self.cwd);
 929                    ipc.appendMessage(gpa, &client.write_buf, .Switch, payload) catch |err| {
 930                        std.log.warn(
 931                            "failed to buffer terminal state for client err={s}",
 932                            .{@errorName(err)},
 933                        );
 934                    };
 935                } else {
 936                    ipc.appendMessage(gpa, &client.write_buf, .Switch, session_name) catch |err| {
 937                        std.log.warn(
 938                            "failed to buffer terminal state for client err={s}",
 939                            .{@errorName(err)},
 940                        );
 941                    };
 942                }
 943                client.has_pending_output = true;
 944                return;
 945            }
 946        }
 947    }
 948
 949    pub fn handleInit(
 950        self: *Daemon,
 951        gpa: std.mem.Allocator,
 952        client: *Client,
 953        pty_fd: i32,
 954        term: *ghostty_vt.Terminal,
 955        payload: []const u8,
 956    ) !void {
 957        if (payload.len != @sizeOf(ipc.Resize)) return;
 958
 959        client.is_terminal = true;
 960
 961        if (self.leader_client_fd == null) {
 962            try self.setLeader(gpa, client);
 963        }
 964        const is_leader = self.leader_client_fd == client.socket_fd;
 965        const resize = std.mem.bytesToValue(ipc.Resize, payload);
 966
 967        // Resize our terminal (not yet the PTY) to the leader's size before
 968        // serializing, so the snapshot is laid out for the width the client
 969        // will render it at instead of being wrapped a second time on arrival.
 970        // Cursor position stays consistent because it is serialized from the
 971        // same, already-resized terminal; the PTY is resized below, so the
 972        // shell's own SIGWINCH redraw still arrives after the snapshot.
 973        if (is_leader) {
 974            try resizeTerm(gpa, term, resize.cols, resize.rows);
 975        }
 976
 977        // Only serialize on re-attach (has_had_client), not first attach, to avoid
 978        // interfering with shell initialization (DA1 queries, etc.)
 979        if (self.has_pty_output and self.has_had_client) {
 980            const cursor = &term.screens.active.cursor;
 981            std.log.debug(
 982                "cursor before serialize: x={d} y={d} pending_wrap={}",
 983                .{ cursor.x, cursor.y, cursor.pending_wrap },
 984            );
 985            if (util.serializeTerminalState(gpa, term)) |term_output| {
 986                std.log.debug("serialize terminal state", .{});
 987                // Rewrite OSC 133;A to include redraw=0 so the outer terminal
 988                // does not clear prompt lines on resize (issue #111).
 989                const restore_data = util.rewritePromptRedraw(gpa, term_output) orelse term_output;
 990                defer gpa.free(term_output);
 991                defer if (restore_data.ptr != term_output.ptr) gpa.free(restore_data);
 992                ipc.appendMessage(gpa, &client.write_buf, .Output, restore_data) catch |err| {
 993                    std.log.warn(
 994                        "failed to buffer terminal state for client err={s}",
 995                        .{@errorName(err)},
 996                    );
 997                };
 998                client.has_pending_output = true;
 999            }
1000        }
1001
1002        // only resize if leader
1003        if (is_leader) {
1004            var ws = resize.winsize();
1005            _ = cross.c.ioctl(pty_fd, cross.c.TIOCSWINSZ, &ws);
1006
1007            // On re-attach, deliver SIGWINCH to the foreground process group so
1008            // incremental renderers (Ink, Claude Code, etc.) know to repaint.
1009            // If the size changed, TIOCSWINSZ above already sent SIGWINCH; if the size
1010            // was unchanged, the kernel suppressed it, so signal the pgrp explicitly.
1011            if (self.has_pty_output and self.has_had_client) {
1012                var pgrp: lib_posix.pid_t = 0;
1013                if (cross.c.ioctl(pty_fd, cross.c.TIOCGPGRP, &pgrp) == 0 and pgrp > 0) {
1014                    lib_posix.kill(-pgrp, .WINCH) catch {};
1015                }
1016            }
1017
1018            // Mark that we've had a client init, so subsequent clients get terminal state
1019            self.has_had_client = true;
1020
1021            std.log.debug("init resize rows={d} cols={d}", .{ resize.rows, resize.cols });
1022        }
1023    }
1024
1025    pub fn handleResize(
1026        self: *Daemon,
1027        gpa: std.mem.Allocator,
1028        client: *Client,
1029        pty_fd: i32,
1030        term: *ghostty_vt.Terminal,
1031        payload: []const u8,
1032    ) !void {
1033        if (payload.len != @sizeOf(ipc.Resize)) return;
1034        if (self.leader_client_fd == null) {
1035            try self.setLeader(gpa, client);
1036        }
1037        // only leader can resize
1038        if (self.leader_client_fd != client.socket_fd) return;
1039
1040        const resize = std.mem.bytesToValue(ipc.Resize, payload);
1041        var ws = resize.winsize();
1042        _ = cross.c.ioctl(pty_fd, cross.c.TIOCSWINSZ, &ws);
1043        try resizeTerm(gpa, term, resize.cols, resize.rows);
1044        std.log.debug("resize rows={d} cols={d}", .{ resize.rows, resize.cols });
1045    }
1046
1047    pub fn handleDetach(self: *Daemon, gpa: std.mem.Allocator, client: *Client, i: usize) void {
1048        std.log.info("client detach session={s} fd={d}", .{ self.session_name, client.socket_fd });
1049        _ = self.closeClient(gpa, client, i, false);
1050    }
1051
1052    pub fn handleDetachAll(self: *Daemon, gpa: std.mem.Allocator) void {
1053        std.log.info("detach all clients={d}", .{self.clients.items.len});
1054        // Go through closeClient so the leader is cleared like any other detach.
1055        while (self.clients.items.len > 0) {
1056            const last = self.clients.items.len - 1;
1057            _ = self.closeClient(gpa, self.clients.items[last], last, false);
1058        }
1059    }
1060
1061    pub fn handleKill(self: *Daemon, gpa: std.mem.Allocator, io: std.Io) void {
1062        std.log.info("kill received session={s}", .{self.session_name});
1063        self.shutdown(gpa);
1064        // gracefully shutdown shell processes, shells tend to ignore SIGTERM so we send SIGHUP
1065        // instead
1066        //   https://www.gnu.org/software/bash/manual/html_node/Signals.html
1067        // negative pid means kill process and children
1068        std.log.info("sending SIGHUP session={s} pid={d}", .{ self.session_name, self.pid });
1069        lib_posix.kill(-self.pid, lib_posix.SIG.HUP) catch |err| {
1070            std.log.warn("failed to send SIGHUP to pty child err={s}", .{@errorName(err)});
1071        };
1072        std.Io.sleep(io, std.Io.Duration.fromMilliseconds(500), .real) catch unreachable;
1073        lib_posix.kill(-self.pid, lib_posix.SIG.KILL) catch |err| {
1074            std.log.warn("failed to send SIGKILL to pty child err={s}", .{@errorName(err)});
1075        };
1076    }
1077
1078    pub fn handleInfo(self: *Daemon, gpa: std.mem.Allocator, client: *Client, term: *ghostty_vt.Terminal) !void {
1079        self.setPwd(term);
1080
1081        // zeroes() so asBytes() doesn't ship struct padding + unused cmd/cwd
1082        // tail bytes (daemon stack contents) to clients.
1083        var info = std.mem.zeroes(ipc.Info);
1084        info.clients_len = self.clients.items.len - 1;
1085        info.pid = self.pid;
1086        info.created_at = self.created_at;
1087        info.task_ended_at = self.task_ended_at orelse 0;
1088        info.task_exit_code = self.task_exit_code orelse 0;
1089
1090        // Build command string from args, re-quoting args that contain
1091        // shell-special characters so the displayed command is copy-pasteable.
1092        const cur_cmd = self.command;
1093        if (cur_cmd) |args| {
1094            for (args, 0..) |arg, i| {
1095                const quoted = if (util.shellNeedsQuoting(arg))
1096                    util.shellQuote(gpa, arg) catch null
1097                else
1098                    null;
1099                defer if (quoted) |q| gpa.free(q);
1100                const src = quoted orelse arg;
1101
1102                const need = src.len + @as(usize, if (i > 0) 1 else 0);
1103                if (info.cmd_len + need > ipc.MAX_CMD_LEN) {
1104                    const ellipsis = "...";
1105                    if (info.cmd_len + ellipsis.len <= ipc.MAX_CMD_LEN) {
1106                        @memcpy(info.cmd[info.cmd_len..][0..ellipsis.len], ellipsis);
1107                        info.cmd_len += ellipsis.len;
1108                    }
1109                    break;
1110                }
1111
1112                if (i > 0) {
1113                    info.cmd[info.cmd_len] = ' ';
1114                    info.cmd_len += 1;
1115                }
1116                @memcpy(info.cmd[info.cmd_len..][0..src.len], src);
1117                info.cmd_len += @intCast(src.len);
1118            }
1119        }
1120
1121        info.cwd_len = @intCast(@min(self.cwd.len, ipc.MAX_CWD_LEN));
1122        @memcpy(info.cwd[0..info.cwd_len], self.cwd[0..info.cwd_len]);
1123
1124        try ipc.appendMessage(gpa, &client.write_buf, .Info, std.mem.asBytes(&info));
1125        client.has_pending_output = true;
1126    }
1127
1128    pub fn handleHistory(
1129        self: *Daemon,
1130        gpa: std.mem.Allocator,
1131        client: *Client,
1132        term: *ghostty_vt.Terminal,
1133        payload: []const u8,
1134    ) !void {
1135        self.setPwd(term);
1136        const format: util.HistoryFormat = if (payload.len > 0)
1137            @enumFromInt(payload[0])
1138        else
1139            .plain;
1140        if (util.serializeTerminal(gpa, term, format)) |output| {
1141            defer gpa.free(output);
1142            try ipc.appendMessage(gpa, &client.write_buf, .History, output);
1143            client.has_pending_output = true;
1144        } else {
1145            try ipc.appendMessage(gpa, &client.write_buf, .History, "");
1146            client.has_pending_output = true;
1147        }
1148    }
1149
1150    pub fn handleRun(self: *Daemon, gpa: std.mem.Allocator, io: std.Io, client: *Client, payload: []const u8) !void {
1151        // Reset task tracking so the new command's exit marker is detected.
1152        // Without this, a second `zmx run` on the same session is ignored
1153        // because task_exit_code is still set from the first run.
1154        self.task_exit_code = null;
1155        self.task_ended_at = null;
1156        self.is_task_mode = true;
1157        self.task_id = util.generateTaskId(io);
1158
1159        if (payload.len == 0) return;
1160
1161        const cmd = payload;
1162
1163        // Chain the exit marker with `;` on the same line. `$?` captures the
1164        // exit code of the command (not the `;`). The sole exception is when
1165        // the command contains a heredoc (`<<`), the delimiter must be alone
1166        // on its line, so the marker goes on the next line instead.
1167        var buf: [1024]u8 = undefined;
1168        const marker = try util.getTaskExitMarker(&buf, self.task_id);
1169        var single_buf: [1024]u8 = undefined;
1170        const single_line_marker = try std.fmt.bufPrint(&single_buf, "; echo {s}$?\r", .{marker});
1171        var here_buf: [1024]u8 = undefined;
1172        const heredoc_marker = try std.fmt.bufPrint(&here_buf, "\r\necho {s}$?\r", .{marker});
1173        const uses_heredoc = std.mem.indexOf(u8, cmd, "<<") != null;
1174
1175        if (cmd.len > 0 and cmd[cmd.len - 1] == '\r') {
1176            self.queuePtyInput(gpa, cmd[0 .. cmd.len - 1]);
1177        } else {
1178            self.queuePtyInput(gpa, cmd);
1179        }
1180        self.queuePtyInput(gpa, if (uses_heredoc) heredoc_marker else single_line_marker);
1181
1182        try ipc.appendMessage(gpa, &client.write_buf, .Ack, "");
1183        client.has_pending_output = true;
1184        self.has_had_client = true;
1185        std.log.debug("run command len={d}", .{payload.len});
1186    }
1187
1188    /// Store the session's working directory as a plain path.
1189    ///
1190    /// Accepts either an OSC 7 value (`file://<host><path>`, percent-encoded)
1191    /// or a path. Decoding here rather than at each use keeps `zmx list`
1192    /// printing a path and lets the chdir on session create find directories
1193    /// whose names needed escaping.
1194    ///
1195    /// The value is copied, so callers may pass a temporary.
1196    pub fn setCwd(self: *Daemon, value: []const u8) void {
1197        var buf: [std.fs.max_path_bytes]u8 = undefined;
1198        var host_buf: [std.posix.HOST_NAME_MAX]u8 = undefined;
1199        const hostname = std.posix.gethostname(&host_buf) catch "";
1200        const cwd = util.parseOsc7Cwd(&buf, value, hostname) orelse {
1201            std.log.warn("ignoring unusable cwd={s}", .{value});
1202            return;
1203        };
1204
1205        // Store the URI form. A caller that handed us a plain path gets one
1206        // built here, so `cwd` has the same shape no matter the source. A value
1207        // that already was a URI is kept verbatim, so `list` shows what the
1208        // shell actually reported.
1209        self.cwd = if (std.fs.path.isAbsolute(value))
1210            util.toOsc7Cwd(&self.cwd_buf, value, hostname) orelse return
1211        else blk: {
1212            if (value.len > self.cwd_buf.len) return;
1213            @memcpy(self.cwd_buf[0..value.len], value);
1214            break :blk self.cwd_buf[0..value.len];
1215        };
1216
1217        // Only keep an openable path when it names a directory on this host.
1218        if (cwd.is_local and cwd.path.len <= self.cwd_path_buf.len) {
1219            @memcpy(self.cwd_path_buf[0..cwd.path.len], cwd.path);
1220            self.cwd_path = self.cwd_path_buf[0..cwd.path.len];
1221        } else {
1222            self.cwd_path = "";
1223        }
1224        std.log.info("set cwd={s} path={s}", .{ self.cwd, self.cwd_path });
1225    }
1226
1227    fn setPwd(self: *Daemon, term: *ghostty_vt.Terminal) void {
1228        const pwd = term.getPwd() orelse return;
1229        if (std.mem.eql(u8, self.cwd, pwd)) return;
1230        self.setCwd(pwd);
1231    }
1232
1233    pub fn handleOutput(self: *Daemon, gpa: std.mem.Allocator, payload: []const u8, term: *ghostty_vt.Terminal, vt_stream: anytype) !void {
1234        vt_stream.nextSlice(payload);
1235        self.setPwd(term);
1236        self.has_pty_output = true;
1237        for (self.clients.items) |client| {
1238            try ipc.appendMessage(gpa, &client.write_buf, .Output, payload);
1239            client.has_pending_output = true;
1240        }
1241        if (self.clients.items.len > 0) {
1242            lib_posix.kill(self.pid, lib_posix.SIG.WINCH) catch |err| {
1243                std.log.warn("failed to send SIGWINCH err={s}", .{@errorName(err)});
1244            };
1245        }
1246    }
1247
1248    pub fn handleWrite(self: *Daemon, gpa: std.mem.Allocator, client: *Client, payload: []const u8) !void {
1249        // Wire format: [u32 path len][path bytes][file content]
1250        if (payload.len < @sizeOf(u32)) return error.InvalidPayload;
1251        const path_len = std.mem.bytesToValue(u32, payload[0..@sizeOf(u32)]);
1252        if (payload.len < @sizeOf(u32) + path_len) return error.InvalidPayload;
1253        const file_path = payload[@sizeOf(u32)..][0..path_len];
1254        const file_content = payload[@sizeOf(u32) + path_len ..];
1255
1256        // Inject file creation through the PTY so it works over SSH.
1257        // Base64-encode content and pipe through printf | base64 -d > file.
1258        // Chunk large files to stay under command-line length limits.
1259        // 48000 is divisible by 3 (clean base64 boundaries) and encodes
1260        // to ~64KB, well under typical ARG_MAX.
1261        const chunk_size = 48000;
1262        var offset: usize = 0;
1263        var is_first = true;
1264
1265        while (offset < file_content.len or is_first) {
1266            const end = @min(offset + chunk_size, file_content.len);
1267            const chunk = file_content[offset..end];
1268
1269            const encoded_len = std.base64.standard.Encoder.calcSize(chunk.len);
1270            const encoded = try gpa.alloc(u8, encoded_len);
1271            defer gpa.free(encoded);
1272            _ = std.base64.standard.Encoder.encode(encoded, chunk);
1273
1274            self.queuePtyInput(gpa, "printf '%s' '");
1275            self.queuePtyInput(gpa, encoded);
1276            if (is_first) {
1277                self.queuePtyInput(gpa, "' | base64 -d > '");
1278            } else {
1279                self.queuePtyInput(gpa, "' | base64 -d >> '");
1280            }
1281            self.queuePtyInput(gpa, file_path);
1282            self.queuePtyInput(gpa, "'");
1283            self.queuePtyInput(gpa, "\r");
1284
1285            offset = end;
1286            is_first = false;
1287        }
1288
1289        try ipc.appendMessage(gpa, &client.write_buf, .Ack, "");
1290        client.has_pending_output = true;
1291        self.has_had_client = true;
1292        std.log.debug(
1293            "write command len={d} file_path={s}",
1294            .{ file_content.len, file_path },
1295        );
1296    }
1297
1298    fn handleEnvGet(self: *Daemon, gpa: std.mem.Allocator, client: *Client) !void {
1299        const leader_opt = self.getLeaderClient();
1300        const payload = if (leader_opt) |leader| leader.env_str orelse "" else "";
1301        try ipc.appendMessage(gpa, &client.write_buf, .EnvData, payload);
1302        client.has_pending_output = true;
1303    }
1304
1305    fn handleEnvSet(_: *Daemon, gpa: std.mem.Allocator, client: *Client, env_str: []const u8) !void {
1306        std.log.info("handle env set payload={s}", .{env_str});
1307        try client.setEnv(gpa, env_str);
1308
1309        try ipc.appendMessage(gpa, &client.write_buf, .Ack, "");
1310        client.has_pending_output = true;
1311    }
1312
1313    fn handleLabelGet(self: *Daemon, gpa: std.mem.Allocator, client: *Client) !void {
1314        const out = try label.labelsToU8(gpa, self.labels);
1315        defer gpa.free(out);
1316        try ipc.appendMessage(gpa, &client.write_buf, .LabelData, out);
1317        client.has_pending_output = true;
1318    }
1319
1320    fn handleLabelSet(self: *Daemon, gpa: std.mem.Allocator, client: *Client, labels: []const u8) !void {
1321        std.log.info("handle label set payload={s}", .{labels});
1322
1323        var kvs = label.LabelIterator.init(labels);
1324        while (kvs.next()) |kv| {
1325            if (kv.value.len == 0) {
1326                if (self.labels.fetchRemove(kv.key)) |existing| {
1327                    gpa.free(existing.key);
1328                    gpa.free(existing.value);
1329                }
1330                continue;
1331            }
1332
1333            const owned_key = try gpa.dupe(u8, kv.key);
1334            errdefer gpa.free(owned_key);
1335            const owned_value = try gpa.dupe(u8, kv.value);
1336            errdefer gpa.free(owned_value);
1337            if (try self.labels.fetchPut(gpa, owned_key, owned_value)) |existing| {
1338                // fetchPut does NOT replace the key in the map, the old
1339                // key pointer stays. So free the new (unused) key and the
1340                // old value.
1341                gpa.free(owned_key);
1342                gpa.free(existing.value);
1343            }
1344        }
1345
1346        try ipc.appendMessage(gpa, &client.write_buf, .Ack, "");
1347        client.has_pending_output = true;
1348    }
1349
1350    fn handleLabelClear(self: *Daemon, gpa: std.mem.Allocator, client: *Client) !void {
1351        var it = self.labels.iterator();
1352        while (it.next()) |entry| {
1353            gpa.free(entry.key_ptr.*);
1354            gpa.free(entry.value_ptr.*);
1355        }
1356        self.labels.clearRetainingCapacity();
1357        try ipc.appendMessage(gpa, &client.write_buf, .Ack, "");
1358        client.has_pending_output = true;
1359    }
1360};
1361
1362test "terminal retains the configured scrollback without the default byte cap" {
1363    const alloc = std.testing.allocator;
1364    const cfg = Cfg{ .socket_dir = "", .log_dir = "" };
1365    var term = try initTerminal(alloc, std.testing.io, .{ .cols = 80, .rows = 24 }, &cfg);
1366    defer term.deinit(alloc);
1367    var stream = term.vtStream();
1368    defer stream.deinit();
1369
1370    stream.nextSlice("first line\r\n");
1371    for (1..cfg.max_scrollback_lines) |_| stream.nextSlice("more output\r\n");
1372
1373    const history = util.serializeTerminal(alloc, &term, .plain) orelse return error.TestUnexpectedNull;
1374    defer alloc.free(history);
1375    try std.testing.expect(std.mem.startsWith(u8, history, "first line\n"));
1376
1377    // Exceed the limit comfortably because Ghostty prunes whole pages.
1378    for (0..cfg.max_scrollback_lines) |_| stream.nextSlice("more output\r\n");
1379    stream.nextSlice("latest line\r\n");
1380
1381    const pruned_history = util.serializeTerminal(alloc, &term, .plain) orelse return error.TestUnexpectedNull;
1382    defer alloc.free(pruned_history);
1383    try std.testing.expect(std.mem.indexOf(u8, pruned_history, "first line") == null);
1384    try std.testing.expect(std.mem.indexOf(u8, pruned_history, "latest line") != null);
1385}
1386
1387fn testDaemon() Daemon {
1388    return .{
1389        .cfg = undefined,
1390        .clients = .empty,
1391        .leader_client_fd = null,
1392        .session_name = "test",
1393        .socket_path = "",
1394        .running = true,
1395        .pid = 0,
1396        .created_at = 0,
1397    };
1398}
1399
1400fn testClient(alloc: std.mem.Allocator, fd: i32, is_terminal: bool) !*Client {
1401    const c = try alloc.create(Client);
1402    c.* = .{ .alloc = alloc, .socket_fd = fd, .read_buf = try ipc.SocketBuffer.init(alloc), .write_buf = .empty };
1403    c.is_terminal = is_terminal;
1404    return c;
1405}
1406
1407test "hasTerminalClient follows attach, detach and detach-all" {
1408    const alloc = std.testing.allocator;
1409    var daemon = testDaemon();
1410    defer daemon.clients.deinit(alloc);
1411    defer daemon.pty_write_buf.deinit(alloc);
1412
1413    // Real fds (pipes) so Client.deinit's close() is legal.
1414    const a = try lib_posix.pipe2(.{});
1415    const b = try lib_posix.pipe2(.{});
1416
1417    try std.testing.expect(!daemon.hasTerminalClient());
1418
1419    // A run/send/tail client doesn't count.
1420    const tail_client = try testClient(alloc, a[0], false);
1421    try daemon.clients.append(alloc, tail_client);
1422    try std.testing.expect(!daemon.hasTerminalClient());
1423
1424    // An attach client does, until it disconnects.
1425    const term_client = try testClient(alloc, a[1], true);
1426    try daemon.clients.append(alloc, term_client);
1427    try std.testing.expect(daemon.hasTerminalClient());
1428    _ = daemon.closeClient(alloc, tail_client, 0, false);
1429    try std.testing.expect(daemon.hasTerminalClient());
1430    _ = daemon.closeClient(alloc, term_client, 0, false);
1431    try std.testing.expect(!daemon.hasTerminalClient());
1432
1433    try daemon.clients.append(alloc, try testClient(alloc, b[0], true));
1434    try daemon.clients.append(alloc, try testClient(alloc, b[1], true));
1435    daemon.handleDetachAll(alloc);
1436    try std.testing.expect(!daemon.hasTerminalClient());
1437    try std.testing.expectEqual(@as(?i32, null), daemon.leader_client_fd);
1438}
1439
1440test "send queues PTY input without changing leader" {
1441    const alloc = std.testing.allocator;
1442    var daemon = testDaemon();
1443    daemon.leader_client_fd = 42;
1444    defer daemon.pty_write_buf.deinit(alloc);
1445
1446    daemon.handleSend(alloc, "hello");
1447
1448    try std.testing.expectEqual(@as(?i32, 42), daemon.leader_client_fd);
1449    try std.testing.expectEqualStrings("hello", daemon.pty_write_buf.items);
1450}
1451
1452test "handleEnvGet returns leader client's environment variables including unsets" {
1453    const alloc = std.testing.allocator;
1454    var cfg = Cfg{
1455        .socket_dir = "/tmp",
1456        .log_dir = "/tmp",
1457    };
1458
1459    var daemon = Daemon{
1460        .cfg = &cfg,
1461        .clients = .empty,
1462        .session_name = "test",
1463        .socket_path = try alloc.dupe(u8, ""),
1464        .running = true,
1465        .pid = 0,
1466        .created_at = 0,
1467    };
1468    defer daemon.deinit(alloc);
1469
1470    const fds1 = try lib_posix.pipe2(.{});
1471    defer lib_posix.close(fds1[1]);
1472    var client1 = Client{
1473        .alloc = alloc,
1474        .socket_fd = fds1[0],
1475        .read_buf = try ipc.SocketBuffer.init(alloc),
1476        .write_buf = std.ArrayList(u8).empty,
1477    };
1478    defer client1.deinit(alloc);
1479    try client1.setEnv(alloc, "DISPLAY=:1\nSSH_AUTH_SOCK=/tmp/ssh-1\nKITTY_LISTEN_ON=unix:/tmp/kitty-1\n-WINDOWID\n");
1480
1481    const fds2 = try lib_posix.pipe2(.{});
1482    defer lib_posix.close(fds2[1]);
1483    var client2 = Client{
1484        .alloc = alloc,
1485        .socket_fd = fds2[0],
1486        .read_buf = try ipc.SocketBuffer.init(alloc),
1487        .write_buf = std.ArrayList(u8).empty,
1488    };
1489    defer client2.deinit(alloc);
1490    try client2.setEnv(alloc, "SSH_AUTH_SOCK=/tmp/ssh-2\nWINDOWID=12345\nKITTY_LISTEN_ON=unix:/tmp/kitty-2\n-DISPLAY\n");
1491
1492    const fds_req = try lib_posix.pipe2(.{});
1493    defer lib_posix.close(fds_req[1]);
1494    var client_req = Client{
1495        .alloc = alloc,
1496        .socket_fd = fds_req[0],
1497        .read_buf = try ipc.SocketBuffer.init(alloc),
1498        .write_buf = std.ArrayList(u8).empty,
1499    };
1500    defer client_req.deinit(alloc);
1501
1502    try daemon.clients.append(alloc, &client1);
1503    try daemon.clients.append(alloc, &client2);
1504
1505    // No leader yet -> handleEnvGet returns empty string payload
1506    try daemon.handleEnvGet(alloc, &client_req);
1507    try std.testing.expect(client_req.write_buf.items.len > 0);
1508    client_req.write_buf.clearRetainingCapacity();
1509
1510    // Set client1 as leader
1511    try daemon.setLeader(alloc, &client1);
1512    try std.testing.expectEqual(@as(?i32, fds1[0]), daemon.leader_client_fd);
1513    try daemon.handleEnvGet(alloc, &client_req);
1514    // Wire message: [Header][Payload]
1515    const pay1 = client_req.write_buf.items[@sizeOf(ipc.Header)..];
1516    try std.testing.expectEqualStrings("DISPLAY=:1\nSSH_AUTH_SOCK=/tmp/ssh-1\nKITTY_LISTEN_ON=unix:/tmp/kitty-1\n-WINDOWID\n", pay1);
1517    client_req.write_buf.clearRetainingCapacity();
1518
1519    // Switch leader to client2
1520    try daemon.setLeader(alloc, &client2);
1521    try std.testing.expectEqual(@as(?i32, fds2[0]), daemon.leader_client_fd);
1522    try daemon.handleEnvGet(alloc, &client_req);
1523    const pay2 = client_req.write_buf.items[@sizeOf(ipc.Header)..];
1524    try std.testing.expectEqualStrings("SSH_AUTH_SOCK=/tmp/ssh-2\nWINDOWID=12345\nKITTY_LISTEN_ON=unix:/tmp/kitty-2\n-DISPLAY\n", pay2);
1525    client_req.write_buf.clearRetainingCapacity();
1526
1527    // Client2 updates env
1528    try daemon.handleEnvSet(alloc, &client2, "DISPLAY=:99\n");
1529    try daemon.handleEnvGet(alloc, &client_req);
1530    const pay3 = client_req.write_buf.items[@sizeOf(ipc.Header)..];
1531    try std.testing.expectEqualStrings("DISPLAY=:99\n", pay3);
1532}