Commit 0d66096
Eric Bower
·
2026-07-27 09:28:10 -0400 EDT
parent 9e5eb55
refactor(daemon): restructure io and allocators Because we perform a double-fork when spawning a daemon we have to be careful about multithreading and mutex locks within io and allocators or else we could end up with a deadlock. I don't fully understand how this happens within the DebugAllocator but I've hit deadlocks previously. So I restructured the Daemon to not store io or an allocator and made it much more clear what is going on and why we have to create a new allocator after the double-fork. I also extracted the daemonize logic into its own file and set of fns mainly because that code doesn't need to change much and it is mostly fork() machinery.
11 files changed,
+1818,
-1743
+10,
-10
| ... | ... | @@ -13,7 +13,7 @@ const macos_targets: []const std.Target.Query = &.{ | |
| 13 | 13 | ||
| 14 | 14 | pub fn build(b: *std.Build) void { | |
| 15 | 15 | const target = b.standardTargetOptions(.{}); | |
| 16 | - | const is_macos = target.result.os.tag == .macos; | |
| 16 | + | // const is_macos = target.result.os.tag == .macos; | |
| 17 | 17 | const optimize = b.standardOptimizeOption(.{}); | |
| 18 | 18 | const version = b.option([]const u8, "version", "Version string for release") orelse | |
| 19 | 19 | @as([]const u8, build_zig_zon.version); |
| ... | ... | @@ -51,8 +51,8 @@ pub fn build(b: *std.Build) void { | |
| 51 | 51 | const run_step = b.step("run", "Run the app"); | |
| 52 | 52 | const exe = b.addExecutable(.{ | |
| 53 | 53 | .name = "zmx", | |
| 54 | - | .use_llvm = true, | |
| 55 | - | .use_lld = !is_macos, | |
| 54 | + | // .use_llvm = true, | |
| 55 | + | // .use_lld = !is_macos, | |
| 56 | 56 | .root_module = exe_mod, | |
| 57 | 57 | }); | |
| 58 | 58 |
| ... | ... | @@ -85,8 +85,8 @@ pub fn build(b: *std.Build) void { | |
| 85 | 85 | ); | |
| 86 | 86 | const exe_unit_tests = b.addTest(.{ | |
| 87 | 87 | .root_module = test_module, | |
| 88 | - | .use_llvm = true, | |
| 89 | - | .use_lld = !is_macos, | |
| 88 | + | // .use_llvm = true, | |
| 89 | + | // .use_lld = !is_macos, | |
| 90 | 90 | }); | |
| 91 | 91 | ||
| 92 | 92 | const run_exe_unit_tests = b.addRunArtifact(exe_unit_tests); |
| ... | ... | @@ -98,8 +98,8 @@ pub fn build(b: *std.Build) void { | |
| 98 | 98 | const check = b.step("check", "Check if zmx compiles"); | |
| 99 | 99 | const exe_check = b.addExecutable(.{ | |
| 100 | 100 | .name = "zmx", | |
| 101 | - | .use_llvm = true, | |
| 102 | - | .use_lld = !is_macos, | |
| 101 | + | // .use_llvm = true, | |
| 102 | + | // .use_lld = !is_macos, | |
| 103 | 103 | .root_module = exe_mod, | |
| 104 | 104 | }); | |
| 105 | 105 |
| ... | ... | @@ -136,11 +136,11 @@ pub fn build(b: *std.Build) void { | |
| 136 | 136 | release_mod.addImport("ghostty-vt", release_dep.module("ghostty-vt")); | |
| 137 | 137 | } | |
| 138 | 138 | ||
| 139 | - | const is_local_macos = resolved.result.os.tag == .macos; | |
| 139 | + | // const is_local_macos = resolved.result.os.tag == .macos; | |
| 140 | 140 | const release_exe = b.addExecutable(.{ | |
| 141 | 141 | .name = "zmx", | |
| 142 | - | .use_llvm = true, | |
| 143 | - | .use_lld = !is_local_macos, | |
| 142 | + | // .use_llvm = true, | |
| 143 | + | // .use_lld = !is_local_macos, | |
| 144 | 144 | .root_module = release_mod, | |
| 145 | 145 | }); | |
| 146 | 146 |
+133,
-0
| ... | ... | @@ -0,0 +1,133 @@ | |
| 1 | + | /// Cfg is zmx's configuration container. | |
| 2 | + | /// | |
| 3 | + | /// The purpose of this container is to hold anything that can be modified by the user. | |
| 4 | + | pub const Cfg = @This(); | |
| 5 | + | ||
| 6 | + | const std = @import("std"); | |
| 7 | + | const lib_posix = @import("posix.zig"); | |
| 8 | + | const cross = @import("cross.zig"); | |
| 9 | + | ||
| 10 | + | socket_dir: []const u8, | |
| 11 | + | log_dir: []const u8, | |
| 12 | + | max_scrollback: usize = 10_000_000, | |
| 13 | + | dir_mode: u32 = 0o750, | |
| 14 | + | log_mode: u32 = 0o640, | |
| 15 | + | ||
| 16 | + | pub fn init(alloc: std.mem.Allocator, io: std.Io) !Cfg { | |
| 17 | + | const socket_dir = try socketDir(alloc); | |
| 18 | + | errdefer alloc.free(socket_dir); | |
| 19 | + | const log_dir = try logDir(alloc); | |
| 20 | + | errdefer alloc.free(log_dir); | |
| 21 | + | ||
| 22 | + | const dir_mode = if (lib_posix.getenv("ZMX_DIR_MODE")) |m| | |
| 23 | + | std.fmt.parseInt(u32, m, 8) catch 0o750 | |
| 24 | + | else | |
| 25 | + | 0o750; | |
| 26 | + | ||
| 27 | + | const log_mode = if (lib_posix.getenv("ZMX_LOG_MODE")) |m| | |
| 28 | + | std.fmt.parseInt(u32, m, 8) catch 0o640 | |
| 29 | + | else | |
| 30 | + | 0o640; | |
| 31 | + | ||
| 32 | + | var cfg = Cfg{ | |
| 33 | + | .socket_dir = socket_dir, | |
| 34 | + | .log_dir = log_dir, | |
| 35 | + | .dir_mode = dir_mode, | |
| 36 | + | .log_mode = log_mode, | |
| 37 | + | }; | |
| 38 | + | ||
| 39 | + | try cfg.mkdir(io); | |
| 40 | + | ||
| 41 | + | return cfg; | |
| 42 | + | } | |
| 43 | + | ||
| 44 | + | fn socketDir(alloc: std.mem.Allocator) ![]const u8 { | |
| 45 | + | const tmpdir = std.mem.trimEnd(u8, lib_posix.getenv("TMPDIR") orelse "/tmp", "/"); | |
| 46 | + | const uid = lib_posix.getuid(); | |
| 47 | + | ||
| 48 | + | const socket_dir: []const u8 = if (lib_posix.getenv("ZMX_DIR")) |zmxdir| | |
| 49 | + | try alloc.dupe(u8, zmxdir) | |
| 50 | + | else if (lib_posix.getenv("XDG_RUNTIME_DIR")) |xdg_runtime| | |
| 51 | + | try std.fmt.allocPrint(alloc, "{s}/zmx", .{xdg_runtime}) | |
| 52 | + | else | |
| 53 | + | try std.fmt.allocPrint(alloc, "{s}/zmx-{d}", .{ tmpdir, uid }); | |
| 54 | + | ||
| 55 | + | return socket_dir; | |
| 56 | + | } | |
| 57 | + | ||
| 58 | + | fn logDir(alloc: std.mem.Allocator) ![]const u8 { | |
| 59 | + | const log_dir = if (lib_posix.getenv("ZMX_DIR")) |zmxdir| | |
| 60 | + | try std.fmt.allocPrint(alloc, "{s}/logs", .{zmxdir}) | |
| 61 | + | else if (lib_posix.getenv("XDG_STATE_HOME")) |xdg_state_home| | |
| 62 | + | try std.fmt.allocPrint(alloc, "{s}/zmx/logs", .{xdg_state_home}) | |
| 63 | + | else if (lib_posix.getenv("HOME")) |home_dir| | |
| 64 | + | try std.fmt.allocPrint(alloc, "{s}/.local/state/zmx/logs", .{home_dir}) | |
| 65 | + | else fallback: { | |
| 66 | + | // This is the last resort: falling back to /tmp/$UID if HOME is unset. | |
| 67 | + | const tmpdir = std.mem.trimEnd(u8, lib_posix.getenv("TMPDIR") orelse "/tmp", "/"); | |
| 68 | + | const uid = lib_posix.getuid(); | |
| 69 | + | break :fallback try std.fmt.allocPrint(alloc, "{s}/zmx-{d}", .{ tmpdir, uid }); | |
| 70 | + | }; | |
| 71 | + | ||
| 72 | + | return log_dir; | |
| 73 | + | } | |
| 74 | + | ||
| 75 | + | pub fn deinit(self: *Cfg, alloc: std.mem.Allocator) void { | |
| 76 | + | if (self.socket_dir.len > 0) alloc.free(self.socket_dir); | |
| 77 | + | if (self.log_dir.len > 0) alloc.free(self.log_dir); | |
| 78 | + | } | |
| 79 | + | ||
| 80 | + | pub fn mkdir(self: *Cfg, io: std.Io) !void { | |
| 81 | + | const sock_perms = std.Io.Dir.Permissions.fromMode(@intCast(self.dir_mode)); | |
| 82 | + | try mkdirAll(io, self.socket_dir, sock_perms); | |
| 83 | + | const log_perms = std.Io.Dir.Permissions.fromMode(@intCast(self.dir_mode)); | |
| 84 | + | try mkdirAll(io, self.log_dir, log_perms); | |
| 85 | + | } | |
| 86 | + | ||
| 87 | + | fn mkdirAll(io: std.Io, sub_dir_path: []const u8, permissions: std.Io.Dir.Permissions) !void { | |
| 88 | + | var it = std.fs.path.componentIterator(sub_dir_path); | |
| 89 | + | var component = it.last() orelse return error.BadPathName; | |
| 90 | + | while (true) { | |
| 91 | + | std.Io.Dir.createDirAbsolute(io, component.path, permissions) catch |err| switch (err) { | |
| 92 | + | error.PathAlreadyExists => {}, | |
| 93 | + | error.FileNotFound => |e| { | |
| 94 | + | component = it.previous() orelse return e; | |
| 95 | + | continue; | |
| 96 | + | }, | |
| 97 | + | else => |e| return e, | |
| 98 | + | }; | |
| 99 | + | component = it.next() orelse return; | |
| 100 | + | } | |
| 101 | + | } | |
| 102 | + | ||
| 103 | + | test "Cfg.init uses default modes when env vars are not set" { | |
| 104 | + | const alloc = std.testing.allocator; | |
| 105 | + | ||
| 106 | + | // Ensure they are not set | |
| 107 | + | _ = cross.c.unsetenv("ZMX_DIR_MODE"); | |
| 108 | + | _ = cross.c.unsetenv("ZMX_LOG_MODE"); | |
| 109 | + | ||
| 110 | + | var cfg = try Cfg.init(alloc, std.testing.io); | |
| 111 | + | defer cfg.deinit(alloc); | |
| 112 | + | ||
| 113 | + | try std.testing.expectEqual(@as(u32, 0o750), cfg.dir_mode); | |
| 114 | + | try std.testing.expectEqual(@as(u32, 0o640), cfg.log_mode); | |
| 115 | + | } | |
| 116 | + | ||
| 117 | + | test "Cfg.init uses custom modes from env vars" { | |
| 118 | + | const alloc = std.testing.allocator; | |
| 119 | + | ||
| 120 | + | // Set custom octal values | |
| 121 | + | _ = cross.c.setenv("ZMX_DIR_MODE", "770", 1); | |
| 122 | + | _ = cross.c.setenv("ZMX_LOG_MODE", "660", 1); | |
| 123 | + | defer { | |
| 124 | + | _ = cross.c.unsetenv("ZMX_DIR_MODE"); | |
| 125 | + | _ = cross.c.unsetenv("ZMX_LOG_MODE"); | |
| 126 | + | } | |
| 127 | + | ||
| 128 | + | var cfg = try Cfg.init(alloc, std.testing.io); | |
| 129 | + | defer cfg.deinit(alloc); | |
| 130 | + | ||
| 131 | + | try std.testing.expectEqual(@as(u32, 0o770), cfg.dir_mode); | |
| 132 | + | try std.testing.expectEqual(@as(u32, 0o660), cfg.log_mode); | |
| 133 | + | } |
+211,
-0
| ... | ... | @@ -0,0 +1,211 @@ | |
| 1 | + | const std = @import("std"); | |
| 2 | + | const lib_posix = @import("posix.zig"); | |
| 3 | + | const Cfg = @import("cfg.zig"); | |
| 4 | + | const socket = @import("socket.zig"); | |
| 5 | + | const ipc = @import("ipc.zig"); | |
| 6 | + | const assert = std.debug.assert; | |
| 7 | + | const log = @import("log.zig"); | |
| 8 | + | const cross = @import("cross.zig"); | |
| 9 | + | ||
| 10 | + | const Cmd = struct { | |
| 11 | + | file: [*:0]const u8, | |
| 12 | + | argv_ptr: [*:null]const ?[*:0]const u8, | |
| 13 | + | }; | |
| 14 | + | ||
| 15 | + | pub fn createCmdZ(def_shell: []const u8, is_task_mode: bool, command: ?[]const []const u8) !Cmd { | |
| 16 | + | const gpa = std.heap.c_allocator; | |
| 17 | + | ||
| 18 | + | if (command) |cmd_args| { | |
| 19 | + | const argv = try gpa.allocSentinel(?[*:0]const u8, cmd_args.len, null); | |
| 20 | + | for (cmd_args, 0..) |arg, i| { | |
| 21 | + | argv[i] = try gpa.dupeZ(u8, arg); | |
| 22 | + | } | |
| 23 | + | return .{ | |
| 24 | + | .file = argv[0].?, | |
| 25 | + | .argv_ptr = argv.ptr, | |
| 26 | + | }; | |
| 27 | + | } | |
| 28 | + | ||
| 29 | + | const z = try std.fmt.allocPrintSentinel(gpa, "{s}", .{def_shell}, 0); | |
| 30 | + | const shell: [:0]const u8 = if (is_task_mode) "bash" else z; | |
| 31 | + | ||
| 32 | + | // Use "-shellname" as argv[0] to signal login shell (traditional method) | |
| 33 | + | const login_shell = try std.fmt.allocPrintSentinel(gpa, "-{s}", .{std.fs.path.basename(shell)}, 0); | |
| 34 | + | const argv = try gpa.allocSentinel(?[*:0]const u8, 1, null); | |
| 35 | + | argv[0] = login_shell.ptr; | |
| 36 | + | ||
| 37 | + | return .{ | |
| 38 | + | .file = shell, | |
| 39 | + | .argv_ptr = argv, | |
| 40 | + | }; | |
| 41 | + | } | |
| 42 | + | ||
| 43 | + | /// Runs in the forked child. Either execs or returns an error (caller | |
| 44 | + | /// must exit on error -- returning would fall through to parent code). | |
| 45 | + | fn exec(sesh_name: []const u8, cmd: Cmd) !noreturn { | |
| 46 | + | const gpa = std.heap.c_allocator; | |
| 47 | + | ||
| 48 | + | // main() set SIGPIPE to SIG_IGN, which (unlike handlers) survives | |
| 49 | + | // exec. Restore the default so the shell and its children behave | |
| 50 | + | // normally (e.g. `yes | head` should exit 141 via SIGPIPE). | |
| 51 | + | const dfl: lib_posix.Sigaction = .{ | |
| 52 | + | .handler = .{ .handler = lib_posix.SIG.DFL }, | |
| 53 | + | .mask = lib_posix.sigemptyset(), | |
| 54 | + | .flags = 0, | |
| 55 | + | }; | |
| 56 | + | lib_posix.sigaction(lib_posix.SIG.PIPE, &dfl, null); | |
| 57 | + | ||
| 58 | + | const session_env = try std.fmt.allocPrintSentinel( | |
| 59 | + | gpa, | |
| 60 | + | "ZMX_SESSION={s}", | |
| 61 | + | .{sesh_name}, | |
| 62 | + | 0, | |
| 63 | + | ); | |
| 64 | + | _ = cross.c.putenv(session_env.ptr); | |
| 65 | + | ||
| 66 | + | const err = lib_posix.execvpeZ(cmd.file, cmd.argv_ptr, std.c.environ); | |
| 67 | + | std.log.err("execvpe failed: cmd={s} err={s}", .{ cmd.file, @errorName(err) }); | |
| 68 | + | lib_posix.exit(1); | |
| 69 | + | } | |
| 70 | + | ||
| 71 | + | pub const PtyInfo = struct { | |
| 72 | + | master_fd: c_int = undefined, | |
| 73 | + | pid: c_int = undefined, | |
| 74 | + | }; | |
| 75 | + | ||
| 76 | + | /// spawnPty runs forkpty() and executes the shell or shell command the user | |
| 77 | + | /// provides. | |
| 78 | + | /// | |
| 79 | + | /// This is the second fork in the double-fork technique explained in the | |
| 80 | + | /// daemonize() comment. | |
| 81 | + | pub fn spawnPty(sesh_name: []const u8, cmd: Cmd) !PtyInfo { | |
| 82 | + | const size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO); | |
| 83 | + | var ws: cross.c.struct_winsize = .{ | |
| 84 | + | .ws_row = size.rows, | |
| 85 | + | .ws_col = size.cols, | |
| 86 | + | .ws_xpixel = size.xpixel, | |
| 87 | + | .ws_ypixel = size.ypixel, | |
| 88 | + | }; | |
| 89 | + | ||
| 90 | + | var master_fd: c_int = undefined; | |
| 91 | + | const pid = cross.forkpty(&master_fd, null, null, &ws); | |
| 92 | + | if (pid < 0) { | |
| 93 | + | return error.ForkPtyFailed; | |
| 94 | + | } | |
| 95 | + | ||
| 96 | + | if (pid == 0) { // child pid code path | |
| 97 | + | // In the forked child, ANY error must exit rather than propagate: | |
| 98 | + | // a returned error falls through to the parent code path below, | |
| 99 | + | // running a second daemon on the same socket (or worse, hitting | |
| 100 | + | // errdefers that delete the parent's socket file). | |
| 101 | + | exec(sesh_name, cmd) catch |err| { | |
| 102 | + | std.log.err("child setup failed: {s}", .{@errorName(err)}); | |
| 103 | + | lib_posix.exit(1); | |
| 104 | + | }; | |
| 105 | + | unreachable; // exec() either execs or exits, never returns ok | |
| 106 | + | } | |
| 107 | + | // master pid code path | |
| 108 | + | std.log.info("pty spawned session={s} pid={d}", .{ sesh_name, pid }); | |
| 109 | + | ||
| 110 | + | // make pty non-blocking | |
| 111 | + | const flags = try lib_posix.fcntl(master_fd, lib_posix.F.GETFL, 0); | |
| 112 | + | _ = try lib_posix.fcntl(master_fd, lib_posix.F.SETFL, flags | lib_posix.O_NONBLOCK); | |
| 113 | + | ||
| 114 | + | return .{ | |
| 115 | + | .master_fd = master_fd, | |
| 116 | + | .pid = pid, | |
| 117 | + | }; | |
| 118 | + | } | |
| 119 | + | ||
| 120 | + | /// daemonize is the first fork in a double-fork technique to create a | |
| 121 | + | /// completely disconnected session (group of processes). | |
| 122 | + | /// | |
| 123 | + | /// When launching a daemon, you normally set the child process of the fork to | |
| 124 | + | /// be the session leader via setsid() which creates a new session that does | |
| 125 | + | /// *not* have a controlling terminal. This is important because we don't want | |
| 126 | + | /// a controlling terminal for our daemon or else our daemon could receive | |
| 127 | + | /// signals to shutdown when the controlling terminal closes. | |
| 128 | + | /// | |
| 129 | + | /// However, if the first fork's child process is also the daemon process, then | |
| 130 | + | /// it's technically possible for the daemon to open a terminal device | |
| 131 | + | /// (e.g. open("/dev/console", O_RDWR)) and then it would acquire a controlling | |
| 132 | + | /// terminal! A controlling terminal would expose the daemon to | |
| 133 | + | /// terminal-generated signals (e.g. SIGINT) or SIGHUP from terminal disconnect | |
| 134 | + | /// which could kill the daemon. | |
| 135 | + | /// | |
| 136 | + | /// By forking a second time, the grandchild process (the daemon) is not the | |
| 137 | + | /// session leader. Per POSIX, only a process that is the session leader can | |
| 138 | + | /// acquire a controlling terminal. | |
| 139 | + | /// | |
| 140 | + | /// Apparently this is considered "being paranoid" but appears to be a standard | |
| 141 | + | /// practice for daemons so we're doing it anyway. | |
| 142 | + | /// | |
| 143 | + | /// https://pubs.opengroup.org/onlinepubs/9699919799/basedefs/V1_chap11.html#tag_11_01_03 | |
| 144 | + | /// https://stackoverflow.com/a/16317668 | |
| 145 | + | pub fn daemonize(sesh_name: []const u8, cmd: Cmd, keep_fds_open: []i32) !PtyInfo { | |
| 146 | + | // creates the daemon | |
| 147 | + | const pid = try lib_posix.fork(); | |
| 148 | + | assert(pid != -1); | |
| 149 | + | ||
| 150 | + | if (pid > 0) { // parent (client) | |
| 151 | + | // cannot use a passed-in io or alloc after a fork so we create what we need | |
| 152 | + | // after the fork() | |
| 153 | + | var threaded: std.Io.Threaded = .init_single_threaded; | |
| 154 | + | defer threaded.deinit(); | |
| 155 | + | const io = threaded.io(); | |
| 156 | + | std.Io.sleep(io, std.Io.Duration.fromMilliseconds(10), .real) catch unreachable; | |
| 157 | + | return error.IsClientProc; | |
| 158 | + | } | |
| 159 | + | ||
| 160 | + | assert(pid == 0); // child (daemon's parent in double-fork) | |
| 161 | + | // becomes the session leader and detaches process from its controlling terminal | |
| 162 | + | _ = try lib_posix.setsid(); | |
| 163 | + | ||
| 164 | + | // Redirect stdin/stdout/stderr to /dev/null. The daemon | |
| 165 | + | // communicates via its unix socket, not stdio. Without | |
| 166 | + | // this, any pipe on FDs 0-2 (e.g. from bats' `run` | |
| 167 | + | // keyword) stays open for the daemon's lifetime, causing | |
| 168 | + | // the caller to hang waiting for EOF. | |
| 169 | + | { | |
| 170 | + | const devnull = lib_posix.open( | |
| 171 | + | "/dev/null", | |
| 172 | + | .{ .ACCMODE = .RDWR }, | |
| 173 | + | 0, | |
| 174 | + | ) catch |err| { | |
| 175 | + | std.log.warn("failed to open /dev/null: {s}", .{@errorName(err)}); | |
| 176 | + | return err; | |
| 177 | + | }; | |
| 178 | + | inline for (.{ lib_posix.STDIN_FILENO, lib_posix.STDOUT_FILENO, lib_posix.STDERR_FILENO }) |fd| { | |
| 179 | + | _ = lib_posix.dup2(devnull, fd) catch |err| { | |
| 180 | + | std.log.warn("dup2 /dev/null -> {d}: {s}", .{ fd, @errorName(err) }); | |
| 181 | + | return err; | |
| 182 | + | }; | |
| 183 | + | } | |
| 184 | + | var found = false; | |
| 185 | + | for (keep_fds_open) |fd| { | |
| 186 | + | if (devnull == fd) found = true; | |
| 187 | + | } | |
| 188 | + | if (devnull > 2 and !found) lib_posix.close(devnull); | |
| 189 | + | } | |
| 190 | + | ||
| 191 | + | // Close file descriptors inherited from the parent that the | |
| 192 | + | // daemon doesn't need. This prevents test harnesses (like | |
| 193 | + | // bats) from hanging: they wait for their internal FDs (3+) | |
| 194 | + | // to close before exiting. | |
| 195 | + | // | |
| 196 | + | // Skip any fds that the caller wants to keep open, e.g. server_sock_fd | |
| 197 | + | // (needed for IPC) and dir.fd (needed to delete the socket file on | |
| 198 | + | // shutdown). | |
| 199 | + | { | |
| 200 | + | var fd: i32 = 3; | |
| 201 | + | while (fd < 64) : (fd += 1) { | |
| 202 | + | var found = false; | |
| 203 | + | for (keep_fds_open) |kfd| { | |
| 204 | + | if (fd == kfd) found = true; | |
| 205 | + | } | |
| 206 | + | if (!found) _ = std.c.close(fd); | |
| 207 | + | } | |
| 208 | + | } | |
| 209 | + | ||
| 210 | + | return spawnPty(sesh_name, cmd); | |
| 211 | + | } |
+2,
-2
| ... | ... | @@ -94,7 +94,7 @@ pub fn send(fd: i32, tag: Tag, data: []const u8) !void { | |
| 94 | 94 | } | |
| 95 | 95 | ||
| 96 | 96 | pub fn appendMessage( | |
| 97 | - | alloc: std.mem.Allocator, | |
| 97 | + | gpa: std.mem.Allocator, | |
| 98 | 98 | list: *std.ArrayList(u8), | |
| 99 | 99 | tag: Tag, | |
| 100 | 100 | data: []const u8, |
| ... | ... | @@ -105,7 +105,7 @@ pub fn appendMessage( | |
| 105 | 105 | }; | |
| 106 | 106 | // Guarantee capacity for header + payload in one check to avoid | |
| 107 | 107 | // intermediate realloc between the two appends on the hot path. | |
| 108 | - | try list.ensureTotalCapacity(alloc, list.items.len + @sizeOf(Header) + data.len); | |
| 108 | + | try list.ensureTotalCapacity(gpa, list.items.len + @sizeOf(Header) + data.len); | |
| 109 | 109 | list.appendSliceAssumeCapacity(std.mem.asBytes(&header)); | |
| 110 | 110 | if (data.len > 0) { | |
| 111 | 111 | list.appendSliceAssumeCapacity(data); |
+18,
-18
| ... | ... | @@ -1,19 +1,28 @@ | |
| 1 | 1 | const std = @import("std"); | |
| 2 | 2 | ||
| 3 | + | pub var log_system = LogSystem{}; | |
| 4 | + | ||
| 5 | + | pub fn zmxLogFn( | |
| 6 | + | comptime level: std.log.Level, | |
| 7 | + | comptime scope: anytype, | |
| 8 | + | comptime format: []const u8, | |
| 9 | + | args: anytype, | |
| 10 | + | ) void { | |
| 11 | + | log_system.log(level, scope, format, args) catch {}; | |
| 12 | + | } | |
| 13 | + | ||
| 3 | 14 | pub const LogSystem = struct { | |
| 4 | 15 | file: ?std.Io.File = null, | |
| 5 | 16 | mutex: std.Io.Mutex = .init, | |
| 6 | 17 | current_size: u64 = 0, | |
| 7 | - | max_size: u64 = 5 * 1024 * 1024, // 5MB | |
| 18 | + | max_size: u64 = 2 * 1024 * 1024, // 2MB | |
| 8 | 19 | path: []const u8 = "", | |
| 9 | - | alloc: std.mem.Allocator = undefined, | |
| 10 | 20 | io: std.Io = undefined, | |
| 11 | 21 | mode: std.Io.File.Permissions = std.Io.File.Permissions.fromMode(0o640), | |
| 12 | 22 | ||
| 13 | - | pub fn init(self: *LogSystem, alloc: std.mem.Allocator, io: std.Io, path: []const u8, mode: std.Io.File.Permissions) !void { | |
| 14 | - | self.alloc = alloc; | |
| 23 | + | pub fn init(self: *LogSystem, io: std.Io, path: []const u8, mode: std.Io.File.Permissions) !void { | |
| 15 | 24 | self.io = io; | |
| 16 | - | self.path = try alloc.dupe(u8, path); | |
| 25 | + | self.path = path; | |
| 17 | 26 | self.mode = mode; | |
| 18 | 27 | ||
| 19 | 28 | const file = std.Io.Dir.openFileAbsolute(self.io, path, .{ .mode = .read_write }) catch |err| switch (err) { |
| ... | ... | @@ -35,7 +44,6 @@ pub const LogSystem = struct { | |
| 35 | 44 | ||
| 36 | 45 | pub fn deinit(self: *LogSystem) void { | |
| 37 | 46 | if (self.file) |f| std.Io.File.close(f, self.io); | |
| 38 | - | if (self.path.len > 0) self.alloc.free(self.path); | |
| 39 | 47 | } | |
| 40 | 48 | ||
| 41 | 49 | pub fn log( |
| ... | ... | @@ -54,12 +62,12 @@ pub const LogSystem = struct { | |
| 54 | 62 | } | |
| 55 | 63 | ||
| 56 | 64 | if (self.current_size >= self.max_size) { | |
| 57 | - | self.rotate() catch |err| { | |
| 58 | - | std.debug.print("Log rotation failed: {s}\n", .{@errorName(err)}); | |
| 65 | + | self.wipe() catch |err| { | |
| 66 | + | std.debug.print("Log wipe failed: {s}\n", .{@errorName(err)}); | |
| 59 | 67 | }; | |
| 60 | 68 | } | |
| 61 | 69 | ||
| 62 | - | const now: i64 = @intCast(@divTrunc(std.Io.Timestamp.now(self.io, .real).nanoseconds, std.time.ns_per_ms)); | |
| 70 | + | const now: std.Io.Timestamp = .now(self.io, .real); | |
| 63 | 71 | const prefix = "[{d}] [{s}] ({s}): "; | |
| 64 | 72 | const scope_name = @tagName(scope); | |
| 65 | 73 | const level_name = level.asText(); |
| ... | ... | @@ -84,20 +92,12 @@ pub const LogSystem = struct { | |
| 84 | 92 | } | |
| 85 | 93 | } | |
| 86 | 94 | ||
| 87 | - | fn rotate(self: *LogSystem) !void { | |
| 95 | + | fn wipe(self: *LogSystem) !void { | |
| 88 | 96 | if (self.file) |f| { | |
| 89 | 97 | std.Io.File.close(f, self.io); | |
| 90 | 98 | self.file = null; | |
| 91 | 99 | } | |
| 92 | 100 | ||
| 93 | - | const old_path = try std.fmt.allocPrint(self.alloc, "{s}.old", .{self.path}); | |
| 94 | - | defer self.alloc.free(old_path); | |
| 95 | - | ||
| 96 | - | std.Io.Dir.renameAbsolute(self.path, old_path, self.io) catch |err| switch (err) { | |
| 97 | - | error.FileNotFound => {}, | |
| 98 | - | else => return err, | |
| 99 | - | }; | |
| 100 | - | ||
| 101 | 101 | self.file = try std.Io.Dir.createFileAbsolute( | |
| 102 | 102 | self.io, | |
| 103 | 103 | self.path, |
+1212,
-0
| ... | ... | @@ -0,0 +1,1212 @@ | |
| 1 | + | const std = @import("std"); | |
| 2 | + | const ghostty_vt = @import("ghostty-vt"); | |
| 3 | + | const ipc = @import("ipc.zig"); | |
| 4 | + | const log = @import("log.zig"); | |
| 5 | + | const util = @import("util.zig"); | |
| 6 | + | const cross = @import("cross.zig"); | |
| 7 | + | const socket = @import("socket.zig"); | |
| 8 | + | const label = @import("label.zig"); | |
| 9 | + | const lib_posix = @import("posix.zig"); | |
| 10 | + | const Cfg = @import("cfg.zig"); | |
| 11 | + | const signal = @import("signal.zig"); | |
| 12 | + | const assert = std.debug.assert; | |
| 13 | + | const daemonize = @import("daemonize.zig"); | |
| 14 | + | const 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. | |
| 18 | + | pub fn clientLoop(client_sock_fd: i32) !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 | + | // Send init message with terminal size (buffered) | |
| 45 | + | const size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO); | |
| 46 | + | try ipc.appendMessage(gpa, &sock_write_buf, .Init, std.mem.asBytes(&size)); | |
| 47 | + | ||
| 48 | + | var poll_fds = try std.ArrayList(lib_posix.pollfd).initCapacity(gpa, 4); | |
| 49 | + | defer poll_fds.deinit(gpa); | |
| 50 | + | ||
| 51 | + | var read_buf = try ipc.SocketBuffer.init(gpa); | |
| 52 | + | defer read_buf.deinit(); | |
| 53 | + | ||
| 54 | + | var stdout_buf = try std.ArrayList(u8).initCapacity(gpa, 4096); | |
| 55 | + | defer stdout_buf.deinit(gpa); | |
| 56 | + | ||
| 57 | + | const stdin_fd = lib_posix.STDIN_FILENO; | |
| 58 | + | ||
| 59 | + | // Make stdin non-blocking. O_NONBLOCK is set on the open file description, | |
| 60 | + | // which is shared with the parent shell; restore on exit to avoid | |
| 61 | + | // corrupting the parent's stdin. | |
| 62 | + | const stdin_orig_flags = try lib_posix.fcntl(stdin_fd, lib_posix.F.GETFL, 0); | |
| 63 | + | _ = try lib_posix.fcntl(stdin_fd, lib_posix.F.SETFL, stdin_orig_flags | lib_posix.O_NONBLOCK); | |
| 64 | + | defer _ = lib_posix.fcntl(stdin_fd, lib_posix.F.SETFL, stdin_orig_flags) catch {}; | |
| 65 | + | ||
| 66 | + | while (true) { | |
| 67 | + | poll_fds.clearRetainingCapacity(); | |
| 68 | + | ||
| 69 | + | try poll_fds.append(gpa, .{ | |
| 70 | + | .fd = stdin_fd, | |
| 71 | + | .events = lib_posix.POLL.IN, | |
| 72 | + | .revents = 0, | |
| 73 | + | }); | |
| 74 | + | ||
| 75 | + | // Poll socket for read, and also for write if we have pending data | |
| 76 | + | var sock_events: i16 = lib_posix.POLL.IN; | |
| 77 | + | if (sock_write_buf.items.len > 0) { | |
| 78 | + | sock_events |= lib_posix.POLL.OUT; | |
| 79 | + | } | |
| 80 | + | try poll_fds.append(gpa, .{ | |
| 81 | + | .fd = client_sock_fd, | |
| 82 | + | .events = sock_events, | |
| 83 | + | .revents = 0, | |
| 84 | + | }); | |
| 85 | + | ||
| 86 | + | try poll_fds.append(gpa, .{ .fd = signal.sig_pipe[0], .events = lib_posix.POLL.IN, .revents = 0 }); | |
| 87 | + | ||
| 88 | + | if (stdout_buf.items.len > 0) { | |
| 89 | + | try poll_fds.append(gpa, .{ | |
| 90 | + | .fd = lib_posix.STDOUT_FILENO, | |
| 91 | + | .events = lib_posix.POLL.OUT, | |
| 92 | + | .revents = 0, | |
| 93 | + | }); | |
| 94 | + | } | |
| 95 | + | ||
| 96 | + | _ = try lib_posix.poll(poll_fds.items, -1); | |
| 97 | + | ||
| 98 | + | if (poll_fds.items[2].revents & lib_posix.POLL.IN != 0) { | |
| 99 | + | signal.drainSignalPipe(); | |
| 100 | + | const next_size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO); | |
| 101 | + | try ipc.appendMessage(gpa, &sock_write_buf, .Resize, std.mem.asBytes(&next_size)); | |
| 102 | + | } | |
| 103 | + | ||
| 104 | + | // Handle stdin -> socket (Input) | |
| 105 | + | const inp_flags = (lib_posix.POLL.IN | lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL); | |
| 106 | + | if (poll_fds.items[0].revents & inp_flags != 0) { | |
| 107 | + | var buf: [4096]u8 = undefined; | |
| 108 | + | const n_opt: ?usize = lib_posix.read(stdin_fd, &buf) catch |err| blk: { | |
| 109 | + | if (err == error.WouldBlock) break :blk null; | |
| 110 | + | return err; | |
| 111 | + | }; | |
| 112 | + | ||
| 113 | + | if (n_opt) |n| { | |
| 114 | + | if (n > 0) { | |
| 115 | + | // Check for detach sequences (ctrl+\ as first byte or Kitty escape sequence) | |
| 116 | + | if (util.isCtrlBackslash(buf[0..n])) { | |
| 117 | + | std.log.info("detach key detected", .{}); | |
| 118 | + | try ipc.appendMessage(gpa, &sock_write_buf, .Detach, ""); | |
| 119 | + | } else { | |
| 120 | + | try ipc.appendMessage(gpa, &sock_write_buf, .Input, buf[0..n]); | |
| 121 | + | } | |
| 122 | + | } else { | |
| 123 | + | std.log.info("eof stdin", .{}); | |
| 124 | + | // EOF on stdin | |
| 125 | + | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 126 | + | } | |
| 127 | + | } | |
| 128 | + | } | |
| 129 | + | ||
| 130 | + | // Handle socket read (incoming Output messages from daemon) | |
| 131 | + | if (poll_fds.items[1].revents & lib_posix.POLL.IN != 0) { | |
| 132 | + | const n = read_buf.read(client_sock_fd) catch |err| { | |
| 133 | + | if (err == error.WouldBlock) continue; | |
| 134 | + | if (err == error.ConnectionResetByPeer or err == error.BrokenPipe) { | |
| 135 | + | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 136 | + | } | |
| 137 | + | std.log.err("daemon read err={s}", .{@errorName(err)}); | |
| 138 | + | return err; | |
| 139 | + | }; | |
| 140 | + | if (n == 0) { | |
| 141 | + | std.log.info("server closed connection", .{}); | |
| 142 | + | // Server closed connection | |
| 143 | + | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 144 | + | } | |
| 145 | + | ||
| 146 | + | while (read_buf.next()) |msg| { | |
| 147 | + | switch (msg.header.tag) { | |
| 148 | + | .Output => { | |
| 149 | + | if (msg.payload.len > 0) { | |
| 150 | + | try stdout_buf.appendSlice(gpa, msg.payload); | |
| 151 | + | } | |
| 152 | + | }, | |
| 153 | + | .Resize => { | |
| 154 | + | // daemon is asking for the client's window size usually in response | |
| 155 | + | // to this client being set as leader. | |
| 156 | + | const next_size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO); | |
| 157 | + | try ipc.appendMessage( | |
| 158 | + | gpa, | |
| 159 | + | &sock_write_buf, | |
| 160 | + | .Resize, | |
| 161 | + | std.mem.asBytes(&next_size), | |
| 162 | + | ); | |
| 163 | + | }, | |
| 164 | + | .Switch => { | |
| 165 | + | std.log.info("switch session", .{}); | |
| 166 | + | return ClientResult{ .kind = .switch_session, .session_name = try gpa.dupe(u8, msg.payload) }; | |
| 167 | + | }, | |
| 168 | + | else => {}, | |
| 169 | + | } | |
| 170 | + | } | |
| 171 | + | } | |
| 172 | + | ||
| 173 | + | // Handle socket write (flush buffered messages to daemon) | |
| 174 | + | if (poll_fds.items[1].revents & lib_posix.POLL.OUT != 0) { | |
| 175 | + | if (sock_write_buf.items.len > 0) { | |
| 176 | + | const n = lib_posix.write(client_sock_fd, sock_write_buf.items) catch |err| blk: { | |
| 177 | + | if (err == error.WouldBlock) break :blk 0; | |
| 178 | + | if (err == error.ConnectionResetByPeer or err == error.BrokenPipe) { | |
| 179 | + | std.log.info("connection reset or broken pipe", .{}); | |
| 180 | + | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 181 | + | } | |
| 182 | + | return err; | |
| 183 | + | }; | |
| 184 | + | if (n > 0) { | |
| 185 | + | try sock_write_buf.replaceRange(gpa, 0, n, &[_]u8{}); | |
| 186 | + | } | |
| 187 | + | } | |
| 188 | + | } | |
| 189 | + | ||
| 190 | + | if (stdout_buf.items.len > 0) { | |
| 191 | + | const n = lib_posix.write(lib_posix.STDOUT_FILENO, stdout_buf.items) catch |err| blk: { | |
| 192 | + | if (err == error.WouldBlock) break :blk 0; | |
| 193 | + | return err; | |
| 194 | + | }; | |
| 195 | + | if (n > 0) { | |
| 196 | + | try stdout_buf.replaceRange(gpa, 0, n, &[_]u8{}); | |
| 197 | + | } | |
| 198 | + | } | |
| 199 | + | ||
| 200 | + | if (poll_fds.items[1].revents & (lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL) != 0) { | |
| 201 | + | std.log.info("poll hup|err|nval", .{}); | |
| 202 | + | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 203 | + | } | |
| 204 | + | } | |
| 205 | + | } | |
| 206 | + | ||
| 207 | + | /// dameonLoop is what the daemon runs to send and receive ipc commands from its corresponding | |
| 208 | + | /// clients. It uses poll() as its non-blocking mechanism. | |
| 209 | + | fn daemonLoop(daemon: *Daemon, gpa: std.mem.Allocator, io: std.Io, server_sock_fd: lib_posix.socket_t, pty_fd: i32) !void { | |
| 210 | + | std.log.info("daemon started session={s} pty_fd={d}", .{ daemon.session_name, pty_fd }); | |
| 211 | + | ||
| 212 | + | try signal.openSignalPipe(); | |
| 213 | + | signal.installWakeHandler(@intFromEnum(lib_posix.SIG.TERM)); | |
| 214 | + | var poll_fds = try std.ArrayList(lib_posix.pollfd).initCapacity(gpa, 8); | |
| 215 | + | defer poll_fds.deinit(gpa); | |
| 216 | + | ||
| 217 | + | const init_size = ipc.getTerminalSize(pty_fd); | |
| 218 | + | var term = try ghostty_vt.Terminal.init(io, gpa, .{ | |
| 219 | + | .cols = init_size.cols, | |
| 220 | + | .rows = init_size.rows, | |
| 221 | + | .max_scrollback = daemon.cfg.max_scrollback, | |
| 222 | + | }); | |
| 223 | + | defer term.deinit(gpa); | |
| 224 | + | var vt_stream = term.vtStream(); | |
| 225 | + | defer vt_stream.deinit(); | |
| 226 | + | ||
| 227 | + | // Carries the tail of the previous PTY read so the task-exit marker | |
| 228 | + | // search below can see across a read() boundary. Sized to comfortably | |
| 229 | + | // hold "ZMX_TASK_COMPLETED:" (19 bytes) plus a u8 exit code and CRLF. | |
| 230 | + | var marker_carry: [32]u8 = undefined; | |
| 231 | + | var marker_carry_len: usize = 0; | |
| 232 | + | ||
| 233 | + | daemon_loop: while (daemon.running) { | |
| 234 | + | poll_fds.clearRetainingCapacity(); | |
| 235 | + | ||
| 236 | + | try poll_fds.append(gpa, .{ | |
| 237 | + | .fd = server_sock_fd, | |
| 238 | + | .events = lib_posix.POLL.IN, | |
| 239 | + | .revents = 0, | |
| 240 | + | }); | |
| 241 | + | ||
| 242 | + | var pty_events: i16 = lib_posix.POLL.IN; | |
| 243 | + | if (daemon.pty_write_buf.items.len > 0) { | |
| 244 | + | pty_events |= lib_posix.POLL.OUT; | |
| 245 | + | } | |
| 246 | + | try poll_fds.append(gpa, .{ | |
| 247 | + | .fd = pty_fd, | |
| 248 | + | .events = pty_events, | |
| 249 | + | .revents = 0, | |
| 250 | + | }); | |
| 251 | + | ||
| 252 | + | try poll_fds.append(gpa, .{ .fd = signal.sig_pipe[0], .events = lib_posix.POLL.IN, .revents = 0 }); | |
| 253 | + | ||
| 254 | + | for (daemon.clients.items) |client| { | |
| 255 | + | var events: i16 = lib_posix.POLL.IN; | |
| 256 | + | if (client.has_pending_output) { | |
| 257 | + | events |= lib_posix.POLL.OUT; | |
| 258 | + | } | |
| 259 | + | try poll_fds.append(gpa, .{ | |
| 260 | + | .fd = client.socket_fd, | |
| 261 | + | .events = events, | |
| 262 | + | .revents = 0, | |
| 263 | + | }); | |
| 264 | + | } | |
| 265 | + | ||
| 266 | + | _ = try lib_posix.poll(poll_fds.items, -1); | |
| 267 | + | ||
| 268 | + | if (poll_fds.items[2].revents & lib_posix.POLL.IN != 0) { | |
| 269 | + | signal.drainSignalPipe(); | |
| 270 | + | std.log.info( | |
| 271 | + | "SIGTERM received, shutting down gracefully session={s}", | |
| 272 | + | .{daemon.session_name}, | |
| 273 | + | ); | |
| 274 | + | break :daemon_loop; | |
| 275 | + | } | |
| 276 | + | ||
| 277 | + | if (poll_fds.items[0].revents & (lib_posix.POLL.ERR | lib_posix.POLL.HUP | lib_posix.POLL.NVAL) != 0) { | |
| 278 | + | std.log.err("server socket error revents={d}", .{poll_fds.items[0].revents}); | |
| 279 | + | break :daemon_loop; | |
| 280 | + | } else if (poll_fds.items[0].revents & lib_posix.POLL.IN != 0) { | |
| 281 | + | const client_fd = try lib_posix.accept( | |
| 282 | + | server_sock_fd, | |
| 283 | + | null, | |
| 284 | + | null, | |
| 285 | + | lib_posix.SOCK.NONBLOCK | lib_posix.SOCK.CLOEXEC, | |
| 286 | + | ); | |
| 287 | + | const client = try gpa.create(Client); | |
| 288 | + | client.* = Client{ | |
| 289 | + | .alloc = gpa, | |
| 290 | + | .socket_fd = client_fd, | |
| 291 | + | .read_buf = try ipc.SocketBuffer.init(gpa), | |
| 292 | + | .write_buf = undefined, | |
| 293 | + | }; | |
| 294 | + | // 64KB initial capacity lets ~15 broadcast cycles (N_TTY_BUF_SIZE reads | |
| 295 | + | // * header) accumulate before the first ArrayList growth. The write | |
| 296 | + | // buffer is userspace-only: it drains via POLLOUT to the client socket, | |
| 297 | + | // which has no corresponding kernel-imposed per-write limit. | |
| 298 | + | client.write_buf = try std.ArrayList(u8).initCapacity(client.alloc, 65536); | |
| 299 | + | try daemon.clients.append(gpa, client); | |
| 300 | + | std.log.info( | |
| 301 | + | "client connected fd={d} total={d}", | |
| 302 | + | .{ client_fd, daemon.clients.items.len }, | |
| 303 | + | ); | |
| 304 | + | } | |
| 305 | + | ||
| 306 | + | const inp_flags = lib_posix.POLL.IN | lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL; | |
| 307 | + | if (poll_fds.items[1].revents & inp_flags != 0) { | |
| 308 | + | // Read from PTY. Buffer is sized to N_TTY_BUF_SIZE (4096): the hard | |
| 309 | + | // kernel limit for the N_TTY line discipline. A larger buffer doesn't | |
| 310 | + | // help: each read() from a PTY master returns at most 4096 bytes | |
| 311 | + | // regardless of the userspace buffer size. | |
| 312 | + | var buf: [4096]u8 = undefined; | |
| 313 | + | const n_opt: ?usize = lib_posix.read(pty_fd, &buf) catch |err| blk: { | |
| 314 | + | if (err == error.WouldBlock) break :blk null; | |
| 315 | + | break :blk 0; | |
| 316 | + | }; | |
| 317 | + | ||
| 318 | + | if (n_opt) |n| { | |
| 319 | + | if (n == 0) { | |
| 320 | + | // EOF: Shell exited | |
| 321 | + | std.log.info("shell exited pty_fd={d}", .{pty_fd}); | |
| 322 | + | // Let the rest of this poll iteration complete so client | |
| 323 | + | // write buffers are flushed via the normal POLLOUT path. | |
| 324 | + | // On the next iteration, daemon.running will be false. | |
| 325 | + | daemon.running = false; | |
| 326 | + | } else { | |
| 327 | + | // Feed PTY output to terminal emulator for state tracking | |
| 328 | + | vt_stream.nextSlice(buf[0..n]); | |
| 329 | + | daemon.has_pty_output = true; | |
| 330 | + | ||
| 331 | + | // When no real terminal client has attached yet, respond to | |
| 332 | + | // terminal queries (e.g. DA1/DA2) on behalf of the terminal. | |
| 333 | + | // This prevents fish from waiting 10s for unanswered queries. | |
| 334 | + | // `has_terminal_client` is only set when a client sends .Init | |
| 335 | + | // (a real zmx attach), not when a `zmx run` tail-only client | |
| 336 | + | // connects. | |
| 337 | + | if (!daemon.has_terminal_client and | |
| 338 | + | daemon.pty_write_buf.items.len < Daemon.PTY_WRITE_BUF_MAX) | |
| 339 | + | { | |
| 340 | + | util.respondToDeviceAttributes(gpa, &daemon.pty_write_buf, buf[0..n]); | |
| 341 | + | } | |
| 342 | + | ||
| 343 | + | // In run mode, scan output for exit code marker. The marker | |
| 344 | + | // can straddle two PTY reads (more likely under a throttled | |
| 345 | + | // scheduler, e.g. containers), so prepend the tail carried | |
| 346 | + | // over from the previous read before searching. | |
| 347 | + | if (daemon.is_task_mode and daemon.task_exit_code == null) { | |
| 348 | + | var scan_buf: [marker_carry.len + buf.len]u8 = undefined; | |
| 349 | + | @memcpy(scan_buf[0..marker_carry_len], marker_carry[0..marker_carry_len]); | |
| 350 | + | @memcpy(scan_buf[marker_carry_len..][0..n], buf[0..n]); | |
| 351 | + | const scan_len = marker_carry_len + n; | |
| 352 | + | ||
| 353 | + | if (util.findTaskExitMarker(scan_buf[0..scan_len])) |exit_code| { | |
| 354 | + | daemon.task_exit_code = exit_code; | |
| 355 | + | daemon.task_ended_at = @intCast(std.Io.Timestamp.now(io, .real).toSeconds()); | |
| 356 | + | ||
| 357 | + | std.log.info("task completed exit_code={d}", .{exit_code}); | |
| 358 | + | ||
| 359 | + | // Notify connected clients | |
| 360 | + | for (daemon.clients.items) |c| { | |
| 361 | + | ipc.appendMessage(gpa, &c.write_buf, .TaskComplete, &[_]u8{exit_code}) catch {}; | |
| 362 | + | c.has_pending_output = true; | |
| 363 | + | } | |
| 364 | + | } | |
| 365 | + | ||
| 366 | + | marker_carry_len = @min(marker_carry.len, scan_len); | |
| 367 | + | @memcpy( | |
| 368 | + | marker_carry[0..marker_carry_len], | |
| 369 | + | scan_buf[scan_len - marker_carry_len .. scan_len], | |
| 370 | + | ); | |
| 371 | + | } | |
| 372 | + | ||
| 373 | + | // Broadcast data to all clients. | |
| 374 | + | // Rewrite OSC 133;A to include redraw=0 so the outer terminal | |
| 375 | + | // does not clear prompt lines on resize (issue #111). | |
| 376 | + | const broadcast_data = util.rewritePromptRedraw(gpa, buf[0..n]) orelse buf[0..n]; | |
| 377 | + | defer if (broadcast_data.ptr != buf[0..n].ptr) gpa.free(broadcast_data); | |
| 378 | + | for (daemon.clients.items) |client| { | |
| 379 | + | ipc.appendMessage(gpa, &client.write_buf, .Output, broadcast_data) catch |err| { | |
| 380 | + | std.log.warn( | |
| 381 | + | "failed to buffer output for client err={s}", | |
| 382 | + | .{@errorName(err)}, | |
| 383 | + | ); | |
| 384 | + | continue; | |
| 385 | + | }; | |
| 386 | + | client.has_pending_output = true; | |
| 387 | + | } | |
| 388 | + | } | |
| 389 | + | } | |
| 390 | + | } | |
| 391 | + | ||
| 392 | + | if (poll_fds.items[1].revents & lib_posix.POLL.OUT != 0) { | |
| 393 | + | while (daemon.pty_write_buf.items.len > 0) { | |
| 394 | + | const n = lib_posix.write(pty_fd, daemon.pty_write_buf.items) catch |err| { | |
| 395 | + | if (err != error.WouldBlock) { | |
| 396 | + | std.log.warn("pty write failed: {s}", .{@errorName(err)}); | |
| 397 | + | daemon.pty_write_buf.clearRetainingCapacity(); | |
| 398 | + | } | |
| 399 | + | break; | |
| 400 | + | }; | |
| 401 | + | if (n == 0) break; | |
| 402 | + | daemon.pty_write_buf.replaceRange(gpa, 0, n, &[_]u8{}) catch unreachable; | |
| 403 | + | } | |
| 404 | + | } | |
| 405 | + | ||
| 406 | + | var i: usize = daemon.clients.items.len; | |
| 407 | + | // Only iterate over clients that were present when poll_fds was constructed | |
| 408 | + | // poll_fds contains [server, pty, sig_pipe, client0, client1, ...] | |
| 409 | + | // So number of clients in poll_fds is poll_fds.items.len - 3 | |
| 410 | + | const num_polled_clients = poll_fds.items.len - 3; | |
| 411 | + | if (i > num_polled_clients) { | |
| 412 | + | // If we have more clients than polled (i.e. we just accepted one), start from the | |
| 413 | + | // polled ones | |
| 414 | + | i = num_polled_clients; | |
| 415 | + | } | |
| 416 | + | ||
| 417 | + | clients_loop: while (i > 0) { | |
| 418 | + | i -= 1; | |
| 419 | + | const client = daemon.clients.items[i]; | |
| 420 | + | const revents = poll_fds.items[i + 3].revents; | |
| 421 | + | ||
| 422 | + | if (revents & lib_posix.POLL.IN != 0) { | |
| 423 | + | const n = client.read_buf.read(client.socket_fd) catch |err| { | |
| 424 | + | if (err == error.WouldBlock) continue; | |
| 425 | + | std.log.debug( | |
| 426 | + | "client read err={s} fd={d}", | |
| 427 | + | .{ @errorName(err), client.socket_fd }, | |
| 428 | + | ); | |
| 429 | + | const last = daemon.closeClient(gpa, client, i, false); | |
| 430 | + | if (last) break :daemon_loop; | |
| 431 | + | continue; | |
| 432 | + | }; | |
| 433 | + | ||
| 434 | + | if (n == 0) { | |
| 435 | + | // Client closed connection | |
| 436 | + | const last = daemon.closeClient(gpa, client, i, false); | |
| 437 | + | if (last) break :daemon_loop; | |
| 438 | + | continue; | |
| 439 | + | } | |
| 440 | + | ||
| 441 | + | while (client.read_buf.next()) |msg| { | |
| 442 | + | switch (msg.header.tag) { | |
| 443 | + | .Input => try daemon.handleInput(gpa, client, msg.payload), | |
| 444 | + | .Send => daemon.handleSend(gpa, msg.payload), | |
| 445 | + | .Output => try daemon.handleOutput(gpa, msg.payload, &vt_stream), | |
| 446 | + | .Init => try daemon.handleInit(gpa, client, pty_fd, &term, msg.payload), | |
| 447 | + | .Switch => try daemon.handleSwitch(gpa, msg.payload), | |
| 448 | + | .Resize => try daemon.handleResize(gpa, client, pty_fd, &term, msg.payload), | |
| 449 | + | .Detach => { | |
| 450 | + | daemon.handleDetach(gpa, client, i); | |
| 451 | + | break :clients_loop; | |
| 452 | + | }, | |
| 453 | + | .DetachAll => { | |
| 454 | + | daemon.handleDetachAll(gpa); | |
| 455 | + | break :clients_loop; | |
| 456 | + | }, | |
| 457 | + | .Kill => { | |
| 458 | + | break :daemon_loop; | |
| 459 | + | }, | |
| 460 | + | .Info => try daemon.handleInfo(gpa, client), | |
| 461 | + | .LabelGet => try daemon.handleLabelGet(gpa, client), | |
| 462 | + | .LabelSet => try daemon.handleLabelSet(gpa, client, msg.payload), | |
| 463 | + | .LabelClear => try daemon.handleLabelClear(gpa, client), | |
| 464 | + | .History => try daemon.handleHistory(gpa, client, &term, msg.payload), | |
| 465 | + | .Run => try daemon.handleRun(gpa, client, msg.payload), | |
| 466 | + | .Ack, .TaskComplete, .LabelData => {}, | |
| 467 | + | .Write => try daemon.handleWrite(gpa, client, msg.payload), | |
| 468 | + | _ => std.log.warn( | |
| 469 | + | "ignoring unknown IPC tag={d}", | |
| 470 | + | .{@intFromEnum(msg.header.tag)}, | |
| 471 | + | ), | |
| 472 | + | } | |
| 473 | + | } | |
| 474 | + | } | |
| 475 | + | ||
| 476 | + | if (revents & lib_posix.POLL.OUT != 0) { | |
| 477 | + | // Flush pending output buffers | |
| 478 | + | const n = lib_posix.write(client.socket_fd, client.write_buf.items) catch |err| blk: { | |
| 479 | + | if (err == error.WouldBlock) break :blk 0; | |
| 480 | + | // Error on write, close client | |
| 481 | + | const last = daemon.closeClient(gpa, client, i, false); | |
| 482 | + | if (last) break :daemon_loop; | |
| 483 | + | continue; | |
| 484 | + | }; | |
| 485 | + | ||
| 486 | + | if (n > 0) { | |
| 487 | + | client.write_buf.replaceRange(gpa, 0, n, &[_]u8{}) catch unreachable; | |
| 488 | + | } | |
| 489 | + | ||
| 490 | + | if (client.write_buf.items.len == 0) { | |
| 491 | + | client.has_pending_output = false; | |
| 492 | + | } | |
| 493 | + | } | |
| 494 | + | ||
| 495 | + | if (revents & (lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL) != 0) { | |
| 496 | + | const last = daemon.closeClient(gpa, client, i, false); | |
| 497 | + | if (last) break :daemon_loop; | |
| 498 | + | } | |
| 499 | + | } | |
| 500 | + | } | |
| 501 | + | } | |
| 502 | + | ||
| 503 | + | const ClientResult = struct { | |
| 504 | + | kind: enum { | |
| 505 | + | detach, | |
| 506 | + | switch_session, | |
| 507 | + | }, | |
| 508 | + | session_name: ?[]const u8, | |
| 509 | + | }; | |
| 510 | + | ||
| 511 | + | /// Client represents each terminal that has connected to a session. | |
| 512 | + | /// | |
| 513 | + | /// Multiple Clients can connect to a single session. | |
| 514 | + | pub const Client = struct { | |
| 515 | + | alloc: std.mem.Allocator, | |
| 516 | + | socket_fd: i32, | |
| 517 | + | has_pending_output: bool = false, | |
| 518 | + | read_buf: ipc.SocketBuffer, | |
| 519 | + | write_buf: std.ArrayList(u8), | |
| 520 | + | ||
| 521 | + | pub fn deinit(self: *Client) void { | |
| 522 | + | lib_posix.close(self.socket_fd); | |
| 523 | + | self.read_buf.deinit(); | |
| 524 | + | self.write_buf.deinit(self.alloc); | |
| 525 | + | } | |
| 526 | + | }; | |
| 527 | + | ||
| 528 | + | /// Daemon is responsible for managing a zmx session. | |
| 529 | + | /// | |
| 530 | + | /// It holds all the state for a running session. Instead of a single daemon for all sessions, we | |
| 531 | + | /// create a daemon for every session. This has some benefits. The ipc communication between | |
| 532 | + | /// session clients and the daemon doesn't need to be tagged with the session name. If a daemon | |
| 533 | + | /// crashes for one session won't crash all the other sessions. | |
| 534 | + | /// | |
| 535 | + | /// Conceptually it's also much simpler to reason about. | |
| 536 | + | pub const Daemon = struct { | |
| 537 | + | cfg: *Cfg, | |
| 538 | + | session_name: []const u8, | |
| 539 | + | socket_path: []const u8, | |
| 540 | + | // === opt === | |
| 541 | + | pty_write_buf: std.ArrayList(u8) = .empty, | |
| 542 | + | clients: std.ArrayList(*Client) = .empty, | |
| 543 | + | labels: std.StringHashMapUnmanaged([]u8) = .empty, | |
| 544 | + | // This control which client is the leader. The leader controls terminal state and | |
| 545 | + | // cols/rows of session. | |
| 546 | + | leader_client_fd: ?i32 = null, | |
| 547 | + | running: bool = true, | |
| 548 | + | pid: i32 = undefined, | |
| 549 | + | command: ?[]const []const u8 = null, | |
| 550 | + | cwd: []const u8 = "", | |
| 551 | + | has_pty_output: bool = false, | |
| 552 | + | has_had_client: bool = false, | |
| 553 | + | has_terminal_client: bool = false, // true only after a real attach (.Init received) | |
| 554 | + | created_at: u64, // unix timestamp (ns) | |
| 555 | + | is_task_mode: bool = false, // flag for when session is run as a task | |
| 556 | + | task_exit_code: ?u8 = null, // null = running or n/a, set when task completes | |
| 557 | + | task_ended_at: ?u64 = null, // timestamp when task exited | |
| 558 | + | pty_fd: i32 = -1, // set by daemonLoop so handleRun can probe the foreground process | |
| 559 | + | shell: []const u8 = "/bin/sh", | |
| 560 | + | ||
| 561 | + | /// Create a Daemon. Caller is responsible for freeing all variables passed | |
| 562 | + | /// into the init fn. | |
| 563 | + | pub fn init(io: std.Io, cfg: *Cfg, sesh_name: []const u8, socket_path: []const u8) Daemon { | |
| 564 | + | return .{ | |
| 565 | + | .cfg = cfg, | |
| 566 | + | .session_name = sesh_name, | |
| 567 | + | .socket_path = socket_path, | |
| 568 | + | .created_at = @intCast(std.Io.Timestamp.now(io, .real).toSeconds()), | |
| 569 | + | }; | |
| 570 | + | } | |
| 571 | + | ||
| 572 | + | pub fn deinit(self: *Daemon, gpa: std.mem.Allocator) void { | |
| 573 | + | self.clients.deinit(gpa); | |
| 574 | + | var it = self.labels.iterator(); | |
| 575 | + | while (it.next()) |entry| { | |
| 576 | + | gpa.free(entry.key_ptr.*); | |
| 577 | + | gpa.free(entry.value_ptr.*); | |
| 578 | + | } | |
| 579 | + | self.labels.deinit(gpa); | |
| 580 | + | self.pty_write_buf.deinit(gpa); | |
| 581 | + | gpa.free(self.socket_path); | |
| 582 | + | } | |
| 583 | + | ||
| 584 | + | pub fn shutdown(self: *Daemon, gpa: std.mem.Allocator) void { | |
| 585 | + | std.log.info("shutting down daemon session={s}", .{self.session_name}); | |
| 586 | + | self.running = false; | |
| 587 | + | ||
| 588 | + | for (self.clients.items) |client| { | |
| 589 | + | client.deinit(); | |
| 590 | + | gpa.destroy(client); | |
| 591 | + | } | |
| 592 | + | self.clients.clearRetainingCapacity(); | |
| 593 | + | } | |
| 594 | + | ||
| 595 | + | pub fn closeClient(self: *Daemon, gpa: std.mem.Allocator, client: *Client, i: usize, shutdown_on_last: bool) bool { | |
| 596 | + | const fd = client.socket_fd; | |
| 597 | + | // leader is disconnected, remove ref and let another client claim leader on input | |
| 598 | + | if (self.leader_client_fd == client.socket_fd) { | |
| 599 | + | std.log.info( | |
| 600 | + | "unsetting leader session={s} fd={d}", | |
| 601 | + | .{ self.session_name, client.socket_fd }, | |
| 602 | + | ); | |
| 603 | + | self.leader_client_fd = null; | |
| 604 | + | } | |
| 605 | + | client.deinit(); | |
| 606 | + | gpa.destroy(client); | |
| 607 | + | _ = self.clients.orderedRemove(i); | |
| 608 | + | std.log.info("client disconnected fd={d} remaining={d}", .{ fd, self.clients.items.len }); | |
| 609 | + | if (shutdown_on_last and self.clients.items.len == 0) { | |
| 610 | + | self.shutdown(gpa); | |
| 611 | + | return true; | |
| 612 | + | } | |
| 613 | + | return false; | |
| 614 | + | } | |
| 615 | + | ||
| 616 | + | /// ensureSession will either create or re-use the daemon used for a session. | |
| 617 | + | /// It will spin up a unix socket, double-fork the process (so it survives | |
| 618 | + | /// the terminal dying), and automatically attach the client to the ipc unix | |
| 619 | + | /// socket. | |
| 620 | + | /// | |
| 621 | + | /// The return bool value indicates if the current process is the daemon | |
| 622 | + | /// or the client since they have different behaviors post-fork. | |
| 623 | + | /// | |
| 624 | + | /// E.g. If it's the client process then we need to connect to the unix socket | |
| 625 | + | /// and run the clientLoop. If it's the daemon then we need to bail since | |
| 626 | + | /// the daemonLoop is created inside this fn and when it returns that means | |
| 627 | + | /// the daemon stopped and needs to exit. | |
| 628 | + | pub fn ensureSession(self: *Daemon, io: std.Io) !bool { | |
| 629 | + | const sesh_name = self.session_name; | |
| 630 | + | std.log.info("ensure session session={s}", .{sesh_name}); | |
| 631 | + | var dir = try std.Io.Dir.openDirAbsolute(io, self.cfg.socket_dir, .{}); | |
| 632 | + | defer dir.close(io); | |
| 633 | + | ||
| 634 | + | const exists = try socket.sessionExists(io, dir, sesh_name); | |
| 635 | + | // if daemon is gone then we flip this to true | |
| 636 | + | var should_create = !exists; | |
| 637 | + | ||
| 638 | + | if (exists) { | |
| 639 | + | if (ipc.connectSession(self.socket_path)) |fd| { | |
| 640 | + | lib_posix.close(fd); | |
| 641 | + | if (self.command != null) { | |
| 642 | + | std.log.warn( | |
| 643 | + | "session already exists, ignoring command session={s}", | |
| 644 | + | .{sesh_name}, | |
| 645 | + | ); | |
| 646 | + | } | |
| 647 | + | } else |err| switch (err) { | |
| 648 | + | // Daemon is definitively gone: safe to replace. | |
| 649 | + | error.ConnectionRefused => { | |
| 650 | + | socket.cleanupStaleSocket(io, dir, sesh_name); | |
| 651 | + | should_create = true; | |
| 652 | + | }, | |
| 653 | + | // Connect failed for an unusual reason. The check is only to | |
| 654 | + | // decide create-vs-attach; the socket file exists, so proceed | |
| 655 | + | // to attach rather than fail or orphan. | |
| 656 | + | else => { | |
| 657 | + | std.log.warn( | |
| 658 | + | "connect failed ({s}), proceeding to attach session={s}", | |
| 659 | + | .{ @errorName(err), sesh_name }, | |
| 660 | + | ); | |
| 661 | + | }, | |
| 662 | + | } | |
| 663 | + | } | |
| 664 | + | ||
| 665 | + | if (!should_create) { | |
| 666 | + | return false; | |
| 667 | + | } | |
| 668 | + | ||
| 669 | + | return self.run(io, dir, sesh_name); | |
| 670 | + | } | |
| 671 | + | ||
| 672 | + | fn run(self: *Daemon, io: std.Io, dir: std.Io.Dir, sesh_name: []const u8) !bool { | |
| 673 | + | std.log.info("creating session={s}", .{sesh_name}); | |
| 674 | + | const server_sock_fd: lib_posix.socket_t = try socket.createSocket(self.socket_path); | |
| 675 | + | const log_fd = log.log_system.file.?.handle; | |
| 676 | + | ||
| 677 | + | var keep_fds_open = [_]i32{ server_sock_fd, dir.handle, log_fd }; | |
| 678 | + | const cmd = try daemonize.createCmdZ(self.shell, self.is_task_mode, self.command); | |
| 679 | + | const pty_info = daemonize.daemonize( | |
| 680 | + | sesh_name, | |
| 681 | + | cmd, | |
| 682 | + | &keep_fds_open, | |
| 683 | + | ) catch |err| { | |
| 684 | + | switch (err) { | |
| 685 | + | error.IsClientProc => { | |
| 686 | + | // send a msg to the client that the session was created. | |
| 687 | + | var w_buf: [2048]u8 = undefined; | |
| 688 | + | var w = std.Io.File.stdout().writer(io, &w_buf); | |
| 689 | + | try w.interface.print("session \"{s}\" created\n", .{sesh_name}); | |
| 690 | + | try w.interface.flush(); | |
| 691 | + | lib_posix.close(server_sock_fd); | |
| 692 | + | return false; | |
| 693 | + | }, | |
| 694 | + | else => { | |
| 695 | + | lib_posix.close(server_sock_fd); | |
| 696 | + | dir.deleteFile(io, self.session_name) catch {}; | |
| 697 | + | return err; | |
| 698 | + | }, | |
| 699 | + | } | |
| 700 | + | }; | |
| 701 | + | // ======= | |
| 702 | + | // WARNING: cannot use upstream allocator or io after this point since | |
| 703 | + | // we forked the process and there's a risk of a mutex (e.g. thread-safe | |
| 704 | + | // allocator) being locked by a thread prior to fork which can cause a | |
| 705 | + | // deadlock. | |
| 706 | + | // ======= | |
| 707 | + | ||
| 708 | + | self.pid = pty_info.pid; | |
| 709 | + | ||
| 710 | + | var threaded: std.Io.Threaded = .init_single_threaded; | |
| 711 | + | defer threaded.deinit(); | |
| 712 | + | const new_io = threaded.io(); | |
| 713 | + | ||
| 714 | + | { // re-initialize logs with the session name as the filename | |
| 715 | + | log.log_system.deinit(); | |
| 716 | + | var log_buf: [4096]u8 = undefined; | |
| 717 | + | const session_log_name = try std.fmt.bufPrint( | |
| 718 | + | &log_buf, | |
| 719 | + | "{s}.log", | |
| 720 | + | .{sesh_name}, | |
| 721 | + | ); | |
| 722 | + | var fba_buf: [4096]u8 = undefined; | |
| 723 | + | var fba = std.heap.FixedBufferAllocator.init(&fba_buf); | |
| 724 | + | const session_log_path = try std.fs.path.join( | |
| 725 | + | fba.allocator(), | |
| 726 | + | &.{ self.cfg.log_dir, session_log_name }, | |
| 727 | + | ); | |
| 728 | + | const log_mode = std.Io.File.Permissions.fromMode(self.cfg.log_mode); | |
| 729 | + | log.log_system.init(new_io, session_log_path, log_mode) catch {}; | |
| 730 | + | } | |
| 731 | + | ||
| 732 | + | const gpa: std.mem.Allocator = blk: { | |
| 733 | + | if (builtin.mode == .Debug) { | |
| 734 | + | const GPA = std.heap.DebugAllocator(.{}); | |
| 735 | + | const Static = struct { | |
| 736 | + | var gpa: GPA = .{}; | |
| 737 | + | }; | |
| 738 | + | break :blk Static.gpa.allocator(); | |
| 739 | + | } | |
| 740 | + | break :blk std.heap.c_allocator; | |
| 741 | + | }; | |
| 742 | + | ||
| 743 | + | defer { | |
| 744 | + | // Close and unlink the listen socket BEFORE handleKill()'s | |
| 745 | + | // 500ms SIGHUP->SIGKILL grace sleep. Otherwise a `zmx run` | |
| 746 | + | // for the same name issued in that window will hang waiting | |
| 747 | + | // for a connect. | |
| 748 | + | lib_posix.close(server_sock_fd); | |
| 749 | + | std.log.info("deleting socket file session={s}", .{sesh_name}); | |
| 750 | + | dir.deleteFile(new_io, sesh_name) catch |err| { | |
| 751 | + | std.log.warn("failed to delete socket file err={s}", .{@errorName(err)}); | |
| 752 | + | }; | |
| 753 | + | self.handleKill(gpa, new_io); | |
| 754 | + | self.deinit(gpa); | |
| 755 | + | lib_posix.close(pty_info.master_fd); | |
| 756 | + | _ = lib_posix.waitpid(self.pid, 0); | |
| 757 | + | } | |
| 758 | + | ||
| 759 | + | try daemonLoop(self, gpa, new_io, server_sock_fd, pty_info.master_fd); | |
| 760 | + | std.log.info("daemon loop shutdown", .{}); | |
| 761 | + | return true; | |
| 762 | + | } | |
| 763 | + | ||
| 764 | + | fn setLeader(self: *Daemon, gpa: std.mem.Allocator, client: *Client) !void { | |
| 765 | + | std.log.info("setting new leader client_fd={d}", .{client.socket_fd}); | |
| 766 | + | self.leader_client_fd = client.socket_fd; | |
| 767 | + | // Send a resize message to the client so it can send us back their window size | |
| 768 | + | // so we can resize the pty and ghostty state. | |
| 769 | + | try ipc.appendMessage(gpa, &client.write_buf, .Resize, ""); | |
| 770 | + | client.has_pending_output = true; | |
| 771 | + | } | |
| 772 | + | ||
| 773 | + | const PTY_WRITE_BUF_MAX = 256 * 1024; | |
| 774 | + | ||
| 775 | + | /// Queue bytes for the PTY's stdin. Flushed by daemonLoop on POLLOUT. | |
| 776 | + | /// Drops the payload if the buffer is over cap -- same failure mode as | |
| 777 | + | /// the old direct-write ptyWrite (drop on EAGAIN), just at a 64x higher | |
| 778 | + | /// threshold. Capping avoids OOM when the shell stops reading; dropping | |
| 779 | + | /// new (not old) bytes avoids tearing a partially-accepted sequence. | |
| 780 | + | fn queuePtyInput(self: *Daemon, gpa: std.mem.Allocator, data: []const u8) void { | |
| 781 | + | if (data.len == 0) return; | |
| 782 | + | if (self.pty_write_buf.items.len + data.len > PTY_WRITE_BUF_MAX) { | |
| 783 | + | std.log.warn( | |
| 784 | + | "pty input dropped {d} bytes (buffer full, shell not reading)", | |
| 785 | + | .{data.len}, | |
| 786 | + | ); | |
| 787 | + | return; | |
| 788 | + | } | |
| 789 | + | std.log.debug("buffering pty input data={x}", .{data}); | |
| 790 | + | self.pty_write_buf.appendSlice(gpa, data) catch |err| { | |
| 791 | + | std.log.warn( | |
| 792 | + | "pty input dropped {d} bytes: {s}", | |
| 793 | + | .{ data.len, @errorName(err) }, | |
| 794 | + | ); | |
| 795 | + | }; | |
| 796 | + | } | |
| 797 | + | ||
| 798 | + | pub fn handleInput(self: *Daemon, gpa: std.mem.Allocator, client: *Client, payload: []const u8) !void { | |
| 799 | + | std.log.debug("buffering pty input data={x}", .{payload}); | |
| 800 | + | // client is leader, send entire payload (ansi escape codes + text) | |
| 801 | + | if (self.leader_client_fd == client.socket_fd) { | |
| 802 | + | self.queuePtyInput(gpa, payload); | |
| 803 | + | return; | |
| 804 | + | } | |
| 805 | + | ||
| 806 | + | // check if leader needs to be updated by detecting any user input | |
| 807 | + | if (util.isUserInput(payload)) { | |
| 808 | + | try self.setLeader(gpa, client); | |
| 809 | + | self.queuePtyInput(gpa, payload); | |
| 810 | + | } | |
| 811 | + | } | |
| 812 | + | ||
| 813 | + | /// Queue input from `zmx send` without changing interactive client leadership. | |
| 814 | + | pub fn handleSend(self: *Daemon, gpa: std.mem.Allocator, payload: []const u8) void { | |
| 815 | + | self.queuePtyInput(gpa, payload); | |
| 816 | + | } | |
| 817 | + | ||
| 818 | + | pub fn handleSwitch(self: *Daemon, gpa: std.mem.Allocator, session_name: []const u8) !void { | |
| 819 | + | for (self.clients.items) |client| { | |
| 820 | + | if (self.leader_client_fd == client.socket_fd) { | |
| 821 | + | ipc.appendMessage( | |
| 822 | + | gpa, | |
| 823 | + | &client.write_buf, | |
| 824 | + | .Switch, | |
| 825 | + | session_name, | |
| 826 | + | ) catch |err| { | |
| 827 | + | std.log.warn( | |
| 828 | + | "failed to buffer terminal state for client err={s}", | |
| 829 | + | .{@errorName(err)}, | |
| 830 | + | ); | |
| 831 | + | }; | |
| 832 | + | client.has_pending_output = true; | |
| 833 | + | return; | |
| 834 | + | } | |
| 835 | + | } | |
| 836 | + | return error.NoLeaderFound; | |
| 837 | + | } | |
| 838 | + | ||
| 839 | + | pub fn handleInit( | |
| 840 | + | self: *Daemon, | |
| 841 | + | gpa: std.mem.Allocator, | |
| 842 | + | client: *Client, | |
| 843 | + | pty_fd: i32, | |
| 844 | + | term: *ghostty_vt.Terminal, | |
| 845 | + | payload: []const u8, | |
| 846 | + | ) !void { | |
| 847 | + | if (payload.len != @sizeOf(ipc.Resize)) return; | |
| 848 | + | ||
| 849 | + | // Serialize terminal state BEFORE resize to capture correct cursor position. | |
| 850 | + | // Resizing triggers reflow which can move the cursor, and the shell's | |
| 851 | + | // SIGWINCH-triggered redraw will run after our snapshot is sent. | |
| 852 | + | // Only serialize on re-attach (has_had_client), not first attach, to avoid | |
| 853 | + | // interfering with shell initialization (DA1 queries, etc.) | |
| 854 | + | if (self.has_pty_output and self.has_had_client) { | |
| 855 | + | const cursor = &term.screens.active.cursor; | |
| 856 | + | std.log.debug( | |
| 857 | + | "cursor before serialize: x={d} y={d} pending_wrap={}", | |
| 858 | + | .{ cursor.x, cursor.y, cursor.pending_wrap }, | |
| 859 | + | ); | |
| 860 | + | if (util.serializeTerminalState(gpa, term)) |term_output| { | |
| 861 | + | std.log.debug("serialize terminal state", .{}); | |
| 862 | + | // Rewrite OSC 133;A to include redraw=0 so the outer terminal | |
| 863 | + | // does not clear prompt lines on resize (issue #111). | |
| 864 | + | const restore_data = util.rewritePromptRedraw(gpa, term_output) orelse term_output; | |
| 865 | + | defer gpa.free(term_output); | |
| 866 | + | defer if (restore_data.ptr != term_output.ptr) gpa.free(restore_data); | |
| 867 | + | ipc.appendMessage(gpa, &client.write_buf, .Output, restore_data) catch |err| { | |
| 868 | + | std.log.warn( | |
| 869 | + | "failed to buffer terminal state for client err={s}", | |
| 870 | + | .{@errorName(err)}, | |
| 871 | + | ); | |
| 872 | + | }; | |
| 873 | + | client.has_pending_output = true; | |
| 874 | + | } | |
| 875 | + | } | |
| 876 | + | ||
| 877 | + | // no leader is set so set one | |
| 878 | + | if (self.leader_client_fd == null) { | |
| 879 | + | try self.setLeader(gpa, client); | |
| 880 | + | } | |
| 881 | + | ||
| 882 | + | // only resize if leader | |
| 883 | + | if (self.leader_client_fd == client.socket_fd) { | |
| 884 | + | const resize = std.mem.bytesToValue(ipc.Resize, payload); | |
| 885 | + | var ws: cross.c.struct_winsize = .{ | |
| 886 | + | .ws_row = resize.rows, | |
| 887 | + | .ws_col = resize.cols, | |
| 888 | + | .ws_xpixel = resize.xpixel, | |
| 889 | + | .ws_ypixel = resize.ypixel, | |
| 890 | + | }; | |
| 891 | + | _ = cross.c.ioctl(pty_fd, cross.c.TIOCSWINSZ, &ws); | |
| 892 | + | // Disable prompt_redraw before resize. The daemon's internal terminal | |
| 893 | + | // would otherwise clear prompt lines expecting the shell to redraw them, | |
| 894 | + | // but the shell's redraw goes to the PTY (forwarded to clients), not to | |
| 895 | + | // this daemon terminal. The clearing corrupts the daemon's snapshot state. | |
| 896 | + | const saved_prompt_redraw = term.flags.shell_redraws_prompt; | |
| 897 | + | term.flags.shell_redraws_prompt = .false; | |
| 898 | + | defer term.flags.shell_redraws_prompt = saved_prompt_redraw; | |
| 899 | + | const opts = ghostty_vt.Terminal.Resize{ | |
| 900 | + | .cols = resize.cols, | |
| 901 | + | .rows = resize.rows, | |
| 902 | + | }; | |
| 903 | + | try term.resize(gpa, opts); | |
| 904 | + | ||
| 905 | + | // Mark that we've had a client init, so subsequent clients get terminal state | |
| 906 | + | self.has_had_client = true; | |
| 907 | + | self.has_terminal_client = true; | |
| 908 | + | ||
| 909 | + | std.log.debug("init resize rows={d} cols={d}", .{ resize.rows, resize.cols }); | |
| 910 | + | } | |
| 911 | + | } | |
| 912 | + | ||
| 913 | + | pub fn handleResize( | |
| 914 | + | self: *Daemon, | |
| 915 | + | gpa: std.mem.Allocator, | |
| 916 | + | client: *Client, | |
| 917 | + | pty_fd: i32, | |
| 918 | + | term: *ghostty_vt.Terminal, | |
| 919 | + | payload: []const u8, | |
| 920 | + | ) !void { | |
| 921 | + | if (payload.len != @sizeOf(ipc.Resize)) return; | |
| 922 | + | if (self.leader_client_fd == null) { | |
| 923 | + | try self.setLeader(gpa, client); | |
| 924 | + | } | |
| 925 | + | // only leader can resize | |
| 926 | + | if (self.leader_client_fd != client.socket_fd) return; | |
| 927 | + | ||
| 928 | + | const resize = std.mem.bytesToValue(ipc.Resize, payload); | |
| 929 | + | var ws: cross.c.struct_winsize = .{ | |
| 930 | + | .ws_row = resize.rows, | |
| 931 | + | .ws_col = resize.cols, | |
| 932 | + | .ws_xpixel = resize.xpixel, | |
| 933 | + | .ws_ypixel = resize.ypixel, | |
| 934 | + | }; | |
| 935 | + | _ = cross.c.ioctl(pty_fd, cross.c.TIOCSWINSZ, &ws); | |
| 936 | + | // Disable prompt_redraw before resize (same rationale as handleInit). | |
| 937 | + | const saved_prompt_redraw = term.flags.shell_redraws_prompt; | |
| 938 | + | term.flags.shell_redraws_prompt = .false; | |
| 939 | + | defer term.flags.shell_redraws_prompt = saved_prompt_redraw; | |
| 940 | + | const opts = ghostty_vt.Terminal.Resize{ | |
| 941 | + | .cols = resize.cols, | |
| 942 | + | .rows = resize.rows, | |
| 943 | + | }; | |
| 944 | + | try term.resize(gpa, opts); | |
| 945 | + | std.log.debug("resize rows={d} cols={d}", .{ resize.rows, resize.cols }); | |
| 946 | + | } | |
| 947 | + | ||
| 948 | + | pub fn handleDetach(self: *Daemon, gpa: std.mem.Allocator, client: *Client, i: usize) void { | |
| 949 | + | std.log.info("client detach session={s} fd={d}", .{ self.session_name, client.socket_fd }); | |
| 950 | + | _ = self.closeClient(gpa, client, i, false); | |
| 951 | + | } | |
| 952 | + | ||
| 953 | + | pub fn handleDetachAll(self: *Daemon, gpa: std.mem.Allocator) void { | |
| 954 | + | std.log.info("detach all clients={d}", .{self.clients.items.len}); | |
| 955 | + | for (self.clients.items) |client_to_close| { | |
| 956 | + | client_to_close.deinit(); | |
| 957 | + | gpa.destroy(client_to_close); | |
| 958 | + | } | |
| 959 | + | self.clients.clearRetainingCapacity(); | |
| 960 | + | } | |
| 961 | + | ||
| 962 | + | pub fn handleKill(self: *Daemon, gpa: std.mem.Allocator, io: std.Io) void { | |
| 963 | + | std.log.info("kill received session={s}", .{self.session_name}); | |
| 964 | + | self.shutdown(gpa); | |
| 965 | + | // gracefully shutdown shell processes, shells tend to ignore SIGTERM so we send SIGHUP | |
| 966 | + | // instead | |
| 967 | + | // https://www.gnu.org/software/bash/manual/html_node/Signals.html | |
| 968 | + | // negative pid means kill process and children | |
| 969 | + | std.log.info("sending SIGHUP session={s} pid={d}", .{ self.session_name, self.pid }); | |
| 970 | + | lib_posix.kill(-self.pid, lib_posix.SIG.HUP) catch |err| { | |
| 971 | + | std.log.warn("failed to send SIGHUP to pty child err={s}", .{@errorName(err)}); | |
| 972 | + | }; | |
| 973 | + | std.Io.sleep(io, std.Io.Duration.fromMilliseconds(500), .real) catch unreachable; | |
| 974 | + | lib_posix.kill(-self.pid, lib_posix.SIG.KILL) catch |err| { | |
| 975 | + | std.log.warn("failed to send SIGKILL to pty child err={s}", .{@errorName(err)}); | |
| 976 | + | }; | |
| 977 | + | } | |
| 978 | + | ||
| 979 | + | pub fn handleInfo(self: *Daemon, gpa: std.mem.Allocator, client: *Client) !void { | |
| 980 | + | // zeroes() so asBytes() doesn't ship struct padding + unused cmd/cwd | |
| 981 | + | // tail bytes (daemon stack contents) to clients. | |
| 982 | + | var info = std.mem.zeroes(ipc.Info); | |
| 983 | + | info.clients_len = self.clients.items.len - 1; | |
| 984 | + | info.pid = self.pid; | |
| 985 | + | info.created_at = self.created_at; | |
| 986 | + | info.task_ended_at = self.task_ended_at orelse 0; | |
| 987 | + | info.task_exit_code = self.task_exit_code orelse 0; | |
| 988 | + | ||
| 989 | + | // Build command string from args, re-quoting args that contain | |
| 990 | + | // shell-special characters so the displayed command is copy-pasteable. | |
| 991 | + | const cur_cmd = self.command; | |
| 992 | + | if (cur_cmd) |args| { | |
| 993 | + | for (args, 0..) |arg, i| { | |
| 994 | + | const quoted = if (util.shellNeedsQuoting(arg)) | |
| 995 | + | util.shellQuote(gpa, arg) catch null | |
| 996 | + | else | |
| 997 | + | null; | |
| 998 | + | defer if (quoted) |q| gpa.free(q); | |
| 999 | + | const src = quoted orelse arg; | |
| 1000 | + | ||
| 1001 | + | const need = src.len + @as(usize, if (i > 0) 1 else 0); | |
| 1002 | + | if (info.cmd_len + need > ipc.MAX_CMD_LEN) { | |
| 1003 | + | const ellipsis = "..."; | |
| 1004 | + | if (info.cmd_len + ellipsis.len <= ipc.MAX_CMD_LEN) { | |
| 1005 | + | @memcpy(info.cmd[info.cmd_len..][0..ellipsis.len], ellipsis); | |
| 1006 | + | info.cmd_len += ellipsis.len; | |
| 1007 | + | } | |
| 1008 | + | break; | |
| 1009 | + | } | |
| 1010 | + | ||
| 1011 | + | if (i > 0) { | |
| 1012 | + | info.cmd[info.cmd_len] = ' '; | |
| 1013 | + | info.cmd_len += 1; | |
| 1014 | + | } | |
| 1015 | + | @memcpy(info.cmd[info.cmd_len..][0..src.len], src); | |
| 1016 | + | info.cmd_len += @intCast(src.len); | |
| 1017 | + | } | |
| 1018 | + | } | |
| 1019 | + | ||
| 1020 | + | info.cwd_len = @intCast(@min(self.cwd.len, ipc.MAX_CWD_LEN)); | |
| 1021 | + | @memcpy(info.cwd[0..info.cwd_len], self.cwd[0..info.cwd_len]); | |
| 1022 | + | ||
| 1023 | + | try ipc.appendMessage(gpa, &client.write_buf, .Info, std.mem.asBytes(&info)); | |
| 1024 | + | client.has_pending_output = true; | |
| 1025 | + | } | |
| 1026 | + | ||
| 1027 | + | pub fn handleHistory( | |
| 1028 | + | _: *Daemon, | |
| 1029 | + | gpa: std.mem.Allocator, | |
| 1030 | + | client: *Client, | |
| 1031 | + | term: *ghostty_vt.Terminal, | |
| 1032 | + | payload: []const u8, | |
| 1033 | + | ) !void { | |
| 1034 | + | const format: util.HistoryFormat = if (payload.len > 0) | |
| 1035 | + | @enumFromInt(payload[0]) | |
| 1036 | + | else | |
| 1037 | + | .plain; | |
| 1038 | + | if (util.serializeTerminal(gpa, term, format)) |output| { | |
| 1039 | + | defer gpa.free(output); | |
| 1040 | + | try ipc.appendMessage(gpa, &client.write_buf, .History, output); | |
| 1041 | + | client.has_pending_output = true; | |
| 1042 | + | } else { | |
| 1043 | + | try ipc.appendMessage(gpa, &client.write_buf, .History, ""); | |
| 1044 | + | client.has_pending_output = true; | |
| 1045 | + | } | |
| 1046 | + | } | |
| 1047 | + | ||
| 1048 | + | pub fn handleRun(self: *Daemon, gpa: std.mem.Allocator, client: *Client, payload: []const u8) !void { | |
| 1049 | + | // Reset task tracking so the new command's exit marker is detected. | |
| 1050 | + | // Without this, a second `zmx run` on the same session is ignored | |
| 1051 | + | // because task_exit_code is still set from the first run. | |
| 1052 | + | self.task_exit_code = null; | |
| 1053 | + | self.task_ended_at = null; | |
| 1054 | + | self.is_task_mode = true; | |
| 1055 | + | ||
| 1056 | + | if (payload.len == 0) return; | |
| 1057 | + | ||
| 1058 | + | const cmd = payload; | |
| 1059 | + | ||
| 1060 | + | // Chain the exit marker with `;` on the same line. `$?` captures the | |
| 1061 | + | // exit code of the command (not the `;`). The sole exception is when | |
| 1062 | + | // the command contains a heredoc (`<<`), the delimiter must be alone | |
| 1063 | + | // on its line, so the marker goes on the next line instead. | |
| 1064 | + | const single_line_marker = "; echo ZMX_TASK_COMPLETED:$?\r"; | |
| 1065 | + | const heredoc_marker = "\r\necho ZMX_TASK_COMPLETED:$?\r"; | |
| 1066 | + | const uses_heredoc = std.mem.indexOf(u8, cmd, "<<") != null; | |
| 1067 | + | ||
| 1068 | + | if (cmd.len > 0 and cmd[cmd.len - 1] == '\r') { | |
| 1069 | + | self.queuePtyInput(gpa, cmd[0 .. cmd.len - 1]); | |
| 1070 | + | } else { | |
| 1071 | + | self.queuePtyInput(gpa, cmd); | |
| 1072 | + | } | |
| 1073 | + | self.queuePtyInput(gpa, if (uses_heredoc) heredoc_marker else single_line_marker); | |
| 1074 | + | ||
| 1075 | + | try ipc.appendMessage(gpa, &client.write_buf, .Ack, ""); | |
| 1076 | + | client.has_pending_output = true; | |
| 1077 | + | self.has_had_client = true; | |
| 1078 | + | std.log.debug("run command len={d}", .{payload.len}); | |
| 1079 | + | } | |
| 1080 | + | ||
| 1081 | + | pub fn handleOutput(self: *Daemon, gpa: std.mem.Allocator, payload: []const u8, vt_stream: anytype) !void { | |
| 1082 | + | vt_stream.nextSlice(payload); | |
| 1083 | + | self.has_pty_output = true; | |
| 1084 | + | for (self.clients.items) |client| { | |
| 1085 | + | try ipc.appendMessage(gpa, &client.write_buf, .Output, payload); | |
| 1086 | + | client.has_pending_output = true; | |
| 1087 | + | } | |
| 1088 | + | if (self.clients.items.len > 0) { | |
| 1089 | + | lib_posix.kill(self.pid, lib_posix.SIG.WINCH) catch |err| { | |
| 1090 | + | std.log.warn("failed to send SIGWINCH err={s}", .{@errorName(err)}); | |
| 1091 | + | }; | |
| 1092 | + | } | |
| 1093 | + | } | |
| 1094 | + | ||
| 1095 | + | pub fn handleWrite(self: *Daemon, gpa: std.mem.Allocator, client: *Client, payload: []const u8) !void { | |
| 1096 | + | // Wire format: [u32 path len][path bytes][file content] | |
| 1097 | + | if (payload.len < @sizeOf(u32)) return error.InvalidPayload; | |
| 1098 | + | const path_len = std.mem.bytesToValue(u32, payload[0..@sizeOf(u32)]); | |
| 1099 | + | if (payload.len < @sizeOf(u32) + path_len) return error.InvalidPayload; | |
| 1100 | + | const file_path = payload[@sizeOf(u32)..][0..path_len]; | |
| 1101 | + | const file_content = payload[@sizeOf(u32) + path_len ..]; | |
| 1102 | + | ||
| 1103 | + | // Inject file creation through the PTY so it works over SSH. | |
| 1104 | + | // Base64-encode content and pipe through printf | base64 -d > file. | |
| 1105 | + | // Chunk large files to stay under command-line length limits. | |
| 1106 | + | // 48000 is divisible by 3 (clean base64 boundaries) and encodes | |
| 1107 | + | // to ~64KB, well under typical ARG_MAX. | |
| 1108 | + | const chunk_size = 48000; | |
| 1109 | + | var offset: usize = 0; | |
| 1110 | + | var is_first = true; | |
| 1111 | + | ||
| 1112 | + | while (offset < file_content.len or is_first) { | |
| 1113 | + | const end = @min(offset + chunk_size, file_content.len); | |
| 1114 | + | const chunk = file_content[offset..end]; | |
| 1115 | + | ||
| 1116 | + | const encoded_len = std.base64.standard.Encoder.calcSize(chunk.len); | |
| 1117 | + | const encoded = try gpa.alloc(u8, encoded_len); | |
| 1118 | + | defer gpa.free(encoded); | |
| 1119 | + | _ = std.base64.standard.Encoder.encode(encoded, chunk); | |
| 1120 | + | ||
| 1121 | + | self.queuePtyInput(gpa, "printf '%s' '"); | |
| 1122 | + | self.queuePtyInput(gpa, encoded); | |
| 1123 | + | if (is_first) { | |
| 1124 | + | self.queuePtyInput(gpa, "' | base64 -d > '"); | |
| 1125 | + | } else { | |
| 1126 | + | self.queuePtyInput(gpa, "' | base64 -d >> '"); | |
| 1127 | + | } | |
| 1128 | + | self.queuePtyInput(gpa, file_path); | |
| 1129 | + | self.queuePtyInput(gpa, "'"); | |
| 1130 | + | self.queuePtyInput(gpa, "\r"); | |
| 1131 | + | ||
| 1132 | + | offset = end; | |
| 1133 | + | is_first = false; | |
| 1134 | + | } | |
| 1135 | + | ||
| 1136 | + | try ipc.appendMessage(gpa, &client.write_buf, .Ack, ""); | |
| 1137 | + | client.has_pending_output = true; | |
| 1138 | + | self.has_had_client = true; | |
| 1139 | + | std.log.debug( | |
| 1140 | + | "write command len={d} file_path={s}", | |
| 1141 | + | .{ file_content.len, file_path }, | |
| 1142 | + | ); | |
| 1143 | + | } | |
| 1144 | + | ||
| 1145 | + | fn handleLabelGet(self: *Daemon, gpa: std.mem.Allocator, client: *Client) !void { | |
| 1146 | + | const out = try label.labelsToU8(gpa, self.labels); | |
| 1147 | + | defer gpa.free(out); | |
| 1148 | + | try ipc.appendMessage(gpa, &client.write_buf, .LabelData, out); | |
| 1149 | + | client.has_pending_output = true; | |
| 1150 | + | } | |
| 1151 | + | ||
| 1152 | + | fn handleLabelSet(self: *Daemon, gpa: std.mem.Allocator, client: *Client, labels: []const u8) !void { | |
| 1153 | + | std.log.info("handle label set payload={s}", .{labels}); | |
| 1154 | + | ||
| 1155 | + | var kvs = label.LabelIterator.init(labels); | |
| 1156 | + | while (kvs.next()) |kv| { | |
| 1157 | + | if (kv.value.len == 0) { | |
| 1158 | + | if (self.labels.fetchRemove(kv.key)) |existing| { | |
| 1159 | + | gpa.free(existing.key); | |
| 1160 | + | gpa.free(existing.value); | |
| 1161 | + | } | |
| 1162 | + | continue; | |
| 1163 | + | } | |
| 1164 | + | ||
| 1165 | + | const owned_key = try gpa.dupe(u8, kv.key); | |
| 1166 | + | errdefer gpa.free(owned_key); | |
| 1167 | + | const owned_value = try gpa.dupe(u8, kv.value); | |
| 1168 | + | errdefer gpa.free(owned_value); | |
| 1169 | + | if (try self.labels.fetchPut(gpa, owned_key, owned_value)) |existing| { | |
| 1170 | + | // fetchPut does NOT replace the key in the map, the old | |
| 1171 | + | // key pointer stays. So free the new (unused) key and the | |
| 1172 | + | // old value. | |
| 1173 | + | gpa.free(owned_key); | |
| 1174 | + | gpa.free(existing.value); | |
| 1175 | + | } | |
| 1176 | + | } | |
| 1177 | + | ||
| 1178 | + | try ipc.appendMessage(gpa, &client.write_buf, .Ack, ""); | |
| 1179 | + | client.has_pending_output = true; | |
| 1180 | + | } | |
| 1181 | + | ||
| 1182 | + | fn handleLabelClear(self: *Daemon, gpa: std.mem.Allocator, client: *Client) !void { | |
| 1183 | + | var it = self.labels.iterator(); | |
| 1184 | + | while (it.next()) |entry| { | |
| 1185 | + | gpa.free(entry.key_ptr.*); | |
| 1186 | + | gpa.free(entry.value_ptr.*); | |
| 1187 | + | } | |
| 1188 | + | self.labels.clearRetainingCapacity(); | |
| 1189 | + | try ipc.appendMessage(gpa, &client.write_buf, .Ack, ""); | |
| 1190 | + | client.has_pending_output = true; | |
| 1191 | + | } | |
| 1192 | + | }; | |
| 1193 | + | ||
| 1194 | + | test "send queues PTY input without changing leader" { | |
| 1195 | + | const alloc = std.testing.allocator; | |
| 1196 | + | var daemon = Daemon{ | |
| 1197 | + | .cfg = undefined, | |
| 1198 | + | .clients = .empty, | |
| 1199 | + | .leader_client_fd = 42, | |
| 1200 | + | .session_name = "test", | |
| 1201 | + | .socket_path = "", | |
| 1202 | + | .running = true, | |
| 1203 | + | .pid = 0, | |
| 1204 | + | .created_at = 0, | |
| 1205 | + | }; | |
| 1206 | + | defer daemon.pty_write_buf.deinit(alloc); | |
| 1207 | + | ||
| 1208 | + | daemon.handleSend(alloc, "hello"); | |
| 1209 | + | ||
| 1210 | + | try std.testing.expectEqual(@as(?i32, 42), daemon.leader_client_fd); | |
| 1211 | + | try std.testing.expectEqualStrings("hello", daemon.pty_write_buf.items); | |
| 1212 | + | } |
+139,
-1711
| ... | ... | @@ -9,88 +9,20 @@ const cross = @import("cross.zig"); | |
| 9 | 9 | const socket = @import("socket.zig"); | |
| 10 | 10 | const label = @import("label.zig"); | |
| 11 | 11 | const lib_posix = @import("posix.zig"); | |
| 12 | - | ||
| 13 | - | pub const version = build_options.version; | |
| 14 | - | pub const ghostty_version = build_options.ghostty_version; | |
| 15 | - | ||
| 16 | - | var log_system = log.LogSystem{}; | |
| 12 | + | const signal = @import("signal.zig"); | |
| 13 | + | const Cfg = @import("cfg.zig"); | |
| 14 | + | const loop = @import("loop.zig"); | |
| 15 | + | const Client = loop.Client; | |
| 16 | + | const Daemon = loop.Daemon; | |
| 17 | + | const version = build_options.version; | |
| 18 | + | const ghostty_version = build_options.ghostty_version; | |
| 17 | 19 | ||
| 18 | 20 | pub const std_options: std.Options = .{ | |
| 19 | - | .logFn = zmxLogFn, | |
| 21 | + | .logFn = log.zmxLogFn, | |
| 20 | 22 | .log_level = .debug, | |
| 21 | 23 | }; | |
| 22 | 24 | ||
| 23 | - | fn zmxLogFn( | |
| 24 | - | comptime level: std.log.Level, | |
| 25 | - | comptime scope: anytype, | |
| 26 | - | comptime format: []const u8, | |
| 27 | - | args: anytype, | |
| 28 | - | ) void { | |
| 29 | - | log_system.log(level, scope, format, args) catch {}; | |
| 30 | - | } | |
| 31 | - | ||
| 32 | - | /// Self-pipe woken by signal handlers. std.posix.poll loops on .INTR internally | |
| 33 | - | /// (PollError has no Interrupted member), so a signal that lands during poll() | |
| 34 | - | /// never surfaces; the handler writes a byte here and poll() wakes on POLLIN. | |
| 35 | - | var sig_pipe: [2]lib_posix.fd_t = .{ -1, -1 }; | |
| 36 | - | ||
| 37 | - | // https://github.com/ziglang/zig/blob/738d2be9d6b6ef3ff3559130c05159ef53336224/lib/std/posix.zig#L3505 | |
| 38 | - | const O_NONBLOCK: usize = 1 << @bitOffsetOf(lib_posix.O, "NONBLOCK"); | |
| 39 | - | ||
| 40 | - | const SessionMatch = struct { | |
| 41 | - | name: []const u8, | |
| 42 | - | is_prefix: bool, | |
| 43 | - | ||
| 44 | - | fn matches(self: SessionMatch, session_name: []const u8) bool { | |
| 45 | - | if (self.is_prefix) return std.mem.startsWith(u8, session_name, self.name); | |
| 46 | - | return std.mem.eql(u8, session_name, self.name); | |
| 47 | - | } | |
| 48 | - | }; | |
| 49 | - | ||
| 50 | - | fn resolveSessionOrEnv(alloc: std.mem.Allocator, io: std.Io, session_name: ?[]const u8) ![]const u8 { | |
| 51 | - | const sesh_env = socket.getSeshNameFromEnv(); | |
| 52 | - | const raw = if (session_name) |name| | |
| 53 | - | if (std.mem.eql(u8, name, ".")) blk: { | |
| 54 | - | if (sesh_env.len > 0) break :blk sesh_env; | |
| 55 | - | var buf: [4096]u8 = undefined; | |
| 56 | - | var w = std.Io.File.stderr().writer(io, &buf); | |
| 57 | - | w.interface.print("error: \".\" requires ZMX_SESSION (are you inside a zmx session?)\n", .{}) catch {}; | |
| 58 | - | w.interface.flush() catch {}; | |
| 59 | - | return error.SessionNameRequired; | |
| 60 | - | } else name | |
| 61 | - | else if (sesh_env.len > 0) | |
| 62 | - | sesh_env | |
| 63 | - | else { | |
| 64 | - | return error.SessionNameRequired; | |
| 65 | - | }; | |
| 66 | - | return socket.getSeshName(alloc, raw); | |
| 67 | - | } | |
| 68 | - | ||
| 69 | - | fn parseSessionArg(alloc: std.mem.Allocator, raw: []const u8) !SessionMatch { | |
| 70 | - | if (raw.len > 0 and raw[raw.len - 1] == '*') { | |
| 71 | - | const name = try socket.getSeshName(alloc, raw[0 .. raw.len - 1]); | |
| 72 | - | return .{ .name = name, .is_prefix = true }; | |
| 73 | - | } | |
| 74 | - | const name = try socket.getSeshName(alloc, raw); | |
| 75 | - | return .{ .name = name, .is_prefix = false }; | |
| 76 | - | } | |
| 77 | - | ||
| 78 | - | fn openSignalPipe() !void { | |
| 79 | - | sig_pipe = try lib_posix.pipe2(.{ .CLOEXEC = true, .NONBLOCK = true }); | |
| 80 | - | } | |
| 81 | - | ||
| 82 | - | fn drainSignalPipe() void { | |
| 83 | - | var b: [16]u8 = undefined; | |
| 84 | - | while (true) { | |
| 85 | - | const n = lib_posix.read(sig_pipe[0], &b) catch return; | |
| 86 | - | if (n == 0) return; | |
| 87 | - | } | |
| 88 | - | } | |
| 89 | - | ||
| 90 | - | fn detectHelp(arg: []const u8) bool { | |
| 91 | - | return (std.mem.eql(u8, arg, "--help") or std.mem.eql(u8, arg, "-h")); | |
| 92 | - | } | |
| 93 | - | ||
| 25 | + | /// This is the entry point for the CLI. | |
| 94 | 26 | pub fn main(init: std.process.Init) !void { | |
| 95 | 27 | const gpa = init.gpa; | |
| 96 | 28 | const io = init.io; |
| ... | ... | @@ -99,7 +31,7 @@ pub fn main(init: std.process.Init) !void { | |
| 99 | 31 | // disappears between probe and send would otherwise kill us before | |
| 100 | 32 | // write() can return BrokenPipe. Inherited across fork, so this also | |
| 101 | 33 | // covers the daemon. | |
| 102 | - | ignoreSigpipe(); | |
| 34 | + | signal.ignoreSigpipe(); | |
| 103 | 35 | ||
| 104 | 36 | var args = init.minimal.args.iterate(); | |
| 105 | 37 | defer args.deinit(); |
| ... | ... | @@ -111,8 +43,8 @@ pub fn main(init: std.process.Init) !void { | |
| 111 | 43 | const log_path = try std.fs.path.join(gpa, &.{ cfg.log_dir, "zmx.log" }); | |
| 112 | 44 | defer gpa.free(log_path); | |
| 113 | 45 | const log_mode = std.Io.File.Permissions.fromMode(cfg.log_mode); | |
| 114 | - | try log_system.init(gpa, io, log_path, log_mode); | |
| 115 | - | defer log_system.deinit(); | |
| 46 | + | try log.log_system.init(io, log_path, log_mode); | |
| 47 | + | defer log.log_system.deinit(); | |
| 116 | 48 | ||
| 117 | 49 | const shell_env = init.environ_map.get("SHELL") orelse "/bin/sh"; | |
| 118 | 50 |
| ... | ... | @@ -134,14 +66,14 @@ pub fn main(init: std.process.Init) !void { | |
| 134 | 66 | } else if (std.mem.eql(u8, cmd, "get") or std.mem.eql(u8, cmd, "g")) { | |
| 135 | 67 | const sesh_name = args.next() orelse return error.SessionNameRequired; | |
| 136 | 68 | if (detectHelp(sesh_name)) return help(io); | |
| 137 | - | const sesh = try resolveSessionOrEnv(gpa, io, sesh_name); | |
| 69 | + | const sesh = try socket.resolveSessionOrEnv(gpa, io, sesh_name); | |
| 138 | 70 | defer gpa.free(sesh); | |
| 139 | 71 | const single_kv = args.next() orelse ""; | |
| 140 | 72 | return labelGet(gpa, io, &cfg, sesh, single_kv); | |
| 141 | 73 | } else if (std.mem.eql(u8, cmd, "set")) { | |
| 142 | 74 | const sesh_name = args.next() orelse return error.SessionNameRequired; | |
| 143 | 75 | if (detectHelp(sesh_name)) return help(io); | |
| 144 | - | const sesh = try resolveSessionOrEnv(gpa, io, sesh_name); | |
| 76 | + | const sesh = try socket.resolveSessionOrEnv(gpa, io, sesh_name); | |
| 145 | 77 | defer gpa.free(sesh); | |
| 146 | 78 | ||
| 147 | 79 | var kvs = std.ArrayList(u8).empty; |
| ... | ... | @@ -156,7 +88,7 @@ pub fn main(init: std.process.Init) !void { | |
| 156 | 88 | } else if (std.mem.eql(u8, cmd, "clear")) { | |
| 157 | 89 | const sesh_name = args.next() orelse return error.SessionNameRequired; | |
| 158 | 90 | if (detectHelp(sesh_name)) return help(io); | |
| 159 | - | const sesh = try resolveSessionOrEnv(gpa, io, sesh_name); | |
| 91 | + | const sesh = try socket.resolveSessionOrEnv(gpa, io, sesh_name); | |
| 160 | 92 | defer gpa.free(sesh); | |
| 161 | 93 | return labelClear(gpa, io, &cfg, sesh); | |
| 162 | 94 | } else if (std.mem.eql(u8, cmd, "completions") or std.mem.eql(u8, cmd, "c")) { |
| ... | ... | @@ -198,7 +130,6 @@ pub fn main(init: std.process.Init) !void { | |
| 198 | 130 | try command_args.append(gpa, arg); | |
| 199 | 131 | } | |
| 200 | 132 | ||
| 201 | - | const clients = try std.ArrayList(*Client).initCapacity(gpa, 10); | |
| 202 | 133 | var command: ?[][]const u8 = null; | |
| 203 | 134 | if (command_args.items.len > 0) { | |
| 204 | 135 | command = command_args.items; |
| ... | ... | @@ -210,27 +141,16 @@ pub fn main(init: std.process.Init) !void { | |
| 210 | 141 | ||
| 211 | 142 | const sesh = try socket.getSeshName(gpa, session_name); | |
| 212 | 143 | defer gpa.free(sesh); | |
| 213 | - | var daemon = Daemon{ | |
| 214 | - | .io = io, | |
| 215 | - | .running = true, | |
| 216 | - | .cfg = &cfg, | |
| 217 | - | .alloc = std.heap.c_allocator, | |
| 218 | - | .clients = clients, | |
| 219 | - | .session_name = sesh, | |
| 220 | - | .socket_path = undefined, | |
| 221 | - | .pid = undefined, | |
| 222 | - | .command = command, | |
| 223 | - | .cwd = cwd, | |
| 224 | - | .created_at = @intCast(std.Io.Timestamp.now(io, .real).nanoseconds), | |
| 225 | - | .leader_client_fd = null, | |
| 226 | - | .shell = shell_env, | |
| 227 | - | }; | |
| 228 | - | daemon.socket_path = socket.getSocketPath(gpa, cfg.socket_dir, sesh) catch |err| switch (err) { | |
| 229 | - | error.NameTooLong => return socket.printSessionNameTooLong(daemon.io, sesh, cfg.socket_dir), | |
| 144 | + | const socket_path = socket.getSocketPath(gpa, cfg.socket_dir, sesh) catch |err| switch (err) { | |
| 145 | + | error.NameTooLong => return socket.printSessionNameTooLong(io, sesh, cfg.socket_dir), | |
| 230 | 146 | error.OutOfMemory => return err, | |
| 231 | 147 | }; | |
| 148 | + | var daemon = Daemon.init(io, &cfg, sesh, socket_path); | |
| 149 | + | daemon.command = command; | |
| 150 | + | daemon.cwd = cwd; | |
| 151 | + | daemon.shell = shell_env; | |
| 232 | 152 | std.log.info("socket path={s}", .{daemon.socket_path}); | |
| 233 | - | return attach(&daemon); | |
| 153 | + | return attach(gpa, io, &daemon); | |
| 234 | 154 | } else if (std.mem.eql(u8, cmd, "run") or std.mem.eql(u8, cmd, "r")) { | |
| 235 | 155 | const session_name = args.next() orelse ""; | |
| 236 | 156 | if (std.mem.eql(u8, session_name, "--help") or std.mem.eql(u8, session_name, "-h")) { |
| ... | ... | @@ -247,7 +167,6 @@ pub fn main(init: std.process.Init) !void { | |
| 247 | 167 | try cmd_args_raw.append(gpa, arg); | |
| 248 | 168 | } | |
| 249 | 169 | } | |
| 250 | - | const clients = try std.ArrayList(*Client).initCapacity(gpa, 10); | |
| 251 | 170 | ||
| 252 | 171 | var cwd_buf: [std.fs.max_path_bytes]u8 = undefined; | |
| 253 | 172 | const cwd_len = std.process.currentPath(io, &cwd_buf) catch 0; |
| ... | ... | @@ -255,28 +174,17 @@ pub fn main(init: std.process.Init) !void { | |
| 255 | 174 | ||
| 256 | 175 | const sesh = try socket.getSeshName(gpa, session_name); | |
| 257 | 176 | defer gpa.free(sesh); | |
| 258 | - | var daemon = Daemon{ | |
| 259 | - | .io = io, | |
| 260 | - | .running = true, | |
| 261 | - | .cfg = &cfg, | |
| 262 | - | .alloc = std.heap.c_allocator, | |
| 263 | - | .clients = clients, | |
| 264 | - | .session_name = sesh, | |
| 265 | - | .socket_path = undefined, | |
| 266 | - | .pid = undefined, | |
| 267 | - | .command = null, | |
| 268 | - | .cwd = cwd, | |
| 269 | - | .created_at = @intCast(std.Io.Timestamp.now(io, .real).nanoseconds), | |
| 270 | - | .is_task_mode = true, | |
| 271 | - | .leader_client_fd = null, | |
| 272 | - | .shell = shell_env, | |
| 273 | - | }; | |
| 274 | - | daemon.socket_path = socket.getSocketPath(gpa, cfg.socket_dir, sesh) catch |err| switch (err) { | |
| 275 | - | error.NameTooLong => return socket.printSessionNameTooLong(daemon.io, sesh, cfg.socket_dir), | |
| 177 | + | const socket_path = socket.getSocketPath(gpa, cfg.socket_dir, sesh) catch |err| switch (err) { | |
| 178 | + | error.NameTooLong => return socket.printSessionNameTooLong(io, sesh, cfg.socket_dir), | |
| 276 | 179 | error.OutOfMemory => return err, | |
| 277 | 180 | }; | |
| 181 | + | defer gpa.free(socket_path); | |
| 182 | + | var daemon = Daemon.init(io, &cfg, sesh, socket_path); | |
| 183 | + | daemon.cwd = cwd; | |
| 184 | + | daemon.is_task_mode = true; | |
| 185 | + | daemon.shell = shell_env; | |
| 278 | 186 | std.log.info("socket path={s}", .{daemon.socket_path}); | |
| 279 | - | return run(&daemon, detached, cmd_args_raw.items); | |
| 187 | + | return run(gpa, io, &daemon, detached, cmd_args_raw.items); | |
| 280 | 188 | } else if (std.mem.eql(u8, cmd, "send") or std.mem.eql(u8, cmd, "s")) { | |
| 281 | 189 | const session_name = args.next() orelse ""; | |
| 282 | 190 | if (std.mem.eql(u8, session_name, "--help") or std.mem.eql(u8, session_name, "-h")) { |
| ... | ... | @@ -322,7 +230,7 @@ pub fn main(init: std.process.Init) !void { | |
| 322 | 230 | var stderr_writer = std.Io.File.stderr().writer(io, &stderr_buffer); | |
| 323 | 231 | const stderr = &stderr_writer.interface; | |
| 324 | 232 | ||
| 325 | - | var matchers: std.ArrayList(SessionMatch) = .empty; | |
| 233 | + | var matchers: std.ArrayList(socket.SessionMatch) = .empty; | |
| 326 | 234 | defer { | |
| 327 | 235 | for (matchers.items) |m| { | |
| 328 | 236 | gpa.free(m.name); |
| ... | ... | @@ -338,7 +246,7 @@ pub fn main(init: std.process.Init) !void { | |
| 338 | 246 | force = true; | |
| 339 | 247 | continue; | |
| 340 | 248 | } | |
| 341 | - | const m = try parseSessionArg(gpa, session_name); | |
| 249 | + | const m = try socket.parseSessionArg(gpa, session_name); | |
| 342 | 250 | try matchers.append(gpa, m); | |
| 343 | 251 | } | |
| 344 | 252 | if (matchers.items.len == 0) { |
| ... | ... | @@ -369,7 +277,7 @@ pub fn main(init: std.process.Init) !void { | |
| 369 | 277 | } | |
| 370 | 278 | } | |
| 371 | 279 | } else if (std.mem.eql(u8, cmd, "wait") or std.mem.eql(u8, cmd, "w")) { | |
| 372 | - | var matchers: std.ArrayList(SessionMatch) = .empty; | |
| 280 | + | var matchers: std.ArrayList(socket.SessionMatch) = .empty; | |
| 373 | 281 | defer { | |
| 374 | 282 | for (matchers.items) |m| { | |
| 375 | 283 | gpa.free(m.name); |
| ... | ... | @@ -380,7 +288,7 @@ pub fn main(init: std.process.Init) !void { | |
| 380 | 288 | if (std.mem.eql(u8, session_name, "--help") or std.mem.eql(u8, session_name, "-h")) { | |
| 381 | 289 | return help(io); | |
| 382 | 290 | } | |
| 383 | - | const m = try parseSessionArg(gpa, session_name); | |
| 291 | + | const m = try socket.parseSessionArg(gpa, session_name); | |
| 384 | 292 | try matchers.append(gpa, m); | |
| 385 | 293 | } | |
| 386 | 294 | if (matchers.items.len == 0) { |
| ... | ... | @@ -388,7 +296,7 @@ pub fn main(init: std.process.Init) !void { | |
| 388 | 296 | } | |
| 389 | 297 | return wait(gpa, io, &cfg, matchers); | |
| 390 | 298 | } else if (std.mem.eql(u8, cmd, "tail") or std.mem.eql(u8, cmd, "t")) { | |
| 391 | - | var matchers: std.ArrayList(SessionMatch) = .empty; | |
| 299 | + | var matchers: std.ArrayList(socket.SessionMatch) = .empty; | |
| 392 | 300 | defer { | |
| 393 | 301 | for (matchers.items) |m| { | |
| 394 | 302 | gpa.free(m.name); |
| ... | ... | @@ -399,7 +307,7 @@ pub fn main(init: std.process.Init) !void { | |
| 399 | 307 | if (std.mem.eql(u8, session_name, "--help") or std.mem.eql(u8, session_name, "-h")) { | |
| 400 | 308 | return help(io); | |
| 401 | 309 | } | |
| 402 | - | const m = try parseSessionArg(gpa, session_name); | |
| 310 | + | const m = try socket.parseSessionArg(gpa, session_name); | |
| 403 | 311 | try matchers.append(gpa, m); | |
| 404 | 312 | } | |
| 405 | 313 | if (matchers.items.len == 0) { |
| ... | ... | @@ -479,962 +387,23 @@ pub fn main(init: std.process.Init) !void { | |
| 479 | 387 | var cwd_buf: [std.fs.max_path_bytes]u8 = undefined; | |
| 480 | 388 | const cwd_len = std.process.currentPath(io, &cwd_buf) catch 0; | |
| 481 | 389 | const cwd = cwd_buf[0..cwd_len]; | |
| 482 | - | const clients = try std.ArrayList(*Client).initCapacity(gpa, 10); | |
| 483 | 390 | const sesh = try socket.getSeshName(gpa, session_name); | |
| 484 | 391 | defer gpa.free(sesh); | |
| 485 | - | var daemon = Daemon{ | |
| 486 | - | .io = io, | |
| 487 | - | .running = true, | |
| 488 | - | .cfg = &cfg, | |
| 489 | - | .alloc = std.heap.c_allocator, | |
| 490 | - | .clients = clients, | |
| 491 | - | .session_name = sesh, | |
| 492 | - | .socket_path = undefined, | |
| 493 | - | .pid = undefined, | |
| 494 | - | .command = null, | |
| 495 | - | .cwd = cwd, | |
| 496 | - | .created_at = @intCast(std.Io.Timestamp.now(io, .real).nanoseconds), | |
| 497 | - | .is_task_mode = true, | |
| 498 | - | .leader_client_fd = null, | |
| 499 | - | .shell = shell_env, | |
| 500 | - | }; | |
| 501 | - | daemon.socket_path = socket.getSocketPath(gpa, cfg.socket_dir, sesh) catch |err| switch (err) { | |
| 502 | - | error.NameTooLong => return socket.printSessionNameTooLong(daemon.io, sesh, cfg.socket_dir), | |
| 392 | + | const socket_path = socket.getSocketPath(gpa, cfg.socket_dir, sesh) catch |err| switch (err) { | |
| 393 | + | error.NameTooLong => return socket.printSessionNameTooLong(io, sesh, cfg.socket_dir), | |
| 503 | 394 | error.OutOfMemory => return err, | |
| 504 | 395 | }; | |
| 396 | + | var daemon = Daemon.init(io, &cfg, sesh, socket_path); | |
| 397 | + | daemon.is_task_mode = true; | |
| 398 | + | daemon.cwd = cwd; | |
| 399 | + | daemon.shell = shell_env; | |
| 505 | 400 | std.log.info("socket path={s}", .{daemon.socket_path}); | |
| 506 | - | try writeFile(&daemon, file_path); | |
| 401 | + | try writeFile(gpa, io, &daemon, file_path); | |
| 507 | 402 | } else { | |
| 508 | 403 | return help(io); | |
| 509 | 404 | } | |
| 510 | 405 | } | |
| 511 | 406 | ||
| 512 | - | /// Client represents each terminal that has connected to a session. | |
| 513 | - | /// | |
| 514 | - | /// Multiple Clients can connect to a single session. | |
| 515 | - | const Client = struct { | |
| 516 | - | alloc: std.mem.Allocator, | |
| 517 | - | socket_fd: i32, | |
| 518 | - | has_pending_output: bool = false, | |
| 519 | - | read_buf: ipc.SocketBuffer, | |
| 520 | - | write_buf: std.ArrayList(u8), | |
| 521 | - | ||
| 522 | - | pub fn deinit(self: *Client) void { | |
| 523 | - | lib_posix.close(self.socket_fd); | |
| 524 | - | self.read_buf.deinit(); | |
| 525 | - | self.write_buf.deinit(self.alloc); | |
| 526 | - | } | |
| 527 | - | }; | |
| 528 | - | ||
| 529 | - | /// Cfg is zmx's configuration container. | |
| 530 | - | /// | |
| 531 | - | /// The purpose of this container is to hold anything that can be modified by the user. | |
| 532 | - | const Cfg = struct { | |
| 533 | - | socket_dir: []const u8, | |
| 534 | - | log_dir: []const u8, | |
| 535 | - | max_scrollback: usize = 10_000_000, | |
| 536 | - | dir_mode: u32 = 0o750, | |
| 537 | - | log_mode: u32 = 0o640, | |
| 538 | - | ||
| 539 | - | pub fn init(alloc: std.mem.Allocator, io: std.Io) !Cfg { | |
| 540 | - | const socket_dir = try socketDir(alloc); | |
| 541 | - | errdefer alloc.free(socket_dir); | |
| 542 | - | const log_dir = try logDir(alloc); | |
| 543 | - | errdefer alloc.free(log_dir); | |
| 544 | - | ||
| 545 | - | const dir_mode = if (lib_posix.getenv("ZMX_DIR_MODE")) |m| | |
| 546 | - | std.fmt.parseInt(u32, m, 8) catch 0o750 | |
| 547 | - | else | |
| 548 | - | 0o750; | |
| 549 | - | ||
| 550 | - | const log_mode = if (lib_posix.getenv("ZMX_LOG_MODE")) |m| | |
| 551 | - | std.fmt.parseInt(u32, m, 8) catch 0o640 | |
| 552 | - | else | |
| 553 | - | 0o640; | |
| 554 | - | ||
| 555 | - | var cfg = Cfg{ | |
| 556 | - | .socket_dir = socket_dir, | |
| 557 | - | .log_dir = log_dir, | |
| 558 | - | .dir_mode = dir_mode, | |
| 559 | - | .log_mode = log_mode, | |
| 560 | - | }; | |
| 561 | - | ||
| 562 | - | try cfg.mkdir(io); | |
| 563 | - | ||
| 564 | - | return cfg; | |
| 565 | - | } | |
| 566 | - | ||
| 567 | - | fn socketDir(alloc: std.mem.Allocator) ![]const u8 { | |
| 568 | - | const tmpdir = std.mem.trimEnd(u8, lib_posix.getenv("TMPDIR") orelse "/tmp", "/"); | |
| 569 | - | const uid = lib_posix.getuid(); | |
| 570 | - | ||
| 571 | - | const socket_dir: []const u8 = if (lib_posix.getenv("ZMX_DIR")) |zmxdir| | |
| 572 | - | try alloc.dupe(u8, zmxdir) | |
| 573 | - | else if (lib_posix.getenv("XDG_RUNTIME_DIR")) |xdg_runtime| | |
| 574 | - | try std.fmt.allocPrint(alloc, "{s}/zmx", .{xdg_runtime}) | |
| 575 | - | else | |
| 576 | - | try std.fmt.allocPrint(alloc, "{s}/zmx-{d}", .{ tmpdir, uid }); | |
| 577 | - | ||
| 578 | - | return socket_dir; | |
| 579 | - | } | |
| 580 | - | ||
| 581 | - | fn logDir(alloc: std.mem.Allocator) ![]const u8 { | |
| 582 | - | const log_dir = if (lib_posix.getenv("ZMX_DIR")) |zmxdir| | |
| 583 | - | try std.fmt.allocPrint(alloc, "{s}/logs", .{zmxdir}) | |
| 584 | - | else if (lib_posix.getenv("XDG_STATE_HOME")) |xdg_state_home| | |
| 585 | - | try std.fmt.allocPrint(alloc, "{s}/zmx/logs", .{xdg_state_home}) | |
| 586 | - | else if (lib_posix.getenv("HOME")) |home_dir| | |
| 587 | - | try std.fmt.allocPrint(alloc, "{s}/.local/state/zmx/logs", .{home_dir}) | |
| 588 | - | else fallback: { | |
| 589 | - | // This is the last resort: falling back to /tmp/$UID if HOME is unset. | |
| 590 | - | const tmpdir = std.mem.trimEnd(u8, lib_posix.getenv("TMPDIR") orelse "/tmp", "/"); | |
| 591 | - | const uid = lib_posix.getuid(); | |
| 592 | - | break :fallback try std.fmt.allocPrint(alloc, "{s}/zmx-{d}", .{ tmpdir, uid }); | |
| 593 | - | }; | |
| 594 | - | ||
| 595 | - | return log_dir; | |
| 596 | - | } | |
| 597 | - | ||
| 598 | - | pub fn deinit(self: *Cfg, alloc: std.mem.Allocator) void { | |
| 599 | - | if (self.socket_dir.len > 0) alloc.free(self.socket_dir); | |
| 600 | - | if (self.log_dir.len > 0) alloc.free(self.log_dir); | |
| 601 | - | } | |
| 602 | - | ||
| 603 | - | pub fn mkdir(self: *Cfg, io: std.Io) !void { | |
| 604 | - | const sock_perms = std.Io.Dir.Permissions.fromMode(@intCast(self.dir_mode)); | |
| 605 | - | try mkdirAll(io, self.socket_dir, sock_perms); | |
| 606 | - | const log_perms = std.Io.Dir.Permissions.fromMode(@intCast(self.dir_mode)); | |
| 607 | - | try mkdirAll(io, self.log_dir, log_perms); | |
| 608 | - | } | |
| 609 | - | ||
| 610 | - | fn mkdirAll(io: std.Io, sub_dir_path: []const u8, permissions: std.Io.Dir.Permissions) !void { | |
| 611 | - | var it = std.fs.path.componentIterator(sub_dir_path); | |
| 612 | - | var component = it.last() orelse return error.BadPathName; | |
| 613 | - | while (true) { | |
| 614 | - | std.Io.Dir.createDirAbsolute(io, component.path, permissions) catch |err| switch (err) { | |
| 615 | - | error.PathAlreadyExists => {}, | |
| 616 | - | error.FileNotFound => |e| { | |
| 617 | - | component = it.previous() orelse return e; | |
| 618 | - | continue; | |
| 619 | - | }, | |
| 620 | - | else => |e| return e, | |
| 621 | - | }; | |
| 622 | - | component = it.next() orelse return; | |
| 623 | - | } | |
| 624 | - | } | |
| 625 | - | }; | |
| 626 | - | ||
| 627 | - | test "Cfg.init uses default modes when env vars are not set" { | |
| 628 | - | const alloc = std.testing.allocator; | |
| 629 | - | ||
| 630 | - | // Ensure they are not set | |
| 631 | - | _ = cross.c.unsetenv("ZMX_DIR_MODE"); | |
| 632 | - | _ = cross.c.unsetenv("ZMX_LOG_MODE"); | |
| 633 | - | ||
| 634 | - | var cfg = try Cfg.init(alloc, std.testing.io); | |
| 635 | - | defer cfg.deinit(alloc); | |
| 636 | - | ||
| 637 | - | try std.testing.expectEqual(@as(u32, 0o750), cfg.dir_mode); | |
| 638 | - | try std.testing.expectEqual(@as(u32, 0o640), cfg.log_mode); | |
| 639 | - | } | |
| 640 | - | ||
| 641 | - | test "Cfg.init uses custom modes from env vars" { | |
| 642 | - | const alloc = std.testing.allocator; | |
| 643 | - | ||
| 644 | - | // Set custom octal values | |
| 645 | - | _ = cross.c.setenv("ZMX_DIR_MODE", "770", 1); | |
| 646 | - | _ = cross.c.setenv("ZMX_LOG_MODE", "660", 1); | |
| 647 | - | defer { | |
| 648 | - | _ = cross.c.unsetenv("ZMX_DIR_MODE"); | |
| 649 | - | _ = cross.c.unsetenv("ZMX_LOG_MODE"); | |
| 650 | - | } | |
| 651 | - | ||
| 652 | - | var cfg = try Cfg.init(alloc, std.testing.io); | |
| 653 | - | defer cfg.deinit(alloc); | |
| 654 | - | ||
| 655 | - | try std.testing.expectEqual(@as(u32, 0o770), cfg.dir_mode); | |
| 656 | - | try std.testing.expectEqual(@as(u32, 0o660), cfg.log_mode); | |
| 657 | - | } | |
| 658 | - | ||
| 659 | - | /// Daemon is responsible for managing a zmx session. | |
| 660 | - | /// | |
| 661 | - | /// It holds all the state for a running session. Instead of a single daemon for all sessions, we | |
| 662 | - | /// create a daemon for every session. This has some benefits. The ipc communication between | |
| 663 | - | /// session clients and the daemon doesn't need to be tagged with the session name. If a daemon | |
| 664 | - | /// crashes for one session won't crash all the other sessions. | |
| 665 | - | /// | |
| 666 | - | /// Conceptually it's also much simpler to reason about. | |
| 667 | - | const Daemon = struct { | |
| 668 | - | io: std.Io, | |
| 669 | - | cfg: *Cfg, | |
| 670 | - | alloc: std.mem.Allocator, | |
| 671 | - | clients: std.ArrayList(*Client), | |
| 672 | - | labels: std.StringHashMapUnmanaged([]u8) = .empty, | |
| 673 | - | // This control which client is the leader. The leader controls terminal state and | |
| 674 | - | // cols/rows of session. | |
| 675 | - | leader_client_fd: ?i32, | |
| 676 | - | session_name: []const u8, | |
| 677 | - | socket_path: []const u8, | |
| 678 | - | running: bool, | |
| 679 | - | pid: i32, | |
| 680 | - | command: ?[]const []const u8 = null, | |
| 681 | - | cwd: []const u8 = "", | |
| 682 | - | has_pty_output: bool = false, | |
| 683 | - | has_had_client: bool = false, | |
| 684 | - | has_terminal_client: bool = false, // true only after a real attach (.Init received) | |
| 685 | - | created_at: u64, // unix timestamp (ns) | |
| 686 | - | is_task_mode: bool = false, // flag for when session is run as a task | |
| 687 | - | task_exit_code: ?u8 = null, // null = running or n/a, set when task completes | |
| 688 | - | task_ended_at: ?u64 = null, // timestamp when task exited | |
| 689 | - | pty_fd: i32 = -1, // set by daemonLoop so handleRun can probe the foreground process | |
| 690 | - | pty_write_buf: std.ArrayList(u8) = .empty, | |
| 691 | - | shell: []const u8 = "/bin/sh", | |
| 692 | - | ||
| 693 | - | const EnsureSessionResult = struct { | |
| 694 | - | created: bool, | |
| 695 | - | is_daemon: bool, | |
| 696 | - | }; | |
| 697 | - | ||
| 698 | - | pub fn deinit(self: *Daemon) void { | |
| 699 | - | self.clients.deinit(self.alloc); | |
| 700 | - | var it = self.labels.iterator(); | |
| 701 | - | while (it.next()) |entry| { | |
| 702 | - | self.alloc.free(entry.key_ptr.*); | |
| 703 | - | self.alloc.free(entry.value_ptr.*); | |
| 704 | - | } | |
| 705 | - | self.labels.deinit(self.alloc); | |
| 706 | - | self.pty_write_buf.deinit(self.alloc); | |
| 707 | - | self.alloc.free(self.socket_path); | |
| 708 | - | } | |
| 709 | - | ||
| 710 | - | fn handleLabelGet(self: *Daemon, client: *Client) !void { | |
| 711 | - | const out = try label.labelsToU8(self.alloc, self.labels); | |
| 712 | - | defer self.alloc.free(out); | |
| 713 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .LabelData, out); | |
| 714 | - | client.has_pending_output = true; | |
| 715 | - | } | |
| 716 | - | ||
| 717 | - | fn handleLabelSet(self: *Daemon, client: *Client, labels: []const u8) !void { | |
| 718 | - | std.log.info("handle label set payload={s}", .{labels}); | |
| 719 | - | ||
| 720 | - | var kvs = label.LabelIterator.init(labels); | |
| 721 | - | while (kvs.next()) |kv| { | |
| 722 | - | if (kv.value.len == 0) { | |
| 723 | - | if (self.labels.fetchRemove(kv.key)) |existing| { | |
| 724 | - | self.alloc.free(existing.key); | |
| 725 | - | self.alloc.free(existing.value); | |
| 726 | - | } | |
| 727 | - | continue; | |
| 728 | - | } | |
| 729 | - | ||
| 730 | - | const owned_key = try self.alloc.dupe(u8, kv.key); | |
| 731 | - | errdefer self.alloc.free(owned_key); | |
| 732 | - | const owned_value = try self.alloc.dupe(u8, kv.value); | |
| 733 | - | errdefer self.alloc.free(owned_value); | |
| 734 | - | if (try self.labels.fetchPut(self.alloc, owned_key, owned_value)) |existing| { | |
| 735 | - | // fetchPut does NOT replace the key in the map, the old | |
| 736 | - | // key pointer stays. So free the new (unused) key and the | |
| 737 | - | // old value. | |
| 738 | - | self.alloc.free(owned_key); | |
| 739 | - | self.alloc.free(existing.value); | |
| 740 | - | } | |
| 741 | - | } | |
| 742 | - | ||
| 743 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .Ack, ""); | |
| 744 | - | client.has_pending_output = true; | |
| 745 | - | } | |
| 746 | - | ||
| 747 | - | fn handleLabelClear(self: *Daemon, client: *Client) !void { | |
| 748 | - | var it = self.labels.iterator(); | |
| 749 | - | while (it.next()) |entry| { | |
| 750 | - | self.alloc.free(entry.key_ptr.*); | |
| 751 | - | self.alloc.free(entry.value_ptr.*); | |
| 752 | - | } | |
| 753 | - | self.labels.clearRetainingCapacity(); | |
| 754 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .Ack, ""); | |
| 755 | - | client.has_pending_output = true; | |
| 756 | - | } | |
| 757 | - | ||
| 758 | - | pub fn shutdown(self: *Daemon) void { | |
| 759 | - | std.log.info("shutting down daemon session={s}", .{self.session_name}); | |
| 760 | - | self.running = false; | |
| 761 | - | ||
| 762 | - | for (self.clients.items) |client| { | |
| 763 | - | client.deinit(); | |
| 764 | - | self.alloc.destroy(client); | |
| 765 | - | } | |
| 766 | - | self.clients.clearRetainingCapacity(); | |
| 767 | - | } | |
| 768 | - | ||
| 769 | - | pub fn closeClient(self: *Daemon, client: *Client, i: usize, shutdown_on_last: bool) bool { | |
| 770 | - | const fd = client.socket_fd; | |
| 771 | - | // leader is disconnected, remove ref and let another client claim leader on input | |
| 772 | - | if (self.leader_client_fd == client.socket_fd) { | |
| 773 | - | std.log.info( | |
| 774 | - | "unsetting leader session={s} fd={d}", | |
| 775 | - | .{ self.session_name, client.socket_fd }, | |
| 776 | - | ); | |
| 777 | - | self.leader_client_fd = null; | |
| 778 | - | } | |
| 779 | - | client.deinit(); | |
| 780 | - | self.alloc.destroy(client); | |
| 781 | - | _ = self.clients.orderedRemove(i); | |
| 782 | - | std.log.info("client disconnected fd={d} remaining={d}", .{ fd, self.clients.items.len }); | |
| 783 | - | if (shutdown_on_last and self.clients.items.len == 0) { | |
| 784 | - | self.shutdown(); | |
| 785 | - | return true; | |
| 786 | - | } | |
| 787 | - | return false; | |
| 788 | - | } | |
| 789 | - | ||
| 790 | - | fn setLeader(self: *Daemon, client: *Client) !void { | |
| 791 | - | std.log.info("setting new leader client_fd={d}", .{client.socket_fd}); | |
| 792 | - | self.leader_client_fd = client.socket_fd; | |
| 793 | - | // Send a resize message to the client so it can send us back their window size | |
| 794 | - | // so we can resize the pty and ghostty state. | |
| 795 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .Resize, ""); | |
| 796 | - | client.has_pending_output = true; | |
| 797 | - | } | |
| 798 | - | ||
| 799 | - | /// Runs in the forked child. Either execs or returns an error (caller | |
| 800 | - | /// must exit on error -- returning would fall through to parent code). | |
| 801 | - | fn execChild(self: *Daemon) !noreturn { | |
| 802 | - | const alloc = std.heap.c_allocator; | |
| 803 | - | ||
| 804 | - | // main() set SIGPIPE to SIG_IGN, which (unlike handlers) survives | |
| 805 | - | // exec. Restore the default so the shell and its children behave | |
| 806 | - | // normally (e.g. `yes | head` should exit 141 via SIGPIPE). | |
| 807 | - | const dfl: lib_posix.Sigaction = .{ | |
| 808 | - | .handler = .{ .handler = lib_posix.SIG.DFL }, | |
| 809 | - | .mask = lib_posix.sigemptyset(), | |
| 810 | - | .flags = 0, | |
| 811 | - | }; | |
| 812 | - | lib_posix.sigaction(lib_posix.SIG.PIPE, &dfl, null); | |
| 813 | - | ||
| 814 | - | const session_env = try std.fmt.allocPrintSentinel( | |
| 815 | - | alloc, | |
| 816 | - | "ZMX_SESSION={s}", | |
| 817 | - | .{self.session_name}, | |
| 818 | - | 0, | |
| 819 | - | ); | |
| 820 | - | _ = cross.c.putenv(session_env.ptr); | |
| 821 | - | ||
| 822 | - | if (self.command) |cmd_args| { | |
| 823 | - | const argv = try alloc.allocSentinel(?[*:0]const u8, cmd_args.len, null); | |
| 824 | - | for (cmd_args, 0..) |arg, i| { | |
| 825 | - | argv[i] = try alloc.dupeZ(u8, arg); | |
| 826 | - | } | |
| 827 | - | const err = lib_posix.execvpeZ(argv[0].?, argv.ptr, std.c.environ); | |
| 828 | - | std.log.err("execvpe failed: cmd={s} err={s}", .{ cmd_args[0], @errorName(err) }); | |
| 829 | - | lib_posix.exit(1); | |
| 830 | - | } | |
| 831 | - | ||
| 832 | - | var buf: [256]u8 = undefined; | |
| 833 | - | const z = try std.fmt.bufPrintZ(&buf, "{s}", .{self.shell}); | |
| 834 | - | const shell: [:0]const u8 = if (self.is_task_mode) "bash" else z; | |
| 835 | - | // Use "-shellname" as argv[0] to signal login shell (traditional method) | |
| 836 | - | const login_shell = try std.fmt.allocPrintSentinel( | |
| 837 | - | alloc, | |
| 838 | - | "-{s}", | |
| 839 | - | .{std.fs.path.basename(shell)}, | |
| 840 | - | 0, | |
| 841 | - | ); | |
| 842 | - | const argv = [_:null]?[*:0]const u8{ login_shell, null }; | |
| 843 | - | const err = lib_posix.execvpeZ(shell, &argv, std.c.environ); | |
| 844 | - | std.log.err("execvpe failed: shell={s} err={s}", .{ shell, @errorName(err) }); | |
| 845 | - | lib_posix.exit(1); | |
| 846 | - | } | |
| 847 | - | ||
| 848 | - | /// spawnPty runs forkpty() and executes the shell or shell command the user provides. | |
| 849 | - | fn spawnPty(self: *Daemon) !c_int { | |
| 850 | - | const size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO); | |
| 851 | - | var ws: cross.c.struct_winsize = .{ | |
| 852 | - | .ws_row = size.rows, | |
| 853 | - | .ws_col = size.cols, | |
| 854 | - | .ws_xpixel = size.xpixel, | |
| 855 | - | .ws_ypixel = size.ypixel, | |
| 856 | - | }; | |
| 857 | - | ||
| 858 | - | var master_fd: c_int = undefined; | |
| 859 | - | const pid = cross.forkpty(&master_fd, null, null, &ws); | |
| 860 | - | if (pid < 0) { | |
| 861 | - | return error.ForkPtyFailed; | |
| 862 | - | } | |
| 863 | - | ||
| 864 | - | if (pid == 0) { // child pid code path | |
| 865 | - | // In the forked child, ANY error must exit rather than propagate: | |
| 866 | - | // a returned error falls through to the parent code path below, | |
| 867 | - | // running a second daemon on the same socket (or worse, hitting | |
| 868 | - | // errdefers that delete the parent's socket file). | |
| 869 | - | execChild(self) catch |err| { | |
| 870 | - | std.log.err("child setup failed: {s}", .{@errorName(err)}); | |
| 871 | - | lib_posix.exit(1); | |
| 872 | - | }; | |
| 873 | - | unreachable; // execChild either execs or exits, never returns ok | |
| 874 | - | } | |
| 875 | - | // master pid code path | |
| 876 | - | self.pid = pid; | |
| 877 | - | std.log.info("pty spawned session={s} pid={d}", .{ self.session_name, pid }); | |
| 878 | - | ||
| 879 | - | // make pty non-blocking | |
| 880 | - | const flags = try lib_posix.fcntl(master_fd, lib_posix.F.GETFL, 0); | |
| 881 | - | _ = try lib_posix.fcntl(master_fd, lib_posix.F.SETFL, flags | O_NONBLOCK); | |
| 882 | - | return master_fd; | |
| 883 | - | } | |
| 884 | - | ||
| 885 | - | /// ensureSession "upserts" a session by checking if the unix socket exists already. | |
| 886 | - | /// If not it creates one and spawns the daemon. | |
| 887 | - | fn ensureSession(self: *Daemon) !EnsureSessionResult { | |
| 888 | - | std.log.info("ensure session session={s}", .{self.session_name}); | |
| 889 | - | var dir = try std.Io.Dir.openDirAbsolute(self.io, self.cfg.socket_dir, .{}); | |
| 890 | - | defer dir.close(self.io); | |
| 891 | - | ||
| 892 | - | const exists = try socket.sessionExists(self.io, dir, self.session_name); | |
| 893 | - | var should_create = !exists; | |
| 894 | - | ||
| 895 | - | if (exists) { | |
| 896 | - | if (ipc.connectSession(self.socket_path)) |fd| { | |
| 897 | - | lib_posix.close(fd); | |
| 898 | - | if (self.command != null) { | |
| 899 | - | std.log.warn( | |
| 900 | - | "session already exists, ignoring command session={s}", | |
| 901 | - | .{self.session_name}, | |
| 902 | - | ); | |
| 903 | - | } | |
| 904 | - | } else |err| switch (err) { | |
| 905 | - | // Daemon is definitively gone: safe to replace. | |
| 906 | - | error.ConnectionRefused => { | |
| 907 | - | socket.cleanupStaleSocket(self.io, dir, self.session_name); | |
| 908 | - | should_create = true; | |
| 909 | - | }, | |
| 910 | - | // Connect failed for an unusual reason. The check is only to | |
| 911 | - | // decide create-vs-attach; the socket file exists, so proceed | |
| 912 | - | // to attach rather than fail or orphan. | |
| 913 | - | else => { | |
| 914 | - | std.log.warn( | |
| 915 | - | "connect failed ({s}), proceeding to attach session={s}", | |
| 916 | - | .{ @errorName(err), self.session_name }, | |
| 917 | - | ); | |
| 918 | - | }, | |
| 919 | - | } | |
| 920 | - | } | |
| 921 | - | ||
| 922 | - | if (should_create) { | |
| 923 | - | std.log.info("creating session={s}", .{self.session_name}); | |
| 924 | - | const server_sock_fd = try socket.createSocket(self.socket_path); | |
| 925 | - | ||
| 926 | - | // creates the daemon | |
| 927 | - | const pid = try lib_posix.fork(); | |
| 928 | - | if (pid == 0) { // child (daemon) | |
| 929 | - | // becomes the session leader and detaches process from its controlling terminal | |
| 930 | - | _ = try lib_posix.setsid(); | |
| 931 | - | ||
| 932 | - | log_system.deinit(); | |
| 933 | - | ||
| 934 | - | // Redirect stdin/stdout/stderr to /dev/null. The daemon | |
| 935 | - | // communicates via its unix socket, not stdio. Without | |
| 936 | - | // this, any pipe on FDs 0-2 (e.g. from bats' `run` | |
| 937 | - | // keyword) stays open for the daemon's lifetime, causing | |
| 938 | - | // the caller to hang waiting for EOF. | |
| 939 | - | { | |
| 940 | - | const devnull = lib_posix.open( | |
| 941 | - | "/dev/null", | |
| 942 | - | .{ .ACCMODE = .RDWR }, | |
| 943 | - | 0, | |
| 944 | - | ) catch |err| { | |
| 945 | - | std.log.warn("failed to open /dev/null: {s}", .{@errorName(err)}); | |
| 946 | - | return err; | |
| 947 | - | }; | |
| 948 | - | inline for (.{ lib_posix.STDIN_FILENO, lib_posix.STDOUT_FILENO, lib_posix.STDERR_FILENO }) |fd| { | |
| 949 | - | _ = lib_posix.dup2(devnull, fd) catch |err| { | |
| 950 | - | std.log.warn("dup2 /dev/null -> {d}: {s}", .{ fd, @errorName(err) }); | |
| 951 | - | return err; | |
| 952 | - | }; | |
| 953 | - | } | |
| 954 | - | if (devnull > 2) lib_posix.close(devnull); | |
| 955 | - | } | |
| 956 | - | ||
| 957 | - | // Close file descriptors inherited from the parent that the | |
| 958 | - | // daemon doesn't need. This prevents test harnesses (like | |
| 959 | - | // bats) from hanging -- they wait for their internal FDs (3+) | |
| 960 | - | // to close before exiting. | |
| 961 | - | // | |
| 962 | - | // Must run BEFORE log_system.init() otherwise the new log | |
| 963 | - | // FD gets closed, and spawnPty() reuses that FD number for | |
| 964 | - | // the PTY master, causing log writes to leak into the terminal. | |
| 965 | - | // | |
| 966 | - | // Skip server_sock_fd (needed for IPC) and dir.fd (needed to | |
| 967 | - | // delete the socket file on shutdown). | |
| 968 | - | { | |
| 969 | - | const dir_fd = @as(i32, @intCast(dir.handle)); | |
| 970 | - | var fd: i32 = 3; | |
| 971 | - | while (fd < 64) : (fd += 1) { | |
| 972 | - | if (fd == server_sock_fd or fd == dir_fd) continue; | |
| 973 | - | _ = std.c.close(fd); | |
| 974 | - | } | |
| 975 | - | } | |
| 976 | - | ||
| 977 | - | const session_log_name = try std.fmt.allocPrint( | |
| 978 | - | self.alloc, | |
| 979 | - | "{s}.log", | |
| 980 | - | .{self.session_name}, | |
| 981 | - | ); | |
| 982 | - | defer self.alloc.free(session_log_name); | |
| 983 | - | const session_log_path = try std.fs.path.join( | |
| 984 | - | self.alloc, | |
| 985 | - | &.{ self.cfg.log_dir, session_log_name }, | |
| 986 | - | ); | |
| 987 | - | defer self.alloc.free(session_log_path); | |
| 988 | - | const log_mode = std.Io.File.Permissions.fromMode(self.cfg.log_mode); | |
| 989 | - | try log_system.init(self.alloc, self.io, session_log_path, log_mode); | |
| 990 | - | ||
| 991 | - | // If spawnPty fails, clean up here. Once it succeeds, | |
| 992 | - | // the inner block's defer takes ownership of cleanup to | |
| 993 | - | // avoid double-closing server_sock_fd on daemonLoop error. | |
| 994 | - | const pty_fd = self.spawnPty() catch |err| { | |
| 995 | - | lib_posix.close(server_sock_fd); | |
| 996 | - | dir.deleteFile(self.io, self.session_name) catch {}; | |
| 997 | - | return err; | |
| 998 | - | }; | |
| 999 | - | ||
| 1000 | - | defer { | |
| 1001 | - | // Close and unlink the listen socket BEFORE handleKill()'s | |
| 1002 | - | // 500ms SIGHUP->SIGKILL grace sleep. Otherwise a `zmx run` | |
| 1003 | - | // for the same name issued in that window will hang waiting | |
| 1004 | - | // for a connect. | |
| 1005 | - | lib_posix.close(server_sock_fd); | |
| 1006 | - | std.log.info("deleting socket file session={s}", .{self.session_name}); | |
| 1007 | - | dir.deleteFile(self.io, self.session_name) catch |err| { | |
| 1008 | - | std.log.warn("failed to delete socket file err={s}", .{@errorName(err)}); | |
| 1009 | - | }; | |
| 1010 | - | self.handleKill(); | |
| 1011 | - | self.deinit(); | |
| 1012 | - | lib_posix.close(pty_fd); | |
| 1013 | - | _ = lib_posix.waitpid(self.pid, 0); | |
| 1014 | - | } | |
| 1015 | - | ||
| 1016 | - | try daemonLoop(self, server_sock_fd, pty_fd); | |
| 1017 | - | std.log.info("daemon loop shutdown", .{}); | |
| 1018 | - | return .{ .created = true, .is_daemon = true }; | |
| 1019 | - | } | |
| 1020 | - | lib_posix.close(server_sock_fd); | |
| 1021 | - | std.Io.sleep(self.io, std.Io.Duration.fromMilliseconds(10), .real) catch unreachable; | |
| 1022 | - | return .{ .created = true, .is_daemon = false }; | |
| 1023 | - | } | |
| 1024 | - | ||
| 1025 | - | return .{ .created = false, .is_daemon = false }; | |
| 1026 | - | } | |
| 1027 | - | ||
| 1028 | - | const PTY_WRITE_BUF_MAX = 256 * 1024; | |
| 1029 | - | ||
| 1030 | - | /// Queue bytes for the PTY's stdin. Flushed by daemonLoop on POLLOUT. | |
| 1031 | - | /// Drops the payload if the buffer is over cap -- same failure mode as | |
| 1032 | - | /// the old direct-write ptyWrite (drop on EAGAIN), just at a 64x higher | |
| 1033 | - | /// threshold. Capping avoids OOM when the shell stops reading; dropping | |
| 1034 | - | /// new (not old) bytes avoids tearing a partially-accepted sequence. | |
| 1035 | - | fn queuePtyInput(self: *Daemon, data: []const u8) void { | |
| 1036 | - | if (data.len == 0) return; | |
| 1037 | - | if (self.pty_write_buf.items.len + data.len > PTY_WRITE_BUF_MAX) { | |
| 1038 | - | std.log.warn( | |
| 1039 | - | "pty input dropped {d} bytes (buffer full, shell not reading)", | |
| 1040 | - | .{data.len}, | |
| 1041 | - | ); | |
| 1042 | - | return; | |
| 1043 | - | } | |
| 1044 | - | std.log.debug("buffering pty input data={x}", .{data}); | |
| 1045 | - | self.pty_write_buf.appendSlice(self.alloc, data) catch |err| { | |
| 1046 | - | std.log.warn( | |
| 1047 | - | "pty input dropped {d} bytes: {s}", | |
| 1048 | - | .{ data.len, @errorName(err) }, | |
| 1049 | - | ); | |
| 1050 | - | }; | |
| 1051 | - | } | |
| 1052 | - | ||
| 1053 | - | pub fn handleInput(self: *Daemon, client: *Client, payload: []const u8) !void { | |
| 1054 | - | std.log.debug("buffering pty input data={x}", .{payload}); | |
| 1055 | - | // client is leader, send entire payload (ansi escape codes + text) | |
| 1056 | - | if (self.leader_client_fd == client.socket_fd) { | |
| 1057 | - | self.queuePtyInput(payload); | |
| 1058 | - | return; | |
| 1059 | - | } | |
| 1060 | - | ||
| 1061 | - | // check if leader needs to be updated by detecting any user input | |
| 1062 | - | if (util.isUserInput(payload)) { | |
| 1063 | - | try self.setLeader(client); | |
| 1064 | - | self.queuePtyInput(payload); | |
| 1065 | - | } | |
| 1066 | - | } | |
| 1067 | - | ||
| 1068 | - | /// Queue input from `zmx send` without changing interactive client leadership. | |
| 1069 | - | pub fn handleSend(self: *Daemon, payload: []const u8) void { | |
| 1070 | - | self.queuePtyInput(payload); | |
| 1071 | - | } | |
| 1072 | - | ||
| 1073 | - | pub fn handleSwitch(self: *Daemon, session_name: []const u8) !void { | |
| 1074 | - | for (self.clients.items) |client| { | |
| 1075 | - | if (self.leader_client_fd == client.socket_fd) { | |
| 1076 | - | ipc.appendMessage( | |
| 1077 | - | self.alloc, | |
| 1078 | - | &client.write_buf, | |
| 1079 | - | .Switch, | |
| 1080 | - | session_name, | |
| 1081 | - | ) catch |err| { | |
| 1082 | - | std.log.warn( | |
| 1083 | - | "failed to buffer terminal state for client err={s}", | |
| 1084 | - | .{@errorName(err)}, | |
| 1085 | - | ); | |
| 1086 | - | }; | |
| 1087 | - | client.has_pending_output = true; | |
| 1088 | - | return; | |
| 1089 | - | } | |
| 1090 | - | } | |
| 1091 | - | return error.NoLeaderFound; | |
| 1092 | - | } | |
| 1093 | - | ||
| 1094 | - | pub fn handleInit( | |
| 1095 | - | self: *Daemon, | |
| 1096 | - | client: *Client, | |
| 1097 | - | pty_fd: i32, | |
| 1098 | - | term: *ghostty_vt.Terminal, | |
| 1099 | - | payload: []const u8, | |
| 1100 | - | ) !void { | |
| 1101 | - | if (payload.len != @sizeOf(ipc.Resize)) return; | |
| 1102 | - | ||
| 1103 | - | // Serialize terminal state BEFORE resize to capture correct cursor position. | |
| 1104 | - | // Resizing triggers reflow which can move the cursor, and the shell's | |
| 1105 | - | // SIGWINCH-triggered redraw will run after our snapshot is sent. | |
| 1106 | - | // Only serialize on re-attach (has_had_client), not first attach, to avoid | |
| 1107 | - | // interfering with shell initialization (DA1 queries, etc.) | |
| 1108 | - | if (self.has_pty_output and self.has_had_client) { | |
| 1109 | - | const cursor = &term.screens.active.cursor; | |
| 1110 | - | std.log.debug( | |
| 1111 | - | "cursor before serialize: x={d} y={d} pending_wrap={}", | |
| 1112 | - | .{ cursor.x, cursor.y, cursor.pending_wrap }, | |
| 1113 | - | ); | |
| 1114 | - | if (util.serializeTerminalState(self.alloc, term)) |term_output| { | |
| 1115 | - | std.log.debug("serialize terminal state", .{}); | |
| 1116 | - | // Rewrite OSC 133;A to include redraw=0 so the outer terminal | |
| 1117 | - | // does not clear prompt lines on resize (issue #111). | |
| 1118 | - | const restore_data = util.rewritePromptRedraw(self.alloc, term_output) orelse term_output; | |
| 1119 | - | defer self.alloc.free(term_output); | |
| 1120 | - | defer if (restore_data.ptr != term_output.ptr) self.alloc.free(restore_data); | |
| 1121 | - | ipc.appendMessage(self.alloc, &client.write_buf, .Output, restore_data) catch |err| { | |
| 1122 | - | std.log.warn( | |
| 1123 | - | "failed to buffer terminal state for client err={s}", | |
| 1124 | - | .{@errorName(err)}, | |
| 1125 | - | ); | |
| 1126 | - | }; | |
| 1127 | - | client.has_pending_output = true; | |
| 1128 | - | } | |
| 1129 | - | } | |
| 1130 | - | ||
| 1131 | - | // no leader is set so set one | |
| 1132 | - | if (self.leader_client_fd == null) { | |
| 1133 | - | try self.setLeader(client); | |
| 1134 | - | } | |
| 1135 | - | ||
| 1136 | - | // only resize if leader | |
| 1137 | - | if (self.leader_client_fd == client.socket_fd) { | |
| 1138 | - | const resize = std.mem.bytesToValue(ipc.Resize, payload); | |
| 1139 | - | var ws: cross.c.struct_winsize = .{ | |
| 1140 | - | .ws_row = resize.rows, | |
| 1141 | - | .ws_col = resize.cols, | |
| 1142 | - | .ws_xpixel = resize.xpixel, | |
| 1143 | - | .ws_ypixel = resize.ypixel, | |
| 1144 | - | }; | |
| 1145 | - | _ = cross.c.ioctl(pty_fd, cross.c.TIOCSWINSZ, &ws); | |
| 1146 | - | // Disable prompt_redraw before resize. The daemon's internal terminal | |
| 1147 | - | // would otherwise clear prompt lines expecting the shell to redraw them, | |
| 1148 | - | // but the shell's redraw goes to the PTY (forwarded to clients), not to | |
| 1149 | - | // this daemon terminal. The clearing corrupts the daemon's snapshot state. | |
| 1150 | - | const saved_prompt_redraw = term.flags.shell_redraws_prompt; | |
| 1151 | - | term.flags.shell_redraws_prompt = .false; | |
| 1152 | - | defer term.flags.shell_redraws_prompt = saved_prompt_redraw; | |
| 1153 | - | const opts = ghostty_vt.Terminal.Resize{ | |
| 1154 | - | .cols = resize.cols, | |
| 1155 | - | .rows = resize.rows, | |
| 1156 | - | }; | |
| 1157 | - | try term.resize(self.alloc, opts); | |
| 1158 | - | ||
| 1159 | - | // Mark that we've had a client init, so subsequent clients get terminal state | |
| 1160 | - | self.has_had_client = true; | |
| 1161 | - | self.has_terminal_client = true; | |
| 1162 | - | ||
| 1163 | - | std.log.debug("init resize rows={d} cols={d}", .{ resize.rows, resize.cols }); | |
| 1164 | - | } | |
| 1165 | - | } | |
| 1166 | - | ||
| 1167 | - | pub fn handleResize( | |
| 1168 | - | self: *Daemon, | |
| 1169 | - | client: *Client, | |
| 1170 | - | pty_fd: i32, | |
| 1171 | - | term: *ghostty_vt.Terminal, | |
| 1172 | - | payload: []const u8, | |
| 1173 | - | ) !void { | |
| 1174 | - | if (payload.len != @sizeOf(ipc.Resize)) return; | |
| 1175 | - | if (self.leader_client_fd == null) { | |
| 1176 | - | try self.setLeader(client); | |
| 1177 | - | } | |
| 1178 | - | // only leader can resize | |
| 1179 | - | if (self.leader_client_fd != client.socket_fd) return; | |
| 1180 | - | ||
| 1181 | - | const resize = std.mem.bytesToValue(ipc.Resize, payload); | |
| 1182 | - | var ws: cross.c.struct_winsize = .{ | |
| 1183 | - | .ws_row = resize.rows, | |
| 1184 | - | .ws_col = resize.cols, | |
| 1185 | - | .ws_xpixel = resize.xpixel, | |
| 1186 | - | .ws_ypixel = resize.ypixel, | |
| 1187 | - | }; | |
| 1188 | - | _ = cross.c.ioctl(pty_fd, cross.c.TIOCSWINSZ, &ws); | |
| 1189 | - | // Disable prompt_redraw before resize (same rationale as handleInit). | |
| 1190 | - | const saved_prompt_redraw = term.flags.shell_redraws_prompt; | |
| 1191 | - | term.flags.shell_redraws_prompt = .false; | |
| 1192 | - | defer term.flags.shell_redraws_prompt = saved_prompt_redraw; | |
| 1193 | - | const opts = ghostty_vt.Terminal.Resize{ | |
| 1194 | - | .cols = resize.cols, | |
| 1195 | - | .rows = resize.rows, | |
| 1196 | - | }; | |
| 1197 | - | try term.resize(self.alloc, opts); | |
| 1198 | - | std.log.debug("resize rows={d} cols={d}", .{ resize.rows, resize.cols }); | |
| 1199 | - | } | |
| 1200 | - | ||
| 1201 | - | pub fn handleDetach(self: *Daemon, client: *Client, i: usize) void { | |
| 1202 | - | std.log.info("client detach session={s} fd={d}", .{ self.session_name, client.socket_fd }); | |
| 1203 | - | _ = self.closeClient(client, i, false); | |
| 1204 | - | } | |
| 1205 | - | ||
| 1206 | - | pub fn handleDetachAll(self: *Daemon) void { | |
| 1207 | - | std.log.info("detach all clients={d}", .{self.clients.items.len}); | |
| 1208 | - | for (self.clients.items) |client_to_close| { | |
| 1209 | - | client_to_close.deinit(); | |
| 1210 | - | self.alloc.destroy(client_to_close); | |
| 1211 | - | } | |
| 1212 | - | self.clients.clearRetainingCapacity(); | |
| 1213 | - | } | |
| 1214 | - | ||
| 1215 | - | pub fn handleKill(self: *Daemon) void { | |
| 1216 | - | std.log.info("kill received session={s}", .{self.session_name}); | |
| 1217 | - | self.shutdown(); | |
| 1218 | - | // gracefully shutdown shell processes, shells tend to ignore SIGTERM so we send SIGHUP | |
| 1219 | - | // instead | |
| 1220 | - | // https://www.gnu.org/software/bash/manual/html_node/Signals.html | |
| 1221 | - | // negative pid means kill process and children | |
| 1222 | - | std.log.info("sending SIGHUP session={s} pid={d}", .{ self.session_name, self.pid }); | |
| 1223 | - | lib_posix.kill(-self.pid, lib_posix.SIG.HUP) catch |err| { | |
| 1224 | - | std.log.warn("failed to send SIGHUP to pty child err={s}", .{@errorName(err)}); | |
| 1225 | - | }; | |
| 1226 | - | std.Io.sleep(self.io, std.Io.Duration.fromMilliseconds(500), .real) catch unreachable; | |
| 1227 | - | lib_posix.kill(-self.pid, lib_posix.SIG.KILL) catch |err| { | |
| 1228 | - | std.log.warn("failed to send SIGKILL to pty child err={s}", .{@errorName(err)}); | |
| 1229 | - | }; | |
| 1230 | - | } | |
| 1231 | - | ||
| 1232 | - | pub fn handleInfo(self: *Daemon, client: *Client) !void { | |
| 1233 | - | // zeroes() so asBytes() doesn't ship struct padding + unused cmd/cwd | |
| 1234 | - | // tail bytes (daemon stack contents) to clients. | |
| 1235 | - | var info = std.mem.zeroes(ipc.Info); | |
| 1236 | - | info.clients_len = self.clients.items.len - 1; | |
| 1237 | - | info.pid = self.pid; | |
| 1238 | - | info.created_at = self.created_at; | |
| 1239 | - | info.task_ended_at = self.task_ended_at orelse 0; | |
| 1240 | - | info.task_exit_code = self.task_exit_code orelse 0; | |
| 1241 | - | ||
| 1242 | - | // Build command string from args, re-quoting args that contain | |
| 1243 | - | // shell-special characters so the displayed command is copy-pasteable. | |
| 1244 | - | const cur_cmd = self.command; | |
| 1245 | - | if (cur_cmd) |args| { | |
| 1246 | - | for (args, 0..) |arg, i| { | |
| 1247 | - | const quoted = if (util.shellNeedsQuoting(arg)) | |
| 1248 | - | util.shellQuote(self.alloc, arg) catch null | |
| 1249 | - | else | |
| 1250 | - | null; | |
| 1251 | - | defer if (quoted) |q| self.alloc.free(q); | |
| 1252 | - | const src = quoted orelse arg; | |
| 1253 | - | ||
| 1254 | - | const need = src.len + @as(usize, if (i > 0) 1 else 0); | |
| 1255 | - | if (info.cmd_len + need > ipc.MAX_CMD_LEN) { | |
| 1256 | - | const ellipsis = "..."; | |
| 1257 | - | if (info.cmd_len + ellipsis.len <= ipc.MAX_CMD_LEN) { | |
| 1258 | - | @memcpy(info.cmd[info.cmd_len..][0..ellipsis.len], ellipsis); | |
| 1259 | - | info.cmd_len += ellipsis.len; | |
| 1260 | - | } | |
| 1261 | - | break; | |
| 1262 | - | } | |
| 1263 | - | ||
| 1264 | - | if (i > 0) { | |
| 1265 | - | info.cmd[info.cmd_len] = ' '; | |
| 1266 | - | info.cmd_len += 1; | |
| 1267 | - | } | |
| 1268 | - | @memcpy(info.cmd[info.cmd_len..][0..src.len], src); | |
| 1269 | - | info.cmd_len += @intCast(src.len); | |
| 1270 | - | } | |
| 1271 | - | } | |
| 1272 | - | ||
| 1273 | - | info.cwd_len = @intCast(@min(self.cwd.len, ipc.MAX_CWD_LEN)); | |
| 1274 | - | @memcpy(info.cwd[0..info.cwd_len], self.cwd[0..info.cwd_len]); | |
| 1275 | - | ||
| 1276 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .Info, std.mem.asBytes(&info)); | |
| 1277 | - | client.has_pending_output = true; | |
| 1278 | - | } | |
| 1279 | - | ||
| 1280 | - | pub fn handleHistory( | |
| 1281 | - | self: *Daemon, | |
| 1282 | - | client: *Client, | |
| 1283 | - | term: *ghostty_vt.Terminal, | |
| 1284 | - | payload: []const u8, | |
| 1285 | - | ) !void { | |
| 1286 | - | const format: util.HistoryFormat = if (payload.len > 0) | |
| 1287 | - | @enumFromInt(payload[0]) | |
| 1288 | - | else | |
| 1289 | - | .plain; | |
| 1290 | - | if (util.serializeTerminal(self.alloc, term, format)) |output| { | |
| 1291 | - | defer self.alloc.free(output); | |
| 1292 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .History, output); | |
| 1293 | - | client.has_pending_output = true; | |
| 1294 | - | } else { | |
| 1295 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .History, ""); | |
| 1296 | - | client.has_pending_output = true; | |
| 1297 | - | } | |
| 1298 | - | } | |
| 1299 | - | ||
| 1300 | - | pub fn handleRun(self: *Daemon, client: *Client, payload: []const u8) !void { | |
| 1301 | - | // Reset task tracking so the new command's exit marker is detected. | |
| 1302 | - | // Without this, a second `zmx run` on the same session is ignored | |
| 1303 | - | // because task_exit_code is still set from the first run. | |
| 1304 | - | self.task_exit_code = null; | |
| 1305 | - | self.task_ended_at = null; | |
| 1306 | - | self.is_task_mode = true; | |
| 1307 | - | ||
| 1308 | - | if (payload.len == 0) return; | |
| 1309 | - | ||
| 1310 | - | const cmd = payload; | |
| 1311 | - | ||
| 1312 | - | // Chain the exit marker with `;` on the same line. `$?` captures the | |
| 1313 | - | // exit code of the command (not the `;`). The sole exception is when | |
| 1314 | - | // the command contains a heredoc (`<<`), the delimiter must be alone | |
| 1315 | - | // on its line, so the marker goes on the next line instead. | |
| 1316 | - | const single_line_marker = "; echo ZMX_TASK_COMPLETED:$?\r"; | |
| 1317 | - | const heredoc_marker = "\r\necho ZMX_TASK_COMPLETED:$?\r"; | |
| 1318 | - | const uses_heredoc = std.mem.indexOf(u8, cmd, "<<") != null; | |
| 1319 | - | ||
| 1320 | - | if (cmd.len > 0 and cmd[cmd.len - 1] == '\r') { | |
| 1321 | - | self.queuePtyInput(cmd[0 .. cmd.len - 1]); | |
| 1322 | - | } else { | |
| 1323 | - | self.queuePtyInput(cmd); | |
| 1324 | - | } | |
| 1325 | - | self.queuePtyInput(if (uses_heredoc) heredoc_marker else single_line_marker); | |
| 1326 | - | ||
| 1327 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .Ack, ""); | |
| 1328 | - | client.has_pending_output = true; | |
| 1329 | - | self.has_had_client = true; | |
| 1330 | - | std.log.debug("run command len={d}", .{payload.len}); | |
| 1331 | - | } | |
| 1332 | - | ||
| 1333 | - | pub fn handleOutput(self: *Daemon, payload: []const u8, vt_stream: anytype) !void { | |
| 1334 | - | vt_stream.nextSlice(payload); | |
| 1335 | - | self.has_pty_output = true; | |
| 1336 | - | for (self.clients.items) |client| { | |
| 1337 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .Output, payload); | |
| 1338 | - | client.has_pending_output = true; | |
| 1339 | - | } | |
| 1340 | - | if (self.clients.items.len > 0) { | |
| 1341 | - | lib_posix.kill(self.pid, lib_posix.SIG.WINCH) catch |err| { | |
| 1342 | - | std.log.warn("failed to send SIGWINCH err={s}", .{@errorName(err)}); | |
| 1343 | - | }; | |
| 1344 | - | } | |
| 1345 | - | } | |
| 1346 | - | ||
| 1347 | - | pub fn handleWrite(self: *Daemon, client: *Client, payload: []const u8) !void { | |
| 1348 | - | // Wire format: [u32 path len][path bytes][file content] | |
| 1349 | - | if (payload.len < @sizeOf(u32)) return error.InvalidPayload; | |
| 1350 | - | const path_len = std.mem.bytesToValue(u32, payload[0..@sizeOf(u32)]); | |
| 1351 | - | if (payload.len < @sizeOf(u32) + path_len) return error.InvalidPayload; | |
| 1352 | - | const file_path = payload[@sizeOf(u32)..][0..path_len]; | |
| 1353 | - | const file_content = payload[@sizeOf(u32) + path_len ..]; | |
| 1354 | - | ||
| 1355 | - | // Inject file creation through the PTY so it works over SSH. | |
| 1356 | - | // Base64-encode content and pipe through printf | base64 -d > file. | |
| 1357 | - | // Chunk large files to stay under command-line length limits. | |
| 1358 | - | // 48000 is divisible by 3 (clean base64 boundaries) and encodes | |
| 1359 | - | // to ~64KB, well under typical ARG_MAX. | |
| 1360 | - | const chunk_size = 48000; | |
| 1361 | - | var offset: usize = 0; | |
| 1362 | - | var is_first = true; | |
| 1363 | - | ||
| 1364 | - | while (offset < file_content.len or is_first) { | |
| 1365 | - | const end = @min(offset + chunk_size, file_content.len); | |
| 1366 | - | const chunk = file_content[offset..end]; | |
| 1367 | - | ||
| 1368 | - | const encoded_len = std.base64.standard.Encoder.calcSize(chunk.len); | |
| 1369 | - | const encoded = try self.alloc.alloc(u8, encoded_len); | |
| 1370 | - | defer self.alloc.free(encoded); | |
| 1371 | - | _ = std.base64.standard.Encoder.encode(encoded, chunk); | |
| 1372 | - | ||
| 1373 | - | self.queuePtyInput("printf '%s' '"); | |
| 1374 | - | self.queuePtyInput(encoded); | |
| 1375 | - | if (is_first) { | |
| 1376 | - | self.queuePtyInput("' | base64 -d > '"); | |
| 1377 | - | } else { | |
| 1378 | - | self.queuePtyInput("' | base64 -d >> '"); | |
| 1379 | - | } | |
| 1380 | - | self.queuePtyInput(file_path); | |
| 1381 | - | self.queuePtyInput("'"); | |
| 1382 | - | self.queuePtyInput("\r"); | |
| 1383 | - | ||
| 1384 | - | offset = end; | |
| 1385 | - | is_first = false; | |
| 1386 | - | } | |
| 1387 | - | ||
| 1388 | - | try ipc.appendMessage(self.alloc, &client.write_buf, .Ack, ""); | |
| 1389 | - | client.has_pending_output = true; | |
| 1390 | - | self.has_had_client = true; | |
| 1391 | - | std.log.debug( | |
| 1392 | - | "write command len={d} file_path={s}", | |
| 1393 | - | .{ file_content.len, file_path }, | |
| 1394 | - | ); | |
| 1395 | - | } | |
| 1396 | - | }; | |
| 1397 | - | ||
| 1398 | - | test "send queues PTY input without changing leader" { | |
| 1399 | - | const alloc = std.testing.allocator; | |
| 1400 | - | var daemon = Daemon{ | |
| 1401 | - | .cfg = undefined, | |
| 1402 | - | .alloc = alloc, | |
| 1403 | - | .clients = .empty, | |
| 1404 | - | .leader_client_fd = 42, | |
| 1405 | - | .session_name = "test", | |
| 1406 | - | .socket_path = "", | |
| 1407 | - | .io = std.testing.io, | |
| 1408 | - | .running = true, | |
| 1409 | - | .pid = 0, | |
| 1410 | - | .created_at = 0, | |
| 1411 | - | }; | |
| 1412 | - | defer daemon.pty_write_buf.deinit(alloc); | |
| 1413 | - | ||
| 1414 | - | daemon.handleSend("hello"); | |
| 1415 | - | ||
| 1416 | - | try std.testing.expectEqual(@as(?i32, 42), daemon.leader_client_fd); | |
| 1417 | - | try std.testing.expectEqualStrings("hello", daemon.pty_write_buf.items); | |
| 1418 | - | } | |
| 1419 | - | ||
| 1420 | - | fn printVersion(io: std.Io, cfg: *Cfg) !void { | |
| 1421 | - | var buf: [256]u8 = undefined; | |
| 1422 | - | var w = std.Io.File.stdout().writer(io, &buf); | |
| 1423 | - | try w.interface.print( | |
| 1424 | - | "zmx\t\t{s}\nghostty_vt\t{s}\nsocket_dir\t{s}\nlog_dir\t\t{s}\n", | |
| 1425 | - | .{ version, ghostty_version, cfg.socket_dir, cfg.log_dir }, | |
| 1426 | - | ); | |
| 1427 | - | try w.interface.flush(); | |
| 1428 | - | } | |
| 1429 | - | ||
| 1430 | - | fn printCompletions(io: std.Io, shell: completions.Shell) !void { | |
| 1431 | - | const script = shell.getCompletionScript(); | |
| 1432 | - | var buf: [8192]u8 = undefined; | |
| 1433 | - | var w = std.Io.File.stdout().writer(io, &buf); | |
| 1434 | - | try w.interface.print("{s}\n", .{script}); | |
| 1435 | - | try w.interface.flush(); | |
| 1436 | - | } | |
| 1437 | - | ||
| 1438 | 407 | fn help(io: std.Io) !void { | |
| 1439 | 408 | const help_text = | |
| 1440 | 409 | \\zmx - session persistence for terminal processes |
| ... | ... | @@ -1578,6 +547,28 @@ fn help(io: std.Io) !void { | |
| 1578 | 547 | try w.interface.flush(); | |
| 1579 | 548 | } | |
| 1580 | 549 | ||
| 550 | + | fn printVersion(io: std.Io, cfg: *Cfg) !void { | |
| 551 | + | var buf: [256]u8 = undefined; | |
| 552 | + | var w = std.Io.File.stdout().writer(io, &buf); | |
| 553 | + | try w.interface.print( | |
| 554 | + | "zmx\t\t{s}\nghostty_vt\t{s}\nsocket_dir\t{s}\nlog_dir\t\t{s}\n", | |
| 555 | + | .{ version, ghostty_version, cfg.socket_dir, cfg.log_dir }, | |
| 556 | + | ); | |
| 557 | + | try w.interface.flush(); | |
| 558 | + | } | |
| 559 | + | ||
| 560 | + | fn printCompletions(io: std.Io, shell: completions.Shell) !void { | |
| 561 | + | const script = shell.getCompletionScript(); | |
| 562 | + | var buf: [8192]u8 = undefined; | |
| 563 | + | var w = std.Io.File.stdout().writer(io, &buf); | |
| 564 | + | try w.interface.print("{s}\n", .{script}); | |
| 565 | + | try w.interface.flush(); | |
| 566 | + | } | |
| 567 | + | ||
| 568 | + | fn detectHelp(arg: []const u8) bool { | |
| 569 | + | return (std.mem.eql(u8, arg, "--help") or std.mem.eql(u8, arg, "-h")); | |
| 570 | + | } | |
| 571 | + | ||
| 1581 | 572 | fn tail(alloc: std.mem.Allocator, client_socket_fds: std.ArrayList(i32), detached: bool, is_run_cmd: bool) !u8 { | |
| 1582 | 573 | var poll_fds = try std.ArrayList(lib_posix.pollfd).initCapacity(alloc, 4); | |
| 1583 | 574 | defer poll_fds.deinit(alloc); |
| ... | ... | @@ -1731,7 +722,7 @@ fn tail(alloc: std.mem.Allocator, client_socket_fds: std.ArrayList(i32), detache | |
| 1731 | 722 | } | |
| 1732 | 723 | } | |
| 1733 | 724 | ||
| 1734 | - | fn wait(alloc: std.mem.Allocator, io: std.Io, cfg: *Cfg, matchers: std.ArrayList(SessionMatch)) !void { | |
| 725 | + | fn wait(alloc: std.mem.Allocator, io: std.Io, cfg: *Cfg, matchers: std.ArrayList(socket.SessionMatch)) !void { | |
| 1735 | 726 | var stdout_buffer: [1024]u8 = undefined; | |
| 1736 | 727 | var stdout_writer = std.Io.File.stdout().writer(io, &stdout_buffer); | |
| 1737 | 728 | const stdout = &stdout_writer.interface; |
| ... | ... | @@ -1747,7 +738,7 @@ fn wait(alloc: std.mem.Allocator, io: std.Io, cfg: *Cfg, matchers: std.ArrayList | |
| 1747 | 738 | var zero_match_iters: u32 = 0; | |
| 1748 | 739 | ||
| 1749 | 740 | var agg_exit_code: u8 = 0; | |
| 1750 | - | var last_print: i96 = 0; | |
| 741 | + | var last_print: std.Io.Timestamp = .zero; | |
| 1751 | 742 | var prev_done: i32 = 0; | |
| 1752 | 743 | while (true) { | |
| 1753 | 744 | agg_exit_code = 0; |
| ... | ... | @@ -1775,7 +766,11 @@ fn wait(alloc: std.mem.Allocator, io: std.Io, cfg: *Cfg, matchers: std.ArrayList | |
| 1775 | 766 | // waiting". Count it as done+failed so wait terminates. | |
| 1776 | 767 | try stderr.print( | |
| 1777 | 768 | "[{d}] task unreachable: {s} ({s})\n", | |
| 1778 | - | .{ std.Io.Timestamp.now(io, .real).nanoseconds, session.name, session.error_name orelse "unknown" }, | |
| 769 | + | .{ | |
| 770 | + | std.Io.Timestamp.now(io, .real).toSeconds(), | |
| 771 | + | session.name, | |
| 772 | + | session.error_name orelse "unknown", | |
| 773 | + | }, | |
| 1779 | 774 | ); | |
| 1780 | 775 | try stderr.flush(); | |
| 1781 | 776 | agg_exit_code = 1; |
| ... | ... | @@ -1783,11 +778,11 @@ fn wait(alloc: std.mem.Allocator, io: std.Io, cfg: *Cfg, matchers: std.ArrayList | |
| 1783 | 778 | continue; | |
| 1784 | 779 | } | |
| 1785 | 780 | if (session.task_ended_at == 0) { | |
| 1786 | - | const now = std.Io.Timestamp.now(io, .real).nanoseconds; | |
| 1787 | - | if (now - last_print >= 5) { | |
| 781 | + | const now = std.Io.Timestamp.now(io, .real); | |
| 782 | + | if (now.toSeconds() - last_print.toSeconds() >= 5) { | |
| 1788 | 783 | try stdout.print( | |
| 1789 | 784 | "[{d}] waiting task={s}\n", | |
| 1790 | - | .{ now, session.name }, | |
| 785 | + | .{ now.toSeconds(), session.name }, | |
| 1791 | 786 | ); | |
| 1792 | 787 | try stdout.flush(); | |
| 1793 | 788 | last_print = now; |
| ... | ... | @@ -2254,32 +1249,32 @@ fn history(alloc: std.mem.Allocator, io: std.Io, cfg: *Cfg, session_name: []cons | |
| 2254 | 1249 | } | |
| 2255 | 1250 | } | |
| 2256 | 1251 | ||
| 2257 | - | fn switchSesh(daemon: *Daemon, current_sesh: []const u8) !void { | |
| 1252 | + | fn switchSesh(gpa: std.mem.Allocator, io: std.Io, daemon: *Daemon, current_sesh: []const u8) !void { | |
| 2258 | 1253 | // we want daemon.session_name because that's the session name the user provided during zmx attach | |
| 2259 | 1254 | // instead of the name of the session they are currently inside of. | |
| 2260 | 1255 | const next_session = daemon.session_name; | |
| 2261 | 1256 | std.log.info("switch session cur={s} next={s}", .{ current_sesh, next_session }); | |
| 2262 | 1257 | ||
| 2263 | - | const socket_path = socket.getSocketPath(daemon.alloc, daemon.cfg.socket_dir, current_sesh) catch |err| switch (err) { | |
| 2264 | - | error.NameTooLong => return socket.printSessionNameTooLong(daemon.io, current_sesh, daemon.cfg.socket_dir), | |
| 1258 | + | const socket_path = socket.getSocketPath(gpa, daemon.cfg.socket_dir, current_sesh) catch |err| switch (err) { | |
| 1259 | + | error.NameTooLong => return socket.printSessionNameTooLong(io, current_sesh, daemon.cfg.socket_dir), | |
| 2265 | 1260 | error.OutOfMemory => return err, | |
| 2266 | 1261 | }; | |
| 2267 | - | defer daemon.alloc.free(socket_path); | |
| 1262 | + | defer gpa.free(socket_path); | |
| 2268 | 1263 | ||
| 2269 | - | var dir = try std.Io.Dir.openDirAbsolute(daemon.io, daemon.cfg.socket_dir, .{}); | |
| 2270 | - | defer dir.close(daemon.io); | |
| 1264 | + | var dir = try std.Io.Dir.openDirAbsolute(io, daemon.cfg.socket_dir, .{}); | |
| 1265 | + | defer dir.close(io); | |
| 2271 | 1266 | ||
| 2272 | - | const exists = try socket.sessionExists(daemon.io, dir, current_sesh); | |
| 1267 | + | const exists = try socket.sessionExists(io, dir, current_sesh); | |
| 2273 | 1268 | if (!exists) { | |
| 2274 | 1269 | var buf: [4096]u8 = undefined; | |
| 2275 | - | var w = std.Io.File.stderr().writer(daemon.io, &buf); | |
| 1270 | + | var w = std.Io.File.stderr().writer(io, &buf); | |
| 2276 | 1271 | w.interface.print("error: session \"{s}\" does not exist\n", .{current_sesh}) catch {}; | |
| 2277 | 1272 | w.interface.flush() catch {}; | |
| 2278 | 1273 | return error.SessionNotFound; | |
| 2279 | 1274 | } | |
| 2280 | 1275 | const fd = ipc.connectSession(socket_path) catch |err| { | |
| 2281 | 1276 | std.log.err("session unresponsive: {s}", .{@errorName(err)}); | |
| 2282 | - | if (err == error.ConnectionRefused) socket.cleanupStaleSocket(daemon.io, dir, current_sesh); | |
| 1277 | + | if (err == error.ConnectionRefused) socket.cleanupStaleSocket(io, dir, current_sesh); | |
| 2283 | 1278 | return; | |
| 2284 | 1279 | }; | |
| 2285 | 1280 | defer lib_posix.close(fd); |
| ... | ... | @@ -2290,14 +1285,14 @@ fn switchSesh(daemon: *Daemon, current_sesh: []const u8) !void { | |
| 2290 | 1285 | }; | |
| 2291 | 1286 | } | |
| 2292 | 1287 | ||
| 2293 | - | fn attach(daemon: *Daemon) !void { | |
| 1288 | + | fn attach(gpa: std.mem.Allocator, io: std.Io, daemon: *Daemon) !void { | |
| 2294 | 1289 | const sesh = socket.getSeshNameFromEnv(); | |
| 2295 | 1290 | if (sesh.len > 0) { | |
| 2296 | - | return switchSesh(daemon, sesh); | |
| 1291 | + | return switchSesh(gpa, io, daemon, sesh); | |
| 2297 | 1292 | } | |
| 2298 | 1293 | ||
| 2299 | - | const result = try daemon.ensureSession(); | |
| 2300 | - | if (result.is_daemon) return; | |
| 1294 | + | const is_daemon_proc = try daemon.ensureSession(io); | |
| 1295 | + | if (is_daemon_proc) return; | |
| 2301 | 1296 | ||
| 2302 | 1297 | const client_sock = try socket.sessionConnect(daemon.socket_path); | |
| 2303 | 1298 | std.log.info("attached session={s}", .{daemon.session_name}); |
| ... | ... | @@ -2347,60 +1342,45 @@ fn attach(daemon: *Daemon) !void { | |
| 2347 | 1342 | const clear_seq = "\x1b[2J\x1b[H"; | |
| 2348 | 1343 | _ = try lib_posix.write(lib_posix.STDOUT_FILENO, clear_seq); | |
| 2349 | 1344 | ||
| 2350 | - | const looper = try clientLoop(client_sock); | |
| 1345 | + | const looper = try loop.clientLoop(client_sock); | |
| 2351 | 1346 | switch (looper.kind) { | |
| 2352 | 1347 | .detach => return, | |
| 2353 | 1348 | .switch_session => { | |
| 2354 | 1349 | if (looper.session_name) |session_name| { | |
| 2355 | 1350 | var cwd_buf: [std.fs.max_path_bytes]u8 = undefined; | |
| 2356 | - | const cwd_len = std.process.currentPath(daemon.io, &cwd_buf) catch 0; | |
| 1351 | + | const cwd_len = std.process.currentPath(io, &cwd_buf) catch 0; | |
| 2357 | 1352 | const cwd = cwd_buf[0..cwd_len]; | |
| 2358 | 1353 | const target_path = socket.getSocketPath( | |
| 2359 | - | daemon.alloc, | |
| 1354 | + | gpa, | |
| 2360 | 1355 | daemon.cfg.socket_dir, | |
| 2361 | 1356 | session_name, | |
| 2362 | 1357 | ) catch |err| switch (err) { | |
| 2363 | 1358 | error.NameTooLong => return socket.printSessionNameTooLong( | |
| 2364 | - | daemon.io, | |
| 1359 | + | io, | |
| 2365 | 1360 | session_name, | |
| 2366 | 1361 | daemon.cfg.socket_dir, | |
| 2367 | 1362 | ), | |
| 2368 | 1363 | error.OutOfMemory => return err, | |
| 2369 | 1364 | }; | |
| 2370 | 1365 | ||
| 2371 | - | const clients = try std.ArrayList(*Client).initCapacity(daemon.alloc, 10); | |
| 2372 | - | var target_daemon = Daemon{ | |
| 2373 | - | .io = daemon.io, | |
| 2374 | - | .running = true, | |
| 2375 | - | .cfg = daemon.cfg, | |
| 2376 | - | .alloc = daemon.alloc, | |
| 2377 | - | .clients = clients, | |
| 2378 | - | .session_name = session_name, | |
| 2379 | - | .socket_path = target_path, | |
| 2380 | - | .pid = undefined, | |
| 2381 | - | .cwd = cwd, | |
| 2382 | - | .created_at = @intCast(std.Io.Timestamp.now(daemon.io, .real).nanoseconds), | |
| 2383 | - | .leader_client_fd = null, | |
| 2384 | - | }; | |
| 2385 | - | return attach(&target_daemon); | |
| 1366 | + | var target_daemon = Daemon.init(io, daemon.cfg, session_name, target_path); | |
| 1367 | + | target_daemon.cwd = cwd; | |
| 1368 | + | return attach(gpa, io, &target_daemon); | |
| 2386 | 1369 | } | |
| 2387 | 1370 | }, | |
| 2388 | 1371 | } | |
| 2389 | 1372 | } | |
| 2390 | 1373 | ||
| 2391 | - | fn writeFile(daemon: *Daemon, file_path: []const u8) !void { | |
| 1374 | + | fn writeFile(gpa: std.mem.Allocator, io: std.Io, daemon: *Daemon, file_path: []const u8) !void { | |
| 1375 | + | const is_daemon_proc = try daemon.ensureSession(io); | |
| 1376 | + | if (is_daemon_proc) return; | |
| 1377 | + | ||
| 2392 | 1378 | var buf: [4096]u8 = undefined; | |
| 2393 | - | var w = std.Io.File.stdout().writer(daemon.io, &buf); | |
| 2394 | - | const sesh_result = try daemon.ensureSession(); | |
| 2395 | - | if (sesh_result.is_daemon) return; | |
| 1379 | + | var w = std.Io.File.stdout().writer(io, &buf); | |
| 2396 | 1380 | ||
| 2397 | - | if (sesh_result.created) { | |
| 2398 | - | try w.interface.print("session \"{s}\" created\n", .{daemon.session_name}); | |
| 2399 | - | try w.interface.flush(); | |
| 2400 | - | } | |
| 2401 | 1381 | const stdin_fd = lib_posix.STDIN_FILENO; | |
| 2402 | - | var stdin_buf = try std.ArrayList(u8).initCapacity(daemon.alloc, 4096); | |
| 2403 | - | defer stdin_buf.deinit(daemon.alloc); | |
| 1382 | + | var stdin_buf = try std.ArrayList(u8).initCapacity(gpa, 4096); | |
| 1383 | + | defer stdin_buf.deinit(gpa); | |
| 2404 | 1384 | ||
| 2405 | 1385 | while (true) { | |
| 2406 | 1386 | var tmp: [4096]u8 = undefined; |
| ... | ... | @@ -2409,28 +1389,28 @@ fn writeFile(daemon: *Daemon, file_path: []const u8) !void { | |
| 2409 | 1389 | return err; | |
| 2410 | 1390 | }; | |
| 2411 | 1391 | if (n == 0) break; | |
| 2412 | - | try stdin_buf.appendSlice(daemon.alloc, tmp[0..n]); | |
| 1392 | + | try stdin_buf.appendSlice(gpa, tmp[0..n]); | |
| 2413 | 1393 | } | |
| 2414 | 1394 | ||
| 2415 | 1395 | const socket_path = socket.getSocketPath( | |
| 2416 | - | daemon.alloc, | |
| 1396 | + | gpa, | |
| 2417 | 1397 | daemon.cfg.socket_dir, | |
| 2418 | 1398 | daemon.session_name, | |
| 2419 | 1399 | ) catch |err| switch (err) { | |
| 2420 | 1400 | error.NameTooLong => return socket.printSessionNameTooLong( | |
| 2421 | - | daemon.io, | |
| 1401 | + | io, | |
| 2422 | 1402 | daemon.session_name, | |
| 2423 | 1403 | daemon.cfg.socket_dir, | |
| 2424 | 1404 | ), | |
| 2425 | 1405 | error.OutOfMemory => return err, | |
| 2426 | 1406 | }; | |
| 2427 | - | var dir = try std.Io.Dir.openDirAbsolute(daemon.io, daemon.cfg.socket_dir, .{}); | |
| 2428 | - | defer dir.close(daemon.io); | |
| 1407 | + | var dir = try std.Io.Dir.openDirAbsolute(io, daemon.cfg.socket_dir, .{}); | |
| 1408 | + | defer dir.close(io); | |
| 2429 | 1409 | ||
| 2430 | - | const result = ipc.probeSession(daemon.alloc, socket_path) catch |err| { | |
| 1410 | + | const result = ipc.probeSession(gpa, socket_path) catch |err| { | |
| 2431 | 1411 | std.log.err("session unresponsive: {s}", .{@errorName(err)}); | |
| 2432 | 1412 | if (err == error.ConnectionRefused) { | |
| 2433 | - | socket.cleanupStaleSocket(daemon.io, dir, daemon.session_name); | |
| 1413 | + | socket.cleanupStaleSocket(io, dir, daemon.session_name); | |
| 2434 | 1414 | w.interface.print("cleaned up stale session {s}\n", .{daemon.session_name}) catch {}; | |
| 2435 | 1415 | } else { | |
| 2436 | 1416 | w.interface.print( |
| ... | ... | @@ -2446,21 +1426,21 @@ fn writeFile(daemon: *Daemon, file_path: []const u8) !void { | |
| 2446 | 1426 | ||
| 2447 | 1427 | // Build wire payload: [u32 path len][path bytes][file content] | |
| 2448 | 1428 | var wire_buf = try std.ArrayList(u8).initCapacity( | |
| 2449 | - | daemon.alloc, | |
| 1429 | + | gpa, | |
| 2450 | 1430 | @sizeOf(u32) + file_path.len + stdin_buf.items.len, | |
| 2451 | 1431 | ); | |
| 2452 | - | defer wire_buf.deinit(daemon.alloc); | |
| 1432 | + | defer wire_buf.deinit(gpa); | |
| 2453 | 1433 | const path_len: u32 = @intCast(file_path.len); | |
| 2454 | - | try wire_buf.appendSlice(daemon.alloc, std.mem.asBytes(&path_len)); | |
| 2455 | - | try wire_buf.appendSlice(daemon.alloc, file_path); | |
| 2456 | - | try wire_buf.appendSlice(daemon.alloc, stdin_buf.items); | |
| 1434 | + | try wire_buf.appendSlice(gpa, std.mem.asBytes(&path_len)); | |
| 1435 | + | try wire_buf.appendSlice(gpa, file_path); | |
| 1436 | + | try wire_buf.appendSlice(gpa, stdin_buf.items); | |
| 2457 | 1437 | ||
| 2458 | 1438 | ipc.send(result.fd, .Write, wire_buf.items) catch |err| switch (err) { | |
| 2459 | 1439 | error.BrokenPipe, error.ConnectionResetByPeer => return, | |
| 2460 | 1440 | else => return err, | |
| 2461 | 1441 | }; | |
| 2462 | 1442 | ||
| 2463 | - | var sb = try ipc.SocketBuffer.init(daemon.alloc); | |
| 1443 | + | var sb = try ipc.SocketBuffer.init(gpa); | |
| 2464 | 1444 | defer sb.deinit(); | |
| 2465 | 1445 | ||
| 2466 | 1446 | const n = sb.read(result.fd) catch return error.ReadFailed; |
| ... | ... | @@ -2539,35 +1519,26 @@ fn send(alloc: std.mem.Allocator, io: std.Io, cfg: *Cfg, session_name: []const u | |
| 2539 | 1519 | }; | |
| 2540 | 1520 | } | |
| 2541 | 1521 | ||
| 2542 | - | fn run(daemon: *Daemon, detached: bool, command_args: [][]const u8) !void { | |
| 2543 | - | const alloc = daemon.alloc; | |
| 2544 | - | var buf: [4096]u8 = undefined; | |
| 2545 | - | var w = std.Io.File.stdout().writer(daemon.io, &buf); | |
| 2546 | - | ||
| 1522 | + | fn run(gpa: std.mem.Allocator, io: std.Io, daemon: *Daemon, detached: bool, command_args: [][]const u8) !void { | |
| 2547 | 1523 | var cmd_to_send: ?[]const u8 = null; | |
| 2548 | 1524 | var allocated_cmd: ?[]u8 = null; | |
| 2549 | - | defer if (allocated_cmd) |cmd| alloc.free(cmd); | |
| 1525 | + | defer if (allocated_cmd) |cmd| gpa.free(cmd); | |
| 2550 | 1526 | ||
| 2551 | - | const result = try daemon.ensureSession(); | |
| 2552 | - | if (result.is_daemon) return; | |
| 2553 | - | ||
| 2554 | - | if (result.created) { | |
| 2555 | - | try w.interface.print("session \"{s}\" created\n", .{daemon.session_name}); | |
| 2556 | - | try w.interface.flush(); | |
| 2557 | - | } | |
| 1527 | + | const is_daemon_proc = try daemon.ensureSession(io); | |
| 1528 | + | if (is_daemon_proc) return; | |
| 2558 | 1529 | ||
| 2559 | 1530 | if (command_args.len > 0) { | |
| 2560 | 1531 | var cmd_list = std.ArrayList(u8).empty; | |
| 2561 | - | defer cmd_list.deinit(alloc); | |
| 1532 | + | defer cmd_list.deinit(gpa); | |
| 2562 | 1533 | ||
| 2563 | 1534 | for (command_args, 0..) |arg, i| { | |
| 2564 | - | if (i > 0) try cmd_list.append(alloc, ' '); | |
| 1535 | + | if (i > 0) try cmd_list.append(gpa, ' '); | |
| 2565 | 1536 | if (util.shellNeedsQuoting(arg)) { | |
| 2566 | - | const quoted = try util.shellQuote(alloc, arg); | |
| 2567 | - | defer alloc.free(quoted); | |
| 2568 | - | try cmd_list.appendSlice(alloc, quoted); | |
| 1537 | + | const quoted = try util.shellQuote(gpa, arg); | |
| 1538 | + | defer gpa.free(quoted); | |
| 1539 | + | try cmd_list.appendSlice(gpa, quoted); | |
| 2569 | 1540 | } else { | |
| 2570 | - | try cmd_list.appendSlice(alloc, arg); | |
| 1541 | + | try cmd_list.appendSlice(gpa, arg); | |
| 2571 | 1542 | } | |
| 2572 | 1543 | } | |
| 2573 | 1544 |
| ... | ... | @@ -2575,24 +1546,24 @@ fn run(daemon: *Daemon, detached: bool, command_args: [][]const u8) !void { | |
| 2575 | 1546 | // raw mode; readline's accept-line binds to CR. The first-ever run | |
| 2576 | 1547 | // works with \n only because it arrives during shell startup while | |
| 2577 | 1548 | // the line discipline is still canonical. | |
| 2578 | - | try cmd_list.append(alloc, '\r'); | |
| 1549 | + | try cmd_list.append(gpa, '\r'); | |
| 2579 | 1550 | ||
| 2580 | - | cmd_to_send = try cmd_list.toOwnedSlice(alloc); | |
| 1551 | + | cmd_to_send = try cmd_list.toOwnedSlice(gpa); | |
| 2581 | 1552 | allocated_cmd = @constCast(cmd_to_send.?); | |
| 2582 | 1553 | } else { | |
| 2583 | 1554 | // Read from stdin when no text arguments provided. | |
| 2584 | 1555 | const stdin_file = std.Io.File.stdin(); | |
| 2585 | - | defer stdin_file.close(daemon.io); | |
| 2586 | - | var stdin_buf = try std.ArrayList(u8).initCapacity(alloc, 4096); | |
| 2587 | - | defer stdin_buf.deinit(alloc); | |
| 1556 | + | defer stdin_file.close(io); | |
| 1557 | + | var stdin_buf = try std.ArrayList(u8).initCapacity(gpa, 4096); | |
| 1558 | + | defer stdin_buf.deinit(gpa); | |
| 2588 | 1559 | var stdbuf: [4096]u8 = undefined; | |
| 2589 | - | var reader = stdin_file.reader(daemon.io, &stdbuf); | |
| 2590 | - | if (!try stdin_file.isTty(daemon.io)) { | |
| 1560 | + | var reader = stdin_file.reader(io, &stdbuf); | |
| 1561 | + | if (!try stdin_file.isTty(io)) { | |
| 2591 | 1562 | while (true) { | |
| 2592 | 1563 | var dest: [1024]u8 = undefined; | |
| 2593 | 1564 | const n = try reader.interface.readSliceShort(&dest); | |
| 2594 | 1565 | if (n == 0) break; // EOF | |
| 2595 | - | try stdin_buf.appendSlice(alloc, dest[0..n]); | |
| 1566 | + | try stdin_buf.appendSlice(gpa, dest[0..n]); | |
| 2596 | 1567 | } | |
| 2597 | 1568 | ||
| 2598 | 1569 | if (stdin_buf.items.len > 0) { |
| ... | ... | @@ -2601,42 +1572,13 @@ fn run(daemon: *Daemon, detached: bool, command_args: [][]const u8) !void { | |
| 2601 | 1572 | if (stdin_buf.items[stdin_buf.items.len - 1] == '\n') { | |
| 2602 | 1573 | stdin_buf.items[stdin_buf.items.len - 1] = '\r'; | |
| 2603 | 1574 | } else { | |
| 2604 | - | try stdin_buf.append(alloc, '\r'); | |
| 1575 | + | try stdin_buf.append(gpa, '\r'); | |
| 2605 | 1576 | } | |
| 2606 | 1577 | ||
| 2607 | - | cmd_to_send = try alloc.dupe(u8, stdin_buf.items); | |
| 1578 | + | cmd_to_send = try gpa.dupe(u8, stdin_buf.items); | |
| 2608 | 1579 | allocated_cmd = @constCast(cmd_to_send.?); | |
| 2609 | 1580 | } | |
| 2610 | 1581 | } | |
| 2611 | - | ||
| 2612 | - | // const stdin_fd = posix.STDIN_FILENO; | |
| 2613 | - | // if (!lib_posix.isatty(stdin_fd)) { | |
| 2614 | - | // var stdin_buf = try std.ArrayList(u8).initCapacity(alloc, 4096); | |
| 2615 | - | // defer stdin_buf.deinit(alloc); | |
| 2616 | - | ||
| 2617 | - | // while (true) { | |
| 2618 | - | // var tmp: [4096]u8 = undefined; | |
| 2619 | - | // const n = posix.read(stdin_fd, &tmp) catch |err| { | |
| 2620 | - | // if (err == error.WouldBlock) break; | |
| 2621 | - | // return err; | |
| 2622 | - | // }; | |
| 2623 | - | // if (n == 0) break; | |
| 2624 | - | // try stdin_buf.appendSlice(alloc, tmp[0..n]); | |
| 2625 | - | // } | |
| 2626 | - | ||
| 2627 | - | // if (stdin_buf.items.len > 0) { | |
| 2628 | - | // // Normalize any trailing newline to CR so readline (raw mode) | |
| 2629 | - | // // accepts each line. | |
| 2630 | - | // if (stdin_buf.items[stdin_buf.items.len - 1] == '\n') { | |
| 2631 | - | // stdin_buf.items[stdin_buf.items.len - 1] = '\r'; | |
| 2632 | - | // } else { | |
| 2633 | - | // try stdin_buf.append(alloc, '\r'); | |
| 2634 | - | // } | |
| 2635 | - | ||
| 2636 | - | // cmd_to_send = try alloc.dupe(u8, stdin_buf.items); | |
| 2637 | - | // allocated_cmd = @constCast(cmd_to_send.?); | |
| 2638 | - | // } | |
| 2639 | - | // } | |
| 2640 | 1582 | } | |
| 2641 | 1583 | ||
| 2642 | 1584 | if (cmd_to_send == null) { |
| ... | ... | @@ -2649,529 +1591,15 @@ fn run(daemon: *Daemon, detached: bool, command_args: [][]const u8) !void { | |
| 2649 | 1591 | }; | |
| 2650 | 1592 | defer lib_posix.close(client_sock); | |
| 2651 | 1593 | ||
| 2652 | - | var fds = try std.ArrayList(i32).initCapacity(alloc, 1); | |
| 2653 | - | defer fds.deinit(alloc); | |
| 2654 | - | try fds.append(alloc, client_sock); | |
| 1594 | + | var fds = try std.ArrayList(i32).initCapacity(gpa, 1); | |
| 1595 | + | defer fds.deinit(gpa); | |
| 1596 | + | try fds.append(gpa, client_sock); | |
| 2655 | 1597 | ||
| 2656 | 1598 | ipc.send(client_sock, .Run, cmd_to_send.?) catch |err| switch (err) { | |
| 2657 | 1599 | error.ConnectionResetByPeer, error.BrokenPipe => return, | |
| 2658 | 1600 | else => return err, | |
| 2659 | 1601 | }; | |
| 2660 | 1602 | ||
| 2661 | - | const exit_code = try tail(daemon.alloc, fds, detached, true); | |
| 1603 | + | const exit_code = try tail(gpa, fds, detached, true); | |
| 2662 | 1604 | lib_posix.exit(exit_code); | |
| 2663 | 1605 | } | |
| 2664 | - | ||
| 2665 | - | const ClientResult = struct { | |
| 2666 | - | kind: enum { | |
| 2667 | - | detach, | |
| 2668 | - | switch_session, | |
| 2669 | - | }, | |
| 2670 | - | session_name: ?[]const u8, | |
| 2671 | - | }; | |
| 2672 | - | ||
| 2673 | - | /// clientLoop sends ipc commands to its corresponding daemon. It uses poll() as its non-blocking | |
| 2674 | - | /// mechanism. It will send stdin to the daemon and receive stdout from the daemon. | |
| 2675 | - | fn clientLoop(client_sock_fd: i32) !ClientResult { | |
| 2676 | - | std.log.info("client loop fd={d}", .{client_sock_fd}); | |
| 2677 | - | // use c_allocator to avoid "reached unreachable code" panic in DebugAllocator when forking | |
| 2678 | - | const alloc = std.heap.c_allocator; | |
| 2679 | - | defer lib_posix.close(client_sock_fd); | |
| 2680 | - | ||
| 2681 | - | try openSignalPipe(); | |
| 2682 | - | installWakeHandler(@intFromEnum(lib_posix.SIG.WINCH)); | |
| 2683 | - | ||
| 2684 | - | // Make socket non-blocking to avoid blocking on writes | |
| 2685 | - | var sock_flags = try lib_posix.fcntl(client_sock_fd, lib_posix.F.GETFL, 0); | |
| 2686 | - | sock_flags |= O_NONBLOCK; | |
| 2687 | - | _ = try lib_posix.fcntl(client_sock_fd, lib_posix.F.SETFL, sock_flags); | |
| 2688 | - | ||
| 2689 | - | // Buffer for outgoing socket writes | |
| 2690 | - | var sock_write_buf = try std.ArrayList(u8).initCapacity(alloc, 4096); | |
| 2691 | - | defer sock_write_buf.deinit(alloc); | |
| 2692 | - | ||
| 2693 | - | // Send init message with terminal size (buffered) | |
| 2694 | - | const size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO); | |
| 2695 | - | try ipc.appendMessage(alloc, &sock_write_buf, .Init, std.mem.asBytes(&size)); | |
| 2696 | - | ||
| 2697 | - | var poll_fds = try std.ArrayList(lib_posix.pollfd).initCapacity(alloc, 4); | |
| 2698 | - | defer poll_fds.deinit(alloc); | |
| 2699 | - | ||
| 2700 | - | var read_buf = try ipc.SocketBuffer.init(alloc); | |
| 2701 | - | defer read_buf.deinit(); | |
| 2702 | - | ||
| 2703 | - | var stdout_buf = try std.ArrayList(u8).initCapacity(alloc, 4096); | |
| 2704 | - | defer stdout_buf.deinit(alloc); | |
| 2705 | - | ||
| 2706 | - | const stdin_fd = lib_posix.STDIN_FILENO; | |
| 2707 | - | ||
| 2708 | - | // Make stdin non-blocking. O_NONBLOCK is set on the open file description, | |
| 2709 | - | // which is shared with the parent shell; restore on exit to avoid | |
| 2710 | - | // corrupting the parent's stdin. | |
| 2711 | - | const stdin_orig_flags = try lib_posix.fcntl(stdin_fd, lib_posix.F.GETFL, 0); | |
| 2712 | - | _ = try lib_posix.fcntl(stdin_fd, lib_posix.F.SETFL, stdin_orig_flags | O_NONBLOCK); | |
| 2713 | - | defer _ = lib_posix.fcntl(stdin_fd, lib_posix.F.SETFL, stdin_orig_flags) catch {}; | |
| 2714 | - | ||
| 2715 | - | while (true) { | |
| 2716 | - | poll_fds.clearRetainingCapacity(); | |
| 2717 | - | ||
| 2718 | - | try poll_fds.append(alloc, .{ | |
| 2719 | - | .fd = stdin_fd, | |
| 2720 | - | .events = lib_posix.POLL.IN, | |
| 2721 | - | .revents = 0, | |
| 2722 | - | }); | |
| 2723 | - | ||
| 2724 | - | // Poll socket for read, and also for write if we have pending data | |
| 2725 | - | var sock_events: i16 = lib_posix.POLL.IN; | |
| 2726 | - | if (sock_write_buf.items.len > 0) { | |
| 2727 | - | sock_events |= lib_posix.POLL.OUT; | |
| 2728 | - | } | |
| 2729 | - | try poll_fds.append(alloc, .{ | |
| 2730 | - | .fd = client_sock_fd, | |
| 2731 | - | .events = sock_events, | |
| 2732 | - | .revents = 0, | |
| 2733 | - | }); | |
| 2734 | - | ||
| 2735 | - | try poll_fds.append(alloc, .{ .fd = sig_pipe[0], .events = lib_posix.POLL.IN, .revents = 0 }); | |
| 2736 | - | ||
| 2737 | - | if (stdout_buf.items.len > 0) { | |
| 2738 | - | try poll_fds.append(alloc, .{ | |
| 2739 | - | .fd = lib_posix.STDOUT_FILENO, | |
| 2740 | - | .events = lib_posix.POLL.OUT, | |
| 2741 | - | .revents = 0, | |
| 2742 | - | }); | |
| 2743 | - | } | |
| 2744 | - | ||
| 2745 | - | _ = try lib_posix.poll(poll_fds.items, -1); | |
| 2746 | - | ||
| 2747 | - | if (poll_fds.items[2].revents & lib_posix.POLL.IN != 0) { | |
| 2748 | - | drainSignalPipe(); | |
| 2749 | - | const next_size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO); | |
| 2750 | - | try ipc.appendMessage(alloc, &sock_write_buf, .Resize, std.mem.asBytes(&next_size)); | |
| 2751 | - | } | |
| 2752 | - | ||
| 2753 | - | // Handle stdin -> socket (Input) | |
| 2754 | - | const inp_flags = (lib_posix.POLL.IN | lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL); | |
| 2755 | - | if (poll_fds.items[0].revents & inp_flags != 0) { | |
| 2756 | - | var buf: [4096]u8 = undefined; | |
| 2757 | - | const n_opt: ?usize = lib_posix.read(stdin_fd, &buf) catch |err| blk: { | |
| 2758 | - | if (err == error.WouldBlock) break :blk null; | |
| 2759 | - | return err; | |
| 2760 | - | }; | |
| 2761 | - | ||
| 2762 | - | if (n_opt) |n| { | |
| 2763 | - | if (n > 0) { | |
| 2764 | - | // Check for detach sequences (ctrl+\ as first byte or Kitty escape sequence) | |
| 2765 | - | if (util.isCtrlBackslash(buf[0..n])) { | |
| 2766 | - | std.log.info("detach key detected", .{}); | |
| 2767 | - | try ipc.appendMessage(alloc, &sock_write_buf, .Detach, ""); | |
| 2768 | - | } else { | |
| 2769 | - | try ipc.appendMessage(alloc, &sock_write_buf, .Input, buf[0..n]); | |
| 2770 | - | } | |
| 2771 | - | } else { | |
| 2772 | - | std.log.info("eof stdin", .{}); | |
| 2773 | - | // EOF on stdin | |
| 2774 | - | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 2775 | - | } | |
| 2776 | - | } | |
| 2777 | - | } | |
| 2778 | - | ||
| 2779 | - | // Handle socket read (incoming Output messages from daemon) | |
| 2780 | - | if (poll_fds.items[1].revents & lib_posix.POLL.IN != 0) { | |
| 2781 | - | const n = read_buf.read(client_sock_fd) catch |err| { | |
| 2782 | - | if (err == error.WouldBlock) continue; | |
| 2783 | - | if (err == error.ConnectionResetByPeer or err == error.BrokenPipe) { | |
| 2784 | - | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 2785 | - | } | |
| 2786 | - | std.log.err("daemon read err={s}", .{@errorName(err)}); | |
| 2787 | - | return err; | |
| 2788 | - | }; | |
| 2789 | - | if (n == 0) { | |
| 2790 | - | std.log.info("server closed connection", .{}); | |
| 2791 | - | // Server closed connection | |
| 2792 | - | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 2793 | - | } | |
| 2794 | - | ||
| 2795 | - | while (read_buf.next()) |msg| { | |
| 2796 | - | switch (msg.header.tag) { | |
| 2797 | - | .Output => { | |
| 2798 | - | if (msg.payload.len > 0) { | |
| 2799 | - | try stdout_buf.appendSlice(alloc, msg.payload); | |
| 2800 | - | } | |
| 2801 | - | }, | |
| 2802 | - | .Resize => { | |
| 2803 | - | // daemon is asking for the client's window size usually in response | |
| 2804 | - | // to this client being set as leader. | |
| 2805 | - | const next_size = ipc.getTerminalSize(lib_posix.STDOUT_FILENO); | |
| 2806 | - | try ipc.appendMessage( | |
| 2807 | - | alloc, | |
| 2808 | - | &sock_write_buf, | |
| 2809 | - | .Resize, | |
| 2810 | - | std.mem.asBytes(&next_size), | |
| 2811 | - | ); | |
| 2812 | - | }, | |
| 2813 | - | .Switch => { | |
| 2814 | - | std.log.info("switch session", .{}); | |
| 2815 | - | return ClientResult{ .kind = .switch_session, .session_name = try alloc.dupe(u8, msg.payload) }; | |
| 2816 | - | }, | |
| 2817 | - | else => {}, | |
| 2818 | - | } | |
| 2819 | - | } | |
| 2820 | - | } | |
| 2821 | - | ||
| 2822 | - | // Handle socket write (flush buffered messages to daemon) | |
| 2823 | - | if (poll_fds.items[1].revents & lib_posix.POLL.OUT != 0) { | |
| 2824 | - | if (sock_write_buf.items.len > 0) { | |
| 2825 | - | const n = lib_posix.write(client_sock_fd, sock_write_buf.items) catch |err| blk: { | |
| 2826 | - | if (err == error.WouldBlock) break :blk 0; | |
| 2827 | - | if (err == error.ConnectionResetByPeer or err == error.BrokenPipe) { | |
| 2828 | - | std.log.info("connection reset or broken pipe", .{}); | |
| 2829 | - | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 2830 | - | } | |
| 2831 | - | return err; | |
| 2832 | - | }; | |
| 2833 | - | if (n > 0) { | |
| 2834 | - | try sock_write_buf.replaceRange(alloc, 0, n, &[_]u8{}); | |
| 2835 | - | } | |
| 2836 | - | } | |
| 2837 | - | } | |
| 2838 | - | ||
| 2839 | - | if (stdout_buf.items.len > 0) { | |
| 2840 | - | const n = lib_posix.write(lib_posix.STDOUT_FILENO, stdout_buf.items) catch |err| blk: { | |
| 2841 | - | if (err == error.WouldBlock) break :blk 0; | |
| 2842 | - | return err; | |
| 2843 | - | }; | |
| 2844 | - | if (n > 0) { | |
| 2845 | - | try stdout_buf.replaceRange(alloc, 0, n, &[_]u8{}); | |
| 2846 | - | } | |
| 2847 | - | } | |
| 2848 | - | ||
| 2849 | - | if (poll_fds.items[1].revents & (lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL) != 0) { | |
| 2850 | - | std.log.info("poll hup|err|nval", .{}); | |
| 2851 | - | return ClientResult{ .kind = .detach, .session_name = null }; | |
| 2852 | - | } | |
| 2853 | - | } | |
| 2854 | - | } | |
| 2855 | - | ||
| 2856 | - | /// dameonLoop is what the daemon runs to send and receive ipc commands from its corresponding | |
| 2857 | - | /// clients. It uses poll() as its non-blocking mechanism. | |
| 2858 | - | fn daemonLoop(daemon: *Daemon, server_sock_fd: i32, pty_fd: i32) !void { | |
| 2859 | - | std.log.info("daemon started session={s} pty_fd={d}", .{ daemon.session_name, pty_fd }); | |
| 2860 | - | daemon.pty_fd = pty_fd; | |
| 2861 | - | try openSignalPipe(); | |
| 2862 | - | installWakeHandler(@intFromEnum(lib_posix.SIG.TERM)); | |
| 2863 | - | var poll_fds = try std.ArrayList(lib_posix.pollfd).initCapacity(daemon.alloc, 8); | |
| 2864 | - | defer poll_fds.deinit(daemon.alloc); | |
| 2865 | - | ||
| 2866 | - | const init_size = ipc.getTerminalSize(pty_fd); | |
| 2867 | - | var term = try ghostty_vt.Terminal.init(daemon.io, daemon.alloc, .{ | |
| 2868 | - | .cols = init_size.cols, | |
| 2869 | - | .rows = init_size.rows, | |
| 2870 | - | .max_scrollback = daemon.cfg.max_scrollback, | |
| 2871 | - | }); | |
| 2872 | - | defer term.deinit(daemon.alloc); | |
| 2873 | - | var vt_stream = term.vtStream(); | |
| 2874 | - | defer vt_stream.deinit(); | |
| 2875 | - | ||
| 2876 | - | // Carries the tail of the previous PTY read so the task-exit marker | |
| 2877 | - | // search below can see across a read() boundary. Sized to comfortably | |
| 2878 | - | // hold "ZMX_TASK_COMPLETED:" (19 bytes) plus a u8 exit code and CRLF. | |
| 2879 | - | var marker_carry: [32]u8 = undefined; | |
| 2880 | - | var marker_carry_len: usize = 0; | |
| 2881 | - | ||
| 2882 | - | daemon_loop: while (daemon.running) { | |
| 2883 | - | poll_fds.clearRetainingCapacity(); | |
| 2884 | - | ||
| 2885 | - | try poll_fds.append(daemon.alloc, .{ | |
| 2886 | - | .fd = server_sock_fd, | |
| 2887 | - | .events = lib_posix.POLL.IN, | |
| 2888 | - | .revents = 0, | |
| 2889 | - | }); | |
| 2890 | - | ||
| 2891 | - | var pty_events: i16 = lib_posix.POLL.IN; | |
| 2892 | - | if (daemon.pty_write_buf.items.len > 0) { | |
| 2893 | - | pty_events |= lib_posix.POLL.OUT; | |
| 2894 | - | } | |
| 2895 | - | try poll_fds.append(daemon.alloc, .{ | |
| 2896 | - | .fd = pty_fd, | |
| 2897 | - | .events = pty_events, | |
| 2898 | - | .revents = 0, | |
| 2899 | - | }); | |
| 2900 | - | ||
| 2901 | - | try poll_fds.append(daemon.alloc, .{ .fd = sig_pipe[0], .events = lib_posix.POLL.IN, .revents = 0 }); | |
| 2902 | - | ||
| 2903 | - | for (daemon.clients.items) |client| { | |
| 2904 | - | var events: i16 = lib_posix.POLL.IN; | |
| 2905 | - | if (client.has_pending_output) { | |
| 2906 | - | events |= lib_posix.POLL.OUT; | |
| 2907 | - | } | |
| 2908 | - | try poll_fds.append(daemon.alloc, .{ | |
| 2909 | - | .fd = client.socket_fd, | |
| 2910 | - | .events = events, | |
| 2911 | - | .revents = 0, | |
| 2912 | - | }); | |
| 2913 | - | } | |
| 2914 | - | ||
| 2915 | - | _ = try lib_posix.poll(poll_fds.items, -1); | |
| 2916 | - | ||
| 2917 | - | if (poll_fds.items[2].revents & lib_posix.POLL.IN != 0) { | |
| 2918 | - | drainSignalPipe(); | |
| 2919 | - | std.log.info( | |
| 2920 | - | "SIGTERM received, shutting down gracefully session={s}", | |
| 2921 | - | .{daemon.session_name}, | |
| 2922 | - | ); | |
| 2923 | - | break :daemon_loop; | |
| 2924 | - | } | |
| 2925 | - | ||
| 2926 | - | if (poll_fds.items[0].revents & (lib_posix.POLL.ERR | lib_posix.POLL.HUP | lib_posix.POLL.NVAL) != 0) { | |
| 2927 | - | std.log.err("server socket error revents={d}", .{poll_fds.items[0].revents}); | |
| 2928 | - | break :daemon_loop; | |
| 2929 | - | } else if (poll_fds.items[0].revents & lib_posix.POLL.IN != 0) { | |
| 2930 | - | const client_fd = try lib_posix.accept( | |
| 2931 | - | server_sock_fd, | |
| 2932 | - | null, | |
| 2933 | - | null, | |
| 2934 | - | lib_posix.SOCK.NONBLOCK | lib_posix.SOCK.CLOEXEC, | |
| 2935 | - | ); | |
| 2936 | - | const client = try daemon.alloc.create(Client); | |
| 2937 | - | client.* = Client{ | |
| 2938 | - | .alloc = daemon.alloc, | |
| 2939 | - | .socket_fd = client_fd, | |
| 2940 | - | .read_buf = try ipc.SocketBuffer.init(daemon.alloc), | |
| 2941 | - | .write_buf = undefined, | |
| 2942 | - | }; | |
| 2943 | - | // 64KB initial capacity lets ~15 broadcast cycles (N_TTY_BUF_SIZE reads | |
| 2944 | - | // * header) accumulate before the first ArrayList growth. The write | |
| 2945 | - | // buffer is userspace-only: it drains via POLLOUT to the client socket, | |
| 2946 | - | // which has no corresponding kernel-imposed per-write limit. | |
| 2947 | - | client.write_buf = try std.ArrayList(u8).initCapacity(client.alloc, 65536); | |
| 2948 | - | try daemon.clients.append(daemon.alloc, client); | |
| 2949 | - | std.log.info( | |
| 2950 | - | "client connected fd={d} total={d}", | |
| 2951 | - | .{ client_fd, daemon.clients.items.len }, | |
| 2952 | - | ); | |
| 2953 | - | } | |
| 2954 | - | ||
| 2955 | - | const inp_flags = lib_posix.POLL.IN | lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL; | |
| 2956 | - | if (poll_fds.items[1].revents & inp_flags != 0) { | |
| 2957 | - | // Read from PTY. Buffer is sized to N_TTY_BUF_SIZE (4096): the hard | |
| 2958 | - | // kernel limit for the N_TTY line discipline. A larger buffer doesn't | |
| 2959 | - | // help: each read() from a PTY master returns at most 4096 bytes | |
| 2960 | - | // regardless of the userspace buffer size. | |
| 2961 | - | var buf: [4096]u8 = undefined; | |
| 2962 | - | const n_opt: ?usize = lib_posix.read(pty_fd, &buf) catch |err| blk: { | |
| 2963 | - | if (err == error.WouldBlock) break :blk null; | |
| 2964 | - | break :blk 0; | |
| 2965 | - | }; | |
| 2966 | - | ||
| 2967 | - | if (n_opt) |n| { | |
| 2968 | - | if (n == 0) { | |
| 2969 | - | // EOF: Shell exited | |
| 2970 | - | std.log.info("shell exited pty_fd={d}", .{pty_fd}); | |
| 2971 | - | // Let the rest of this poll iteration complete so client | |
| 2972 | - | // write buffers are flushed via the normal POLLOUT path. | |
| 2973 | - | // On the next iteration, daemon.running will be false. | |
| 2974 | - | daemon.running = false; | |
| 2975 | - | } else { | |
| 2976 | - | // Feed PTY output to terminal emulator for state tracking | |
| 2977 | - | vt_stream.nextSlice(buf[0..n]); | |
| 2978 | - | daemon.has_pty_output = true; | |
| 2979 | - | ||
| 2980 | - | // When no real terminal client has attached yet, respond to | |
| 2981 | - | // terminal queries (e.g. DA1/DA2) on behalf of the terminal. | |
| 2982 | - | // This prevents fish from waiting 10s for unanswered queries. | |
| 2983 | - | // `has_terminal_client` is only set when a client sends .Init | |
| 2984 | - | // (a real zmx attach), not when a `zmx run` tail-only client | |
| 2985 | - | // connects. | |
| 2986 | - | if (!daemon.has_terminal_client and | |
| 2987 | - | daemon.pty_write_buf.items.len < Daemon.PTY_WRITE_BUF_MAX) | |
| 2988 | - | { | |
| 2989 | - | util.respondToDeviceAttributes(daemon.alloc, &daemon.pty_write_buf, buf[0..n]); | |
| 2990 | - | } | |
| 2991 | - | ||
| 2992 | - | // In run mode, scan output for exit code marker. The marker | |
| 2993 | - | // can straddle two PTY reads (more likely under a throttled | |
| 2994 | - | // scheduler, e.g. containers), so prepend the tail carried | |
| 2995 | - | // over from the previous read before searching. | |
| 2996 | - | if (daemon.is_task_mode and daemon.task_exit_code == null) { | |
| 2997 | - | var scan_buf: [marker_carry.len + buf.len]u8 = undefined; | |
| 2998 | - | @memcpy(scan_buf[0..marker_carry_len], marker_carry[0..marker_carry_len]); | |
| 2999 | - | @memcpy(scan_buf[marker_carry_len..][0..n], buf[0..n]); | |
| 3000 | - | const scan_len = marker_carry_len + n; | |
| 3001 | - | ||
| 3002 | - | if (util.findTaskExitMarker(scan_buf[0..scan_len])) |exit_code| { | |
| 3003 | - | daemon.task_exit_code = exit_code; | |
| 3004 | - | daemon.task_ended_at = @intCast(std.Io.Timestamp.now(daemon.io, .real).nanoseconds); | |
| 3005 | - | ||
| 3006 | - | std.log.info("task completed exit_code={d}", .{exit_code}); | |
| 3007 | - | ||
| 3008 | - | // Notify connected clients | |
| 3009 | - | for (daemon.clients.items) |c| { | |
| 3010 | - | ipc.appendMessage(daemon.alloc, &c.write_buf, .TaskComplete, &[_]u8{exit_code}) catch {}; | |
| 3011 | - | c.has_pending_output = true; | |
| 3012 | - | } | |
| 3013 | - | } | |
| 3014 | - | ||
| 3015 | - | marker_carry_len = @min(marker_carry.len, scan_len); | |
| 3016 | - | @memcpy( | |
| 3017 | - | marker_carry[0..marker_carry_len], | |
| 3018 | - | scan_buf[scan_len - marker_carry_len .. scan_len], | |
| 3019 | - | ); | |
| 3020 | - | } | |
| 3021 | - | ||
| 3022 | - | // Broadcast data to all clients. | |
| 3023 | - | // Rewrite OSC 133;A to include redraw=0 so the outer terminal | |
| 3024 | - | // does not clear prompt lines on resize (issue #111). | |
| 3025 | - | const broadcast_data = util.rewritePromptRedraw(daemon.alloc, buf[0..n]) orelse buf[0..n]; | |
| 3026 | - | defer if (broadcast_data.ptr != buf[0..n].ptr) daemon.alloc.free(broadcast_data); | |
| 3027 | - | for (daemon.clients.items) |client| { | |
| 3028 | - | ipc.appendMessage(daemon.alloc, &client.write_buf, .Output, broadcast_data) catch |err| { | |
| 3029 | - | std.log.warn( | |
| 3030 | - | "failed to buffer output for client err={s}", | |
| 3031 | - | .{@errorName(err)}, | |
| 3032 | - | ); | |
| 3033 | - | continue; | |
| 3034 | - | }; | |
| 3035 | - | client.has_pending_output = true; | |
| 3036 | - | } | |
| 3037 | - | } | |
| 3038 | - | } | |
| 3039 | - | } | |
| 3040 | - | ||
| 3041 | - | if (poll_fds.items[1].revents & lib_posix.POLL.OUT != 0) { | |
| 3042 | - | while (daemon.pty_write_buf.items.len > 0) { | |
| 3043 | - | const n = lib_posix.write(pty_fd, daemon.pty_write_buf.items) catch |err| { | |
| 3044 | - | if (err != error.WouldBlock) { | |
| 3045 | - | std.log.warn("pty write failed: {s}", .{@errorName(err)}); | |
| 3046 | - | daemon.pty_write_buf.clearRetainingCapacity(); | |
| 3047 | - | } | |
| 3048 | - | break; | |
| 3049 | - | }; | |
| 3050 | - | if (n == 0) break; | |
| 3051 | - | daemon.pty_write_buf.replaceRange(daemon.alloc, 0, n, &[_]u8{}) catch unreachable; | |
| 3052 | - | } | |
| 3053 | - | } | |
| 3054 | - | ||
| 3055 | - | var i: usize = daemon.clients.items.len; | |
| 3056 | - | // Only iterate over clients that were present when poll_fds was constructed | |
| 3057 | - | // poll_fds contains [server, pty, sig_pipe, client0, client1, ...] | |
| 3058 | - | // So number of clients in poll_fds is poll_fds.items.len - 3 | |
| 3059 | - | const num_polled_clients = poll_fds.items.len - 3; | |
| 3060 | - | if (i > num_polled_clients) { | |
| 3061 | - | // If we have more clients than polled (i.e. we just accepted one), start from the | |
| 3062 | - | // polled ones | |
| 3063 | - | i = num_polled_clients; | |
| 3064 | - | } | |
| 3065 | - | ||
| 3066 | - | clients_loop: while (i > 0) { | |
| 3067 | - | i -= 1; | |
| 3068 | - | const client = daemon.clients.items[i]; | |
| 3069 | - | const revents = poll_fds.items[i + 3].revents; | |
| 3070 | - | ||
| 3071 | - | if (revents & lib_posix.POLL.IN != 0) { | |
| 3072 | - | const n = client.read_buf.read(client.socket_fd) catch |err| { | |
| 3073 | - | if (err == error.WouldBlock) continue; | |
| 3074 | - | std.log.debug( | |
| 3075 | - | "client read err={s} fd={d}", | |
| 3076 | - | .{ @errorName(err), client.socket_fd }, | |
| 3077 | - | ); | |
| 3078 | - | const last = daemon.closeClient(client, i, false); | |
| 3079 | - | if (last) break :daemon_loop; | |
| 3080 | - | continue; | |
| 3081 | - | }; | |
| 3082 | - | ||
| 3083 | - | if (n == 0) { | |
| 3084 | - | // Client closed connection | |
| 3085 | - | const last = daemon.closeClient(client, i, false); | |
| 3086 | - | if (last) break :daemon_loop; | |
| 3087 | - | continue; | |
| 3088 | - | } | |
| 3089 | - | ||
| 3090 | - | while (client.read_buf.next()) |msg| { | |
| 3091 | - | switch (msg.header.tag) { | |
| 3092 | - | .Input => try daemon.handleInput(client, msg.payload), | |
| 3093 | - | .Send => daemon.handleSend(msg.payload), | |
| 3094 | - | .Output => try daemon.handleOutput(msg.payload, &vt_stream), | |
| 3095 | - | .Init => try daemon.handleInit(client, pty_fd, &term, msg.payload), | |
| 3096 | - | .Switch => try daemon.handleSwitch(msg.payload), | |
| 3097 | - | .Resize => try daemon.handleResize(client, pty_fd, &term, msg.payload), | |
| 3098 | - | .Detach => { | |
| 3099 | - | daemon.handleDetach(client, i); | |
| 3100 | - | break :clients_loop; | |
| 3101 | - | }, | |
| 3102 | - | .DetachAll => { | |
| 3103 | - | daemon.handleDetachAll(); | |
| 3104 | - | break :clients_loop; | |
| 3105 | - | }, | |
| 3106 | - | .Kill => { | |
| 3107 | - | break :daemon_loop; | |
| 3108 | - | }, | |
| 3109 | - | .Info => try daemon.handleInfo(client), | |
| 3110 | - | .LabelGet => try daemon.handleLabelGet(client), | |
| 3111 | - | .LabelSet => try daemon.handleLabelSet(client, msg.payload), | |
| 3112 | - | .LabelClear => try daemon.handleLabelClear(client), | |
| 3113 | - | .History => try daemon.handleHistory(client, &term, msg.payload), | |
| 3114 | - | .Run => try daemon.handleRun(client, msg.payload), | |
| 3115 | - | .Ack, .TaskComplete, .LabelData => {}, | |
| 3116 | - | .Write => try daemon.handleWrite(client, msg.payload), | |
| 3117 | - | _ => std.log.warn( | |
| 3118 | - | "ignoring unknown IPC tag={d}", | |
| 3119 | - | .{@intFromEnum(msg.header.tag)}, | |
| 3120 | - | ), | |
| 3121 | - | } | |
| 3122 | - | } | |
| 3123 | - | } | |
| 3124 | - | ||
| 3125 | - | if (revents & lib_posix.POLL.OUT != 0) { | |
| 3126 | - | // Flush pending output buffers | |
| 3127 | - | const n = lib_posix.write(client.socket_fd, client.write_buf.items) catch |err| blk: { | |
| 3128 | - | if (err == error.WouldBlock) break :blk 0; | |
| 3129 | - | // Error on write, close client | |
| 3130 | - | const last = daemon.closeClient(client, i, false); | |
| 3131 | - | if (last) break :daemon_loop; | |
| 3132 | - | continue; | |
| 3133 | - | }; | |
| 3134 | - | ||
| 3135 | - | if (n > 0) { | |
| 3136 | - | client.write_buf.replaceRange(daemon.alloc, 0, n, &[_]u8{}) catch unreachable; | |
| 3137 | - | } | |
| 3138 | - | ||
| 3139 | - | if (client.write_buf.items.len == 0) { | |
| 3140 | - | client.has_pending_output = false; | |
| 3141 | - | } | |
| 3142 | - | } | |
| 3143 | - | ||
| 3144 | - | if (revents & (lib_posix.POLL.HUP | lib_posix.POLL.ERR | lib_posix.POLL.NVAL) != 0) { | |
| 3145 | - | const last = daemon.closeClient(client, i, false); | |
| 3146 | - | if (last) break :daemon_loop; | |
| 3147 | - | } | |
| 3148 | - | } | |
| 3149 | - | } | |
| 3150 | - | } | |
| 3151 | - | ||
| 3152 | - | fn wakeSignalPipe(_: std.os.linux.SIG, _: *const lib_posix.siginfo_t, _: ?*anyopaque) callconv(.c) void { | |
| 3153 | - | const saved = std.c._errno().*; | |
| 3154 | - | _ = std.c.write(sig_pipe[1], "x", 1); | |
| 3155 | - | std.c._errno().* = saved; | |
| 3156 | - | } | |
| 3157 | - | ||
| 3158 | - | // std.posix.poll retries EINTR internally, so SA_RESTART is moot -- neither | |
| 3159 | - | // setting wakes the loop. The handler writes to sig_pipe instead; poll() | |
| 3160 | - | // wakes on its read end. | |
| 3161 | - | fn installWakeHandler(sig: u6) void { | |
| 3162 | - | const act: lib_posix.Sigaction = .{ | |
| 3163 | - | .handler = .{ .sigaction = wakeSignalPipe }, | |
| 3164 | - | .mask = lib_posix.sigemptyset(), | |
| 3165 | - | .flags = lib_posix.SA.SIGINFO, | |
| 3166 | - | }; | |
| 3167 | - | lib_posix.sigaction(@as(lib_posix.SIG, @enumFromInt(sig)), &act, null); | |
| 3168 | - | } | |
| 3169 | - | ||
| 3170 | - | fn ignoreSigpipe() void { | |
| 3171 | - | const act: lib_posix.Sigaction = .{ | |
| 3172 | - | .handler = .{ .handler = lib_posix.SIG.IGN }, | |
| 3173 | - | .mask = lib_posix.sigemptyset(), | |
| 3174 | - | .flags = 0, | |
| 3175 | - | }; | |
| 3176 | - | lib_posix.sigaction(lib_posix.SIG.PIPE, &act, null); | |
| 3177 | - | } |
+4,
-1
| ... | ... | @@ -30,8 +30,8 @@ const pid_t = system.pid_t; | |
| 30 | 30 | const lfs64_abi = native_os == .linux and builtin.link_libc and (builtin.abi.isGnu() or builtin.abi.isAndroid()); | |
| 31 | 31 | const uid_t = system.uid_t; | |
| 32 | 32 | const mode_t = system.mode_t; | |
| 33 | - | const socket_t = fd_t; | |
| 34 | 33 | const FD_CLOEXEC = system.FD_CLOEXEC; | |
| 34 | + | pub const socket_t = fd_t; | |
| 35 | 35 | pub const SA = system.SA; | |
| 36 | 36 | pub const fd_t = system.fd_t; | |
| 37 | 37 | pub const O = system.O; |
| ... | ... | @@ -52,6 +52,9 @@ pub const Sigaction = system.Sigaction; | |
| 52 | 52 | pub const SIG = system.SIG; | |
| 53 | 53 | pub const siginfo_t = system.siginfo_t; | |
| 54 | 54 | ||
| 55 | + | // https://github.com/ziglang/zig/blob/738d2be9d6b6ef3ff3559130c05159ef53336224/lib/std/posix.zig#L3505 | |
| 56 | + | pub const O_NONBLOCK: usize = 1 << @bitOffsetOf(O, "NONBLOCK"); | |
| 57 | + | ||
| 55 | 58 | pub fn getuid() uid_t { | |
| 56 | 59 | return system.getuid(); | |
| 57 | 60 | } |
+46,
-0
| ... | ... | @@ -0,0 +1,46 @@ | |
| 1 | + | const std = @import("std"); | |
| 2 | + | const lib_posix = @import("posix.zig"); | |
| 3 | + | ||
| 4 | + | /// Self-pipe woken by signal handlers. std.posix.poll loops on .INTR internally | |
| 5 | + | /// (PollError has no Interrupted member), so a signal that lands during poll() | |
| 6 | + | /// never surfaces; the handler writes a byte here and poll() wakes on POLLIN. | |
| 7 | + | pub var sig_pipe: [2]lib_posix.fd_t = .{ -1, -1 }; | |
| 8 | + | ||
| 9 | + | pub fn wakeSignalPipe(_: std.os.linux.SIG, _: *const lib_posix.siginfo_t, _: ?*anyopaque) callconv(.c) void { | |
| 10 | + | const saved = std.c._errno().*; | |
| 11 | + | _ = std.c.write(sig_pipe[1], "x", 1); | |
| 12 | + | std.c._errno().* = saved; | |
| 13 | + | } | |
| 14 | + | ||
| 15 | + | // std.posix.poll retries EINTR internally, so SA_RESTART is moot -- neither | |
| 16 | + | // setting wakes the loop. The handler writes to sig_pipe instead; poll() | |
| 17 | + | // wakes on its read end. | |
| 18 | + | pub fn installWakeHandler(sig: u6) void { | |
| 19 | + | const act: lib_posix.Sigaction = .{ | |
| 20 | + | .handler = .{ .sigaction = wakeSignalPipe }, | |
| 21 | + | .mask = lib_posix.sigemptyset(), | |
| 22 | + | .flags = lib_posix.SA.SIGINFO, | |
| 23 | + | }; | |
| 24 | + | lib_posix.sigaction(@as(lib_posix.SIG, @enumFromInt(sig)), &act, null); | |
| 25 | + | } | |
| 26 | + | ||
| 27 | + | pub fn ignoreSigpipe() void { | |
| 28 | + | const act: lib_posix.Sigaction = .{ | |
| 29 | + | .handler = .{ .handler = lib_posix.SIG.IGN }, | |
| 30 | + | .mask = lib_posix.sigemptyset(), | |
| 31 | + | .flags = 0, | |
| 32 | + | }; | |
| 33 | + | lib_posix.sigaction(lib_posix.SIG.PIPE, &act, null); | |
| 34 | + | } | |
| 35 | + | ||
| 36 | + | pub fn openSignalPipe() !void { | |
| 37 | + | sig_pipe = try lib_posix.pipe2(.{ .CLOEXEC = true, .NONBLOCK = true }); | |
| 38 | + | } | |
| 39 | + | ||
| 40 | + | pub fn drainSignalPipe() void { | |
| 41 | + | var b: [16]u8 = undefined; | |
| 42 | + | while (true) { | |
| 43 | + | const n = lib_posix.read(sig_pipe[0], &b) catch return; | |
| 44 | + | if (n == 0) return; | |
| 45 | + | } | |
| 46 | + | } |
+39,
-1
| ... | ... | @@ -28,6 +28,44 @@ pub fn getSeshName(alloc: std.mem.Allocator, sesh: []const u8) ![]const u8 { | |
| 28 | 28 | return full; | |
| 29 | 29 | } | |
| 30 | 30 | ||
| 31 | + | pub fn resolveSessionOrEnv(alloc: std.mem.Allocator, io: std.Io, session_name: ?[]const u8) ![]const u8 { | |
| 32 | + | const sesh_env = getSeshNameFromEnv(); | |
| 33 | + | const raw = if (session_name) |name| | |
| 34 | + | if (std.mem.eql(u8, name, ".")) blk: { | |
| 35 | + | if (sesh_env.len > 0) break :blk sesh_env; | |
| 36 | + | var buf: [4096]u8 = undefined; | |
| 37 | + | var w = std.Io.File.stderr().writer(io, &buf); | |
| 38 | + | w.interface.print("error: \".\" requires ZMX_SESSION (are you inside a zmx session?)\n", .{}) catch {}; | |
| 39 | + | w.interface.flush() catch {}; | |
| 40 | + | return error.SessionNameRequired; | |
| 41 | + | } else name | |
| 42 | + | else if (sesh_env.len > 0) | |
| 43 | + | sesh_env | |
| 44 | + | else { | |
| 45 | + | return error.SessionNameRequired; | |
| 46 | + | }; | |
| 47 | + | return getSeshName(alloc, raw); | |
| 48 | + | } | |
| 49 | + | ||
| 50 | + | pub const SessionMatch = struct { | |
| 51 | + | name: []const u8, | |
| 52 | + | is_prefix: bool, | |
| 53 | + | ||
| 54 | + | pub fn matches(self: SessionMatch, session_name: []const u8) bool { | |
| 55 | + | if (self.is_prefix) return std.mem.startsWith(u8, session_name, self.name); | |
| 56 | + | return std.mem.eql(u8, session_name, self.name); | |
| 57 | + | } | |
| 58 | + | }; | |
| 59 | + | ||
| 60 | + | pub fn parseSessionArg(alloc: std.mem.Allocator, raw: []const u8) !SessionMatch { | |
| 61 | + | if (raw.len > 0 and raw[raw.len - 1] == '*') { | |
| 62 | + | const name = try getSeshName(alloc, raw[0 .. raw.len - 1]); | |
| 63 | + | return .{ .name = name, .is_prefix = true }; | |
| 64 | + | } | |
| 65 | + | const name = try getSeshName(alloc, raw); | |
| 66 | + | return .{ .name = name, .is_prefix = false }; | |
| 67 | + | } | |
| 68 | + | ||
| 31 | 69 | pub fn sessionConnect(sesh: []const u8) !i32 { | |
| 32 | 70 | var unix_addr = try lib_posix.initUnix(sesh); | |
| 33 | 71 | const socket_fd = try lib_posix.socket(lib_posix.AF.UNIX, lib_posix.SOCK.STREAM | lib_posix.SOCK.CLOEXEC, 0); |
| ... | ... | @@ -56,7 +94,7 @@ pub fn sessionExists(io: std.Io, dir: std.Io.Dir, name: []const u8) !bool { | |
| 56 | 94 | return true; | |
| 57 | 95 | } | |
| 58 | 96 | ||
| 59 | - | pub fn createSocket(sesh: []const u8) !i32 { | |
| 97 | + | pub fn createSocket(sesh: []const u8) !lib_posix.socket_t { | |
| 60 | 98 | // AF.UNIX: Unix domain socket for local IPC with client processes | |
| 61 | 99 | // SOCK.STREAM: Reliable, bidirectional communication | |
| 62 | 100 | // SOCK.NONBLOCK: Set socket to non-blocking |
+4,
-0