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}