Commit d48b642
Eric Bower
·
2026-04-03 11:26:35 -0400 EDT
parent 77d03a5
feat: new client leader policy The last client to send user input bytes (non-ansi escape codes) becomes the leader. The client leader controls resizing and any other terminal state changes. Non-leader clients are read-only until they send user input bytes and takeover leadership. Closes: https://github.com/neurosnap/zmx/issues/73
2 files changed,
+81,
-7
+6,
-0
| ... | ... | @@ -4,6 +4,12 @@ Use spec: https://common-changelog.org/ | |
| 4 | 4 | ||
| 5 | 5 | ## Staged | |
| 6 | 6 | ||
| 7 | + | ### Added | |
| 8 | + | ||
| 9 | + | - New client leader policy: last client to send user input bytes (non-ansi escape codes) becomes the leader | |
| 10 | + | - The client leader controls resizing and any other terminal state changes | |
| 11 | + | - Non-leader clients are read-only until they send user input bytes and takeover leadership | |
| 12 | + | ||
| 7 | 13 | ### Changed | |
| 8 | 14 | ||
| 9 | 15 | - `zmx kill` now supports multiple args and it will kill sessions that match a prefix |
+75,
-7
| ... | ... | @@ -122,6 +122,7 @@ pub fn main() !void { | |
| 122 | 122 | .command = command, | |
| 123 | 123 | .cwd = cwd, | |
| 124 | 124 | .created_at = @intCast(std.time.timestamp()), | |
| 125 | + | .leader_client_fd = null, | |
| 125 | 126 | }; | |
| 126 | 127 | daemon.socket_path = socket.getSocketPath(alloc, cfg.socket_dir, sesh) catch |err| switch (err) { | |
| 127 | 128 | error.NameTooLong => return socket.printSessionNameTooLong(sesh, cfg.socket_dir), |
| ... | ... | @@ -157,6 +158,7 @@ pub fn main() !void { | |
| 157 | 158 | .created_at = @intCast(std.time.timestamp()), | |
| 158 | 159 | .is_task_mode = true, | |
| 159 | 160 | .task_command = cmd_args_raw.items, | |
| 161 | + | .leader_client_fd = undefined, | |
| 160 | 162 | }; | |
| 161 | 163 | daemon.socket_path = socket.getSocketPath(alloc, cfg.socket_dir, sesh) catch |err| switch (err) { | |
| 162 | 164 | error.NameTooLong => return socket.printSessionNameTooLong(sesh, cfg.socket_dir), |
| ... | ... | @@ -334,6 +336,7 @@ const Daemon = struct { | |
| 334 | 336 | cfg: *Cfg, | |
| 335 | 337 | alloc: std.mem.Allocator, | |
| 336 | 338 | clients: std.ArrayList(*Client), | |
| 339 | + | leader_client_fd: ?i32, | |
| 337 | 340 | session_name: []const u8, | |
| 338 | 341 | socket_path: []const u8, | |
| 339 | 342 | running: bool, |
| ... | ... | @@ -356,7 +359,7 @@ const Daemon = struct { | |
| 356 | 359 | } | |
| 357 | 360 | ||
| 358 | 361 | pub fn shutdown(self: *Daemon) void { | |
| 359 | - | std.log.info("shutting down daemon session_name={s}", .{self.session_name}); | |
| 362 | + | std.log.info("shutting down daemon session={s}", .{self.session_name}); | |
| 360 | 363 | self.running = false; | |
| 361 | 364 | ||
| 362 | 365 | for (self.clients.items) |client| { |
| ... | ... | @@ -538,7 +541,7 @@ const Daemon = struct { | |
| 538 | 541 | posix.close(pty_fd); | |
| 539 | 542 | _ = posix.waitpid(self.pid, 0); | |
| 540 | 543 | posix.close(server_sock_fd); | |
| 541 | - | std.log.info("deleting socket file session_name={s}", .{self.session_name}); | |
| 544 | + | std.log.info("deleting socket file session={s}", .{self.session_name}); | |
| 542 | 545 | dir.deleteFile(self.session_name) catch |err| { | |
| 543 | 546 | std.log.warn("failed to delete socket file err={s}", .{@errorName(err)}); | |
| 544 | 547 | }; |
| ... | ... | @@ -571,6 +574,7 @@ const Daemon = struct { | |
| 571 | 574 | ); | |
| 572 | 575 | return; | |
| 573 | 576 | } | |
| 577 | + | std.log.debug("buffering pty input data={x}", .{data}); | |
| 574 | 578 | self.pty_write_buf.appendSlice(self.alloc, data) catch |err| { | |
| 575 | 579 | std.log.warn( | |
| 576 | 580 | "pty input dropped {d} bytes: {s}", |
| ... | ... | @@ -579,8 +583,50 @@ const Daemon = struct { | |
| 579 | 583 | }; | |
| 580 | 584 | } | |
| 581 | 585 | ||
| 582 | - | pub fn handleInput(self: *Daemon, payload: []const u8) void { | |
| 583 | - | self.queuePtyInput(payload); | |
| 586 | + | pub fn handleInput(self: *Daemon, client: *Client, payload: []const u8) !void { | |
| 587 | + | // client is leader, send entire payload (ansi escape codes + text) | |
| 588 | + | if (self.leader_client_fd == client.socket_fd) { | |
| 589 | + | self.queuePtyInput(payload); | |
| 590 | + | return; | |
| 591 | + | } | |
| 592 | + | ||
| 593 | + | // quick check to see if a newline happened so we can set that client to leader | |
| 594 | + | // without creating a ghostty vt | |
| 595 | + | if (std.mem.indexOfScalar(u8, payload, '\r')) |_| { | |
| 596 | + | std.log.info( | |
| 597 | + | "setting new leader session={s} client_fd={d}", | |
| 598 | + | .{ self.session_name, client.socket_fd }, | |
| 599 | + | ); | |
| 600 | + | self.leader_client_fd = client.socket_fd; | |
| 601 | + | self.queuePtyInput(payload); | |
| 602 | + | return; | |
| 603 | + | } | |
| 604 | + | ||
| 605 | + | // check if leader needs to be updated | |
| 606 | + | // this is probably really ineffecient but it was the easiest and most robust way | |
| 607 | + | // to strip ansi escape codes and only detect plain text to determine if we need | |
| 608 | + | // to set a new leader | |
| 609 | + | var termx = try ghostty_vt.Terminal.init(client.alloc, .{ | |
| 610 | + | .cols = 80, | |
| 611 | + | .rows = 24, | |
| 612 | + | }); | |
| 613 | + | defer termx.deinit(client.alloc); | |
| 614 | + | var vt_stream = termx.vtStream(); | |
| 615 | + | defer vt_stream.deinit(); | |
| 616 | + | try vt_stream.nextSlice(payload); | |
| 617 | + | if (util.serializeTerminal(client.alloc, &termx, .plain)) |output| { | |
| 618 | + | defer client.alloc.free(output); | |
| 619 | + | // if there's no text output then this client is effectively read-only until they type | |
| 620 | + | if (output.len > 0) { | |
| 621 | + | std.log.info( | |
| 622 | + | "setting new leader session={s} client_fd={d}", | |
| 623 | + | .{ self.session_name, client.socket_fd }, | |
| 624 | + | ); | |
| 625 | + | self.leader_client_fd = client.socket_fd; | |
| 626 | + | // new leader is set to this client so send *entire* payload | |
| 627 | + | self.queuePtyInput(payload); | |
| 628 | + | } | |
| 629 | + | } | |
| 584 | 630 | } | |
| 585 | 631 | ||
| 586 | 632 | pub fn handleInit( |
| ... | ... | @@ -591,6 +637,10 @@ const Daemon = struct { | |
| 591 | 637 | payload: []const u8, | |
| 592 | 638 | ) !void { | |
| 593 | 639 | if (payload.len != @sizeOf(ipc.Resize)) return; | |
| 640 | + | // no leader is set so set one | |
| 641 | + | if (self.leader_client_fd == null) { | |
| 642 | + | self.leader_client_fd = client.socket_fd; | |
| 643 | + | } | |
| 594 | 644 | ||
| 595 | 645 | const resize = std.mem.bytesToValue(ipc.Resize, payload); | |
| 596 | 646 |
| ... | ... | @@ -635,11 +685,21 @@ const Daemon = struct { | |
| 635 | 685 | ||
| 636 | 686 | pub fn handleResize( | |
| 637 | 687 | self: *Daemon, | |
| 688 | + | client: *Client, | |
| 638 | 689 | pty_fd: i32, | |
| 639 | 690 | term: *ghostty_vt.Terminal, | |
| 640 | 691 | payload: []const u8, | |
| 641 | 692 | ) !void { | |
| 642 | 693 | if (payload.len != @sizeOf(ipc.Resize)) return; | |
| 694 | + | if (self.leader_client_fd == null) { | |
| 695 | + | std.log.info( | |
| 696 | + | "setting new leader session={s} client_fd={d}", | |
| 697 | + | .{ self.session_name, client.socket_fd }, | |
| 698 | + | ); | |
| 699 | + | self.leader_client_fd = client.socket_fd; | |
| 700 | + | } | |
| 701 | + | // only leader can resize | |
| 702 | + | if (self.leader_client_fd != client.socket_fd) return; | |
| 643 | 703 | ||
| 644 | 704 | const resize = std.mem.bytesToValue(ipc.Resize, payload); | |
| 645 | 705 | var ws: cross.c.struct_winsize = .{ |
| ... | ... | @@ -654,7 +714,15 @@ const Daemon = struct { | |
| 654 | 714 | } | |
| 655 | 715 | ||
| 656 | 716 | pub fn handleDetach(self: *Daemon, client: *Client, i: usize) void { | |
| 657 | - | std.log.info("client detach fd={d}", .{client.socket_fd}); | |
| 717 | + | std.log.info("client detach session={s} fd={d}", .{ self.session_name, client.socket_fd }); | |
| 718 | + | // leader is trying to disconnect, remove ref and let another client claim leader on input | |
| 719 | + | if (self.leader_client_fd == client.socket_fd) { | |
| 720 | + | std.log.info( | |
| 721 | + | "unsetting leader session={s} fd={d}", | |
| 722 | + | .{ self.session_name, client.socket_fd }, | |
| 723 | + | ); | |
| 724 | + | self.leader_client_fd = null; | |
| 725 | + | } | |
| 658 | 726 | _ = self.closeClient(client, i, false); | |
| 659 | 727 | } | |
| 660 | 728 |
| ... | ... | @@ -1680,9 +1748,9 @@ fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 1680 | 1748 | ||
| 1681 | 1749 | while (client.read_buf.next()) |msg| { | |
| 1682 | 1750 | switch (msg.header.tag) { | |
| 1683 | - | .Input => daemon.handleInput(msg.payload), | |
| 1751 | + | .Input => try daemon.handleInput(client, msg.payload), | |
| 1684 | 1752 | .Init => try daemon.handleInit(client, pty_fd, &term, msg.payload), | |
| 1685 | - | .Resize => try daemon.handleResize(pty_fd, &term, msg.payload), | |
| 1753 | + | .Resize => try daemon.handleResize(client, pty_fd, &term, msg.payload), | |
| 1686 | 1754 | .Detach => { | |
| 1687 | 1755 | daemon.handleDetach(client, i); | |
| 1688 | 1756 | break :clients_loop; |