const filesystem = @import("../fs.zig"); const std = @import("std"); const builtin = @import("builtin"); const libc = std.c; const posix = std.posix; const pardes = @import("../pardes.zig"); const dump = @import("../dump.zig"); const wire = @import("wire.zig"); const host_io = @import("../host_io.zig"); const selection_pipe = @import("../selection_pipe.zig"); const file_watch = @import("../file_watch.zig"); const shader_build = @import("../shader_build.zig"); const ninep_io = @import("../9p_io.zig"); const look = @import("../look.zig"); const message = pardes.Messages.Message; const TIOCSWINSZ: c_int = @bitCast(@as(u32, if (@hasDecl(posix.T, "IOCSWINSZ")) posix.T.IOCSWINSZ else 0x80087467)); const log = std.log.scoped(.detached); pub const darwin = ninep_io.darwin; const supported = ninep_io.supported; const sun_path_len = ninep_io.sun_path_len; pub const max_clients = 32; const out_backlog = 1 << 20; const read_chunk = 16 * 1024; const pty_chunk = 64 * 1024; const poll_slots = 2 + max_clients + pardes.MAX_PANES + 2 + @as(usize, @intFromBool(ninep_io.quic_enabled)); const reload_retries = 4; const session_backlog = 2 * @as(usize, wire.max_payload); const pty_backlog = 1 << 20; pub const greet_deadline_default_ms: u32 = 5_000; const idle_retain = read_chunk; const accept_pause_ms = 100; const Closed = enum { bye, peer, protocol, backlog, silent, write, read, oom, refused, quitting }; const Client = struct { fd: c_int = -1, attached: bool = false, greet: bool = false, cols: u16 = 0, rows: u16 = 0, in: std.ArrayListUnmanaged(u8) = .empty, out: std.ArrayListUnmanaged(u8) = .empty, mirror: std.ArrayListUnmanaged(pardes.Cell) = .empty, need_full: bool = true, accepted_ms: i64 = 0, /// The post chain last sent (Session.sendPost); null: none yet. post_sent: ?u64 = null, }; const Pty = struct { fd: c_int = -1, pid: posix.pid_t = 0, kill_at: i64 = 0, out: std.ArrayListUnmanaged(u8) = .empty, /// A command pane's child, watched to its exit (host_io.watchExit); /// the pty stays open, unpolled after its end of file, until both its /// end of file and its exit. cmd: host_io.CommandWatch = .{}, }; /// A watched child exited: wake the loop, which takes it (`takeExits`). fn wakeForExit(ctx: ?*anyopaque) void { const box: *Mailbox = @ptrCast(@alignCast(ctx.?)); box.signal(); } const RetiredShell = struct { pid: posix.pid_t = 0, kill_at: i64 = 0 }; const Completion = union(enum) { lsp: struct { id: u32, rows: ?[]u8 }, pipe: selection_pipe.Response, fn deinit(completion: *Completion, gpa: std.mem.Allocator) void { switch (completion.*) { .lsp => |result| if (result.rows) |rows| gpa.free(rows), .pipe => |*result| result.deinit(gpa), } } }; const Mailbox = struct { const capacity = selection_pipe.Tasks.capacity + 1; const Batch = struct { items: [capacity]Completion = undefined, len: usize = 0, status: ?[]u8 = null, }; mutex: std.atomic.Mutex = .unlocked, batch: Batch = .{}, wake: [2]c_int = .{ -1, -1 }, fn post(box: *Mailbox, completion: Completion) void { while (!box.mutex.tryLock()) std.atomic.spinLoopHint(); // Tasks retain their slot until the owner consumes their one completion. std.debug.assert(box.batch.len < capacity); box.batch.items[box.batch.len] = completion; box.batch.len += 1; box.mutex.unlock(); box.signal(); } fn signal(box: *Mailbox) void { while (libc.send(box.wake[1], "w", 1, nosignal) < 0) { if (libc.errno(-1) != .INTR) break; } } fn take(box: *Mailbox) Batch { var bytes: [128]u8 = undefined; while (libc.recv(box.wake[0], &bytes, bytes.len, 0) > 0) {} while (!box.mutex.tryLock()) std.atomic.spinLoopHint(); defer box.mutex.unlock(); const batch = box.batch; box.batch.len = 0; box.batch.status = null; return batch; } }; const Source = union(enum) { listener, completion, client: u8, pty: u8, inotify, ninep_quic, }; pub const Session = struct { gpa: std.mem.Allocator, worker_gpa: std.mem.Allocator, io: std.Io, core: *pardes.Pardes, /// PARDES_TEST_CLOCK's virtual time (host_io.testClock). test_clock: ?u64 = null, listener: c_int = -1, path_buf: [sun_path_len]u8 = undefined, path_len: usize = 0, clients: [max_clients]Client = @splat(.{}), cols: u16, rows: u16, origin: ?u8 = null, scratch: std.ArrayListUnmanaged(u8) = .empty, ptys: [pardes.MAX_PANES]Pty = @splat(.{}), retired_shells: [pardes.MAX_PANES]RetiredShell = @splat(.{}), mailbox: Mailbox = .{}, lsp_task: ?host_io.Lsp.Task = null, pipe_tasks: selection_pipe.Tasks = .{}, prompt_rcs: host_io.Shell.PromptFiles = .{}, inotify_fd: c_int = -1, ninep: ?*ninep_io.Listener = null, ninep_pending: bool = false, /// A 9P request changed the core since the wake byte was last drained: /// the next poll returns at once, to draw it. ninep_wake: std.atomic.Value(bool) = .init(false), watches: file_watch.Table = @splat(null), /// The post chain's files, compiled here for the GUIs attached. shaders: shader_build = .{}, check_files: bool = false, in_loop: bool = false, greet_deadline_ms: u32 = greet_deadline_default_ms, accept_paused_ms: i64 = 0, pub fn initAsync(s: *Session) !void { var pair: [2]c_int = undefined; if (libc.socketpair(libc.AF.UNIX, libc.SOCK.STREAM, 0, &pair) != 0) return error.SocketFailed; errdefer for (pair) |fd| { _ = libc.close(fd); }; for (pair) |fd| { if (libc.fcntl(fd, libc.F.SETFD, @as(c_int, 1)) < 0) return error.SocketOptionFailed; const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0)); if (flags < 0) return error.SocketOptionFailed; var options: libc.O = @bitCast(@as(u32, @bitCast(flags))); options.NONBLOCK = true; if (libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(options))))) < 0) return error.SocketOptionFailed; if (comptime darwin) { const on: c_int = 1; if (libc.setsockopt(fd, libc.SOL.SOCKET, libc.SO.NOSIGPIPE, &on, @sizeOf(c_int)) != 0) return error.SocketOptionFailed; } } s.mailbox.wake = pair; pardes.lsp.setStatusSink(s, lspStatus); host_io.setExitWake(&s.mailbox, wakeForExit); } fn cancelWorkers(s: *Session) void { if (s.lsp_task) |*task| { task.future.cancel(s.io) catch {}; s.lsp_task = null; } s.pipe_tasks.cancelAll(s.io); _ = s.drainCompletions(false); } fn drainCompletions(s: *Session, apply_results: bool) bool { var batch = s.mailbox.take(); for (batch.items[0..batch.len]) |*completion| { defer completion.deinit(s.worker_gpa); if (!apply_results) continue; switch (completion.*) { .lsp => |result| { s.core.update(.{ .lsp_resp = .{ .id = result.id, .rows = result.rows } }); if (s.lsp_task) |*task| if (task.id == result.id) { task.future.await(s.io) catch {}; s.lsp_task = null; }; }, .pipe => |result| { s.core.update(.{ .pipe_resp = .{ .id = result.id, .success = result.success, .outputs = result.outputs, .failure = result.failure, } }); s.pipe_tasks.finish(s.io, result.id); }, } } if (batch.status) |status| { defer s.worker_gpa.free(status); if (apply_results) { var buf: [256]u8 = undefined; s.core.setStatus(s.core.active, message.stamp(&buf, "lsp", status)); } } return batch.len != 0 or batch.status != null; } pub fn restore(s: *Session, bytes: []const u8, from: []const u8) !void { const replacement = try dump.restore(s.core, bytes, from); s.cancelWorkers(); for (0..s.ptys.len) |pane| s.closePty(@intCast(pane)); s.harvest(); for (0..pardes.MAX_PANES) |pane| file_watch.watchPane(s.inotify_fd, &s.watches, @intCast(pane), null, 0, .{ .text = 0 }); _ = file_watch.applyThemeEffect(s.core, s.gpa, s.inotify_fd, &s.watches, 0, false, false); // The new core's chain generations are its own. s.shaders.generation = null; s.check_files = false; if (s.ninep) |listener| listener.reset(replacement); s.ninep_pending = false; replacement.host = s.host(); s.core.deinit(); s.core = replacement; for (&s.clients) |*client| if (client.attached) { client.need_full = true; }; s.mailbox.signal(); } pub fn deinit(s: *Session) void { // First: its compile thread wakes the mailbox, closed below. s.shaders.deinit(s.gpa, s.inotify_fd, &s.watches); if (s.mailbox.wake[0] >= 0) pardes.lsp.setStatusSink(null, null); if (s.mailbox.wake[0] >= 0) host_io.setExitWake(null, null); s.cancelWorkers(); for (s.mailbox.wake) |fd| if (fd >= 0) { _ = libc.close(fd); }; s.mailbox.wake = .{ -1, -1 }; if (s.ninep) |l| { l.deinit(s.gpa); s.ninep = null; } for (&s.clients) |*c| if (c.attached) s.send(c, .quit); for (&s.clients) |*c| if (c.fd >= 0) s.close(c, .quitting); s.unlisten(); s.scratch.deinit(s.gpa); for (0..s.ptys.len) |pane| s.closePty(@intCast(pane)); for (s.retired_shells) |shell| if (shell.pid > 0) { _ = libc.kill(shell.pid, libc.SIG.KILL); }; for (s.ptys) |pty| if (pty.pid > 0) { _ = libc.kill(pty.pid, libc.SIG.KILL); }; for (&s.retired_shells) |*shell| if (shell.pid > 0) { while (libc.waitpid(shell.pid, null, 0) < 0 and libc.errno(-1) == .INTR) {} shell.* = .{}; }; for (&s.ptys) |*pty| if (pty.pid > 0) { while (libc.waitpid(pty.pid, null, 0) < 0 and libc.errno(-1) == .INTR) {} pty.pid = 0; }; if (s.inotify_fd >= 0) { _ = libc.close(s.inotify_fd); s.inotify_fd = -1; } s.prompt_rcs.deinit(); } pub fn listen(s: *Session, name: []const u8) bool { if (comptime !supported) return false; var dir_buf: [sun_path_len:0]u8 = undefined; const dir = ninep_io.socketDir(&dir_buf) orelse return false; if (!ninep_io.ensureSocketDir(dir)) return false; sweep(dir); const path = socketPath(&s.path_buf, dir, name) orelse return false; var addr: libc.sockaddr.un = .{ .path = @splat(0) }; @memcpy(addr.path[0 .. path.len + 1], path[0 .. path.len + 1]); const fd = libc.socket(libc.AF.UNIX, libc.SOCK.STREAM, 0); if (fd < 0) return false; ninep_io.setCloexec(fd); if (libc.bind(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) != 0) { if (alive(path)) { log.debug("a detached session is already listening on {s}", .{path}); _ = libc.close(fd); return false; } _ = libc.unlink(path); if (libc.bind(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) != 0) { _ = libc.close(fd); return false; } } _ = libc.chmod(path, 0o600); if (libc.listen(fd, max_clients) != 0) { _ = libc.close(fd); return false; } setNonblock(fd); s.listener = fd; s.path_len = path.len; return true; } fn unlisten(s: *Session) void { if (s.listener < 0) return; _ = libc.close(s.listener); s.listener = -1; var z: [sun_path_len:0]u8 = undefined; @memcpy(z[0..s.path_len], s.path_buf[0..s.path_len]); z[s.path_len] = 0; _ = libc.unlink(z[0..s.path_len :0]); } pub fn host(s: *Session) host_io.Host { return .{ .ctx = s, .vtable = &vtable }; } fn of(ctx: ?*anyopaque) *Session { return @ptrCast(@alignCast(ctx.?)); } fn clockNow(ctx: ?*anyopaque) u64 { const s = of(ctx); return s.test_clock orelse host_io.monotonicNs(); } const vtable: host_io.Host.VTable = .{ .wait_input = waitInput, .now = clockNow, .present = present, .poll_frame = pollFrame, .spawn = spawn, .pty_write = ptyWrite, .pty_resize = ptyResize, .pty_signal = ptySignal, .close_pty = closePaneShell, .tty_taken = ttyTaken, .fg_name = fgName, .kill_job = killJob, .write_file = writeFile, .write_dump = writeDump, .watch_file = watchFile, .watch_theme = watchTheme, .dump_themes = dumpThemes, .set_clipboard = setClipboard, .read_clipboard = readClipboard, .open_link = openLink, .detach = detach, .lsp = lspRequest, .pipe = pipeRequest, }; fn lspRequest(ctx: ?*anyopaque, request: host_io.Lsp.Request) void { const s = of(ctx); _ = s.drainCompletions(true); if (s.lsp_task) |*task| { task.future.cancel(s.io) catch {}; s.lsp_task = null; _ = s.drainCompletions(true); } const job = host_io.Lsp.snapshot(s.worker_gpa, s.core, request) catch |err| { s.core.update(.{ .lsp_resp = .{ .id = request.id, .rows = null } }); return s.core.reportError(request.pane, "lsp", err); }; const future = s.io.concurrent(lspWorker, .{ s.worker_gpa, job, &s.mailbox }) catch |err| { job.free(s.worker_gpa); s.core.update(.{ .lsp_resp = .{ .id = request.id, .rows = null } }); return s.core.reportError(request.pane, "lsp", err); }; s.lsp_task = .{ .id = request.id, .future = future }; } fn lspWorker(gpa: std.mem.Allocator, job: *host_io.Lsp.Job, box: *Mailbox) anyerror!void { host_io.Lsp.work(gpa, job, box, deliverLspRows); } fn deliverLspRows(ctx: ?*anyopaque, id: u32, rows: ?[]u8) void { const box: *Mailbox = @ptrCast(@alignCast(ctx.?)); box.post(.{ .lsp = .{ .id = id, .rows = rows } }); } fn lspStatus(ctx: ?*anyopaque, text: []const u8) void { const s = of(ctx); const copy = s.worker_gpa.dupe(u8, text) catch return; while (!s.mailbox.mutex.tryLock()) std.atomic.spinLoopHint(); if (s.mailbox.batch.status) |old| s.worker_gpa.free(old); s.mailbox.batch.status = copy; s.mailbox.mutex.unlock(); s.mailbox.signal(); } fn pipeRequest(ctx: ?*anyopaque, id: u32) void { const s = of(ctx); _ = s.drainCompletions(true); if (s.pipe_tasks.full()) { s.core.update(.{ .pipe_resp = .{ .id = id, .success = false, .outputs = &.{} } }); return; } const request = s.core.pipe.pipeRequest(id) orelse return; const job = selection_pipe.Job.copy(s.worker_gpa, request) catch |err| { s.core.update(.{ .pipe_resp = .{ .id = id, .success = false, .outputs = &.{} } }); return s.core.reportError(s.core.active, "pipe", err); }; const future = s.io.concurrent(pipeWorker, .{ s.io, s.worker_gpa, job, &s.mailbox }) catch |err| { job.deinit(s.worker_gpa); s.core.update(.{ .pipe_resp = .{ .id = id, .success = false, .outputs = &.{} } }); return s.core.reportError(s.core.active, "pipe", err); }; std.debug.assert(s.pipe_tasks.add(.{ .id = id, .future = future })); } fn pipeWorker(io: std.Io, gpa: std.mem.Allocator, job: *selection_pipe.Job, box: *Mailbox) anyerror!void { defer job.deinit(gpa); box.post(.{ .pipe = selection_pipe.runJob(gpa, io, job) }); } fn primary(s: *Session) ?*Client { for (&s.clients) |*c| if (c.attached) return c; return null; } fn origins(s: *Session) ?*Client { if (s.origin) |i| { const c = &s.clients[i]; if (c.attached) return c; } return s.primary(); } fn slotOf(s: *Session, c: *const Client) u8 { return @intCast(@divExact(@intFromPtr(c) - @intFromPtr(&s.clients[0]), @sizeOf(Client))); } fn broadcast(s: *Session, msg: wire.ServerMsg) void { for (&s.clients) |*c| if (c.attached) s.send(c, msg); } fn spawn(ctx: ?*anyopaque, pane: u8, cwd: []const u8) void { const s = of(ctx); if (pane >= s.ptys.len) return; // the core indexes its own panes s.closePty(pane); s.harvest(); if (s.ptys[pane].pid != 0) return s.core.reportError(pane, "shell", error.ShellClosing); for (s.retired_shells) |shell| { if (shell.pid == 0) break; } else return s.core.reportError(pane, "shell", error.ShellClosing); const child = host_io.forkShell( s.core, pane, &s.prompt_rcs, s.core.shellBin(), cwd, s.core.screen_h, s.core.screen_w, s.ninep, ) catch |err| return s.core.reportError(pane, "shell", err); const command = if (s.core.panes[pane]) |pn| pn.command != null else false; s.ptys[pane] = .{ .fd = child.file.handle, .pid = child.pid, .cmd = .{ .watched = command and host_io.watchExit(child.pid), .fd = child.file.handle } }; setNonblock(child.file.handle); var lbuf: [pardes.memory.limits.host_path_cap + 1]u8 = undefined; if (host_io.shellCwd(child.pid, &lbuf)) |wd| s.core.setCwd(pane, wd); } fn ptyWrite(ctx: ?*anyopaque, pane: u8, bytes: []const u8) void { const s = of(ctx); if (pane >= s.ptys.len) return; const pt = &s.ptys[pane]; if (pt.fd < 0) return; if (pt.out.items.len > pty_backlog) { var mbuf: [256]u8 = undefined; const text = std.fmt.bufPrint( &mbuf, "input refused: pane not reading ({d} bytes queued)", .{pt.out.items.len}, ) catch "input refused: pane not reading"; return s.core.setMessage(pane, text); } pt.out.appendSlice(s.gpa, bytes) catch { return s.core.setMessage(pane, "input refused: out of memory"); }; s.flushPty(pane); } fn ptyResize(ctx: ?*anyopaque, pane: u8, cols: u16, rows: u16) void { const s = of(ctx); if (pane >= s.ptys.len) return; const fd = s.ptys[pane].fd; if (fd < 0) return; const ws: posix.winsize = .{ .row = rows, .col = cols, .xpixel = 0, .ypixel = 0 }; _ = posix.system.ioctl(fd, TIOCSWINSZ, @intFromPtr(&ws)); } fn ptySignal(ctx: ?*anyopaque, pane: u8, sig: pardes.PtySignal) void { const s = of(ctx); if (pane >= s.ptys.len) return; const pt = s.ptys[pane]; // A command whose exit is recorded is reaped: its group may be another's. if (pt.fd < 0 or pt.cmd.exited) return; host_io.signalTty(pt.pid, pt.fd, sig, pt.cmd.watched); } fn ttyTaken(ctx: ?*anyopaque, pane: u8) bool { const s = of(ctx); if (pane >= s.ptys.len) return false; const pt = s.ptys[pane]; if (pt.fd < 0) return false; return host_io.ttyTaken(pt.pid, pt.fd); } fn fgName(ctx: ?*anyopaque, pane: u8, buf: []u8) ?[]const u8 { const s = of(ctx); if (pane >= s.ptys.len) return null; const pt = s.ptys[pane]; if (pt.fd < 0) return null; return host_io.foregroundName(pt.pid, pt.fd, buf); } fn killJob(ctx: ?*anyopaque, pane: u8) bool { const s = of(ctx); if (pane >= s.ptys.len) return false; const pt = s.ptys[pane]; if (pt.fd < 0) return false; return host_io.killJob(pt.pid, pt.fd); } fn writeFile(ctx: ?*anyopaque, pane: u8, path: []const u8, bytes: []const u8) void { const s = of(ctx); filesystem.write(s.core, path, bytes) catch |err| return s.core.saveFailed(pane, path, err); if (s.core.panes[pane]) |pn| if (pn.file) |f| if (std.mem.eql(u8, f.path, path)) { if (s.watches[pane]) |*w| if (w.serial == pn.serial) switch (w.generation) { .text => w.generation = .{ .text = std.hash.Wyhash.hash(0, bytes) }, .pdf => {}, }; }; var mbuf: [256]u8 = undefined; s.core.setMessage(pane, message.stamp(&mbuf, "saved", path)); } fn writeDump(ctx: ?*anyopaque, bytes: []const u8) void { const s = of(ctx); var pbuf: [1024:0]u8 = undefined; const path = pardes.dump.outPath(&pbuf, s.core.settings.dump_dir.get()) orelse return s.core.dumpFailed(s.core.settings.dump_dir.get(), error.NoDumpDirectory); filesystem.write(s.core, path, bytes) catch |err| return s.core.dumpFailed(path, err); s.core.setLastDump(path); } fn watchFile(ctx: ?*anyopaque, pane: u8, _: []const u8, on: bool, mode: pardes.WatchMode) void { const s = of(ctx); if (file_watch.applyEffect(s.core, s.io, s.inotify(), &s.watches, pane, on, mode)) s.check_files = true; } fn watchTheme(ctx: ?*anyopaque, generation: u32, on: bool) void { const s = of(ctx); if (file_watch.applyThemeEffect(s.core, s.gpa, s.inotify(), &s.watches, generation, on, s.in_loop)) s.check_files = true; } fn dumpThemes(ctx: ?*anyopaque, pane: u8) void { const s = of(ctx); const config_dir = s.core.opts.config_dir orelse return; const out_dir = pardes.config.User.dumpThemes(s.io, s.gpa, config_dir, pardes.themes) catch |err| { s.core.reportError(pane, "dump themes", err); return; }; defer s.gpa.free(out_dir); var mbuf: [256]u8 = undefined; s.core.setMessage(pane, message.stamp(&mbuf, "dumped themes", out_dir)); } fn setClipboard(ctx: ?*anyopaque, text: []const u8) void { const s = of(ctx); s.core.fallback.setClipboard(text); s.broadcast(.{ .set_clipboard = text }); } fn readClipboard(ctx: ?*anyopaque) void { const s = of(ctx); if (s.origins()) |c| return s.send(c, .read_clipboard); s.core.update(.{ .paste = s.core.fallback.clipboard.items }); } fn openLink(ctx: ?*anyopaque, url: []const u8) void { const s = of(ctx); if (s.origins()) |c| return s.send(c, .{ .open_link = url }); s.core.fallback.setLink(url); } fn detach(ctx: ?*anyopaque) void { const s = of(ctx); if (s.origins()) |c| s.send(c, .detach); } /// A 9P request changed the core: wake the poll loop to draw it. fn wakeNinep(ctx: ?*anyopaque) void { const s = of(ctx); s.ninep_wake.store(true, .release); s.mailbox.signal(); } fn pollFrame(ctx: ?*anyopaque) void { const s = of(ctx); if (s.ninep) |l| s.ninep_pending = l.tick().pending; for (&s.ptys, 0..) |*pt, pane| { if (pt.fd < 0) continue; var lbuf: [pardes.memory.limits.host_path_cap + 1]u8 = undefined; if (host_io.shellCwd(pt.pid, &lbuf)) |cwd| s.core.setCwd(pane, cwd); } } fn closePaneShell(ctx: ?*anyopaque, pane: u8) void { const s = of(ctx); if (pane < s.ptys.len) s.closePty(pane); } fn closePty(s: *Session, pane: u8) void { const pt = &s.ptys[pane]; if (pt.fd < 0) return; _ = libc.close(pt.fd); pt.fd = -1; pt.out.deinit(s.gpa); pt.out = .empty; if (pt.pid == 0) return; _ = libc.kill(pt.pid, libc.SIG.HUP); pt.kill_at = monotonicMs() + 100; s.harvest(); if (pt.pid == 0) return; for (&s.retired_shells) |*shell| if (shell.pid == 0) { shell.* = .{ .pid = pt.pid, .kill_at = pt.kill_at }; pt.pid = 0; pt.kill_at = 0; return; }; } fn flushPty(s: *Session, pane: u8) void { const pt = &s.ptys[pane]; var off: usize = 0; while (off < pt.out.items.len) { const n = libc.write(pt.fd, pt.out.items.ptr + off, pt.out.items.len - off); if (n < 0) switch (libc.errno(n)) { .INTR => continue, .AGAIN => break, else => { s.readPty(pane); if (s.ptys[pane].fd >= 0) s.paneEof(pane); return; }, }; if (n == 0) break; off += @intCast(n); } if (off == 0) return; if (off == pt.out.items.len) { pt.out.clearRetainingCapacity(); return retire(s.gpa, &pt.out); } std.mem.copyForwards(u8, pt.out.items, pt.out.items[off..]); pt.out.items.len -= off; } fn readPty(s: *Session, pane: u8) void { var buf: [pty_chunk]u8 = undefined; const got = libc.read(s.ptys[pane].fd, &buf, buf.len); if (got == 0) return s.paneEof(pane); if (got < 0) return switch (libc.errno(got)) { .INTR, .AGAIN => {}, else => s.paneEof(pane), }; s.core.update(.{ .output = .{ .pane = pane, .bytes = buf[0..@intCast(got)] } }); } fn paneEof(s: *Session, pane: u8) void { const pt = &s.ptys[pane]; // A command's pty stays open until its child has exited: closing // it would hang up one that runs on without it. // Its exit, if it came first, is told now its output is in. if (pt.cmd.watched) return host_io.commandEof(s.core, &s.ptys, pane, s, closeWatched); // Unwatched, a command's exit is read here, as its end, and a // shell's, for a run waiting on it and the log. const unwatched = if (s.core.panes[pane]) |pn| pn.command != null else false; const status = host_io.exitStatus(pt.pid, 100); if (status != null) pt.pid = 0; // reaped: no shell to retire s.closePty(pane); s.harvest(); if (unwatched or status != null) s.core.update(.{ .exited = .{ .pane = pane, .status = status } }); s.core.update(.{ .eof = .{ .pane = pane } }); } /// The command panes' exits, told once their output is in (host_io). fn takeExits(s: *Session) void { host_io.takeExits(s.core, &s.ptys, s, closeWatched); } fn closeWatched(s: *Session, id: usize) void { s.closePty(@intCast(id)); } fn harvest(s: *Session) void { const now = monotonicMs(); for (&s.retired_shells) |*shell| reapShell(&shell.pid, &shell.kill_at, now); for (&s.ptys) |*pty| if (pty.fd < 0) reapShell(&pty.pid, &pty.kill_at, now); } fn reapShell(pid: *posix.pid_t, kill_at: *i64, now: i64) void { if (pid.* == 0) return; const result = libc.waitpid(pid.*, null, libc.W.NOHANG); if (result > 0 or (result < 0 and libc.errno(result) == .CHILD)) { pid.* = 0; kill_at.* = 0; } else if (kill_at.* != 0 and now >= kill_at.*) { _ = libc.kill(pid.*, libc.SIG.KILL); kill_at.* = 0; } } fn inotify(s: *Session) c_int { if (s.inotify_fd >= 0) return s.inotify_fd; s.inotify_fd = file_watch.init(true); return s.inotify_fd; } fn drainInotify(s: *Session) void { if (file_watch.drain(s.inotify_fd)) s.check_files = true; } fn reloadWatched(s: *Session) void { if (!s.check_files) return; s.check_files = false; s.shaders.recheck(); for (0..reload_retries) |_| { if (!file_watch.reloadChanged(s.core, s.io, s.gpa, &s.watches)) return; } } fn present(ctx: ?*anyopaque, surface: *const pardes.Surface) void { const s = of(ctx); for (&s.clients) |*c| { if (!c.attached) continue; if (c.out.items.len != 0) continue; s.sendFrame(c, surface); } } fn sendFrame(s: *Session, c: *Client, surface: *const pardes.Surface) void { const cells = surface.cells; const want = wire.frameBound(surface.cols, surface.rows) + wire.layersBound(&surface.body_layers, &surface.tag_layers, surface.regionList()); s.scratch.ensureTotalCapacity(s.gpa, want) catch return s.close(c, .oom); const prev: []const pardes.Cell = if (c.need_full or c.mirror.items.len != cells.len) &.{} else c.mirror.items; const cursor: ?wire.Cursor = if (surface.cursor) |cur| .{ .x = cur.x, .y = cur.y, .bar = cur.bar } else null; const bytes = wire.encodeFrameLayers( s.scratch.allocatedSlice()[0..want], surface.cols, surface.rows, cursor, surface.pointer_shape, cells, prev, &surface.body_layers, &surface.tag_layers, surface.regionList(), surface.chrome, ) catch |err| { log.debug("frame {d}x{d} not encodable: {t}", .{ surface.cols, surface.rows, err }); return; }; s.queue(c, bytes); if (c.fd < 0) return; // the queue closed it; the mirror went with it c.mirror.resize(s.gpa, cells.len) catch return s.close(c, .oom); @memcpy(c.mirror.items, cells); c.need_full = false; } fn waitInput(ctx: ?*anyopaque, timeout_ms: u32) void { const s = of(ctx); // Both before the poll: a wake that landed since the last frame has // had its byte drained here and must not be waited for again. const drained = s.drainCompletions(true); const completed = s.ninep_wake.swap(false, .acq_rel) or drained; for (&s.clients) |*c| if (c.fd >= 0) s.flush(c); const regridded = s.reconcile(); // After the greeting reconcile sends: a chain before it is refused. const compiled = s.sendPost(); // Nobody attached, nobody to dismiss a message: the loop expires // them itself, between the editor's steps (expireUnattended). s.core.unattended = s.primary() == null; pardes.Messages.expireUnattended(s.core); const now = monotonicMs(); var fds: [poll_slots]libc.pollfd = undefined; var src: [poll_slots]Source = undefined; var n: usize = 0; if (s.mailbox.wake[0] >= 0) { fds[n] = .{ .fd = s.mailbox.wake[0], .events = poll_in, .revents = 0 }; src[n] = .completion; n += 1; } const watching_listener = s.listener >= 0 and now >= s.accept_paused_ms; if (watching_listener) { fds[n] = .{ .fd = s.listener, .events = poll_in, .revents = 0 }; src[n] = .listener; n += 1; } for (&s.clients, 0..) |*c, i| { if (c.fd < 0) continue; fds[n] = .{ .fd = c.fd, .events = if (c.out.items.len != 0) poll_in | poll_out else poll_in, .revents = 0, }; src[n] = .{ .client = @intCast(i) }; n += 1; } for (&s.ptys, 0..) |*pt, pane| { if (pt.fd < 0 or pt.cmd.eof) continue; fds[n] = .{ .fd = pt.fd, .events = if (pt.out.items.len != 0) poll_in | poll_out else poll_in, .revents = 0, }; src[n] = .{ .pty = @intCast(pane) }; n += 1; } if (s.inotify_fd >= 0) { fds[n] = .{ .fd = s.inotify_fd, .events = poll_in, .revents = 0 }; src[n] = .inotify; n += 1; } if (s.ninep) |l| { // Unix and TCP connections run on cloud9's runner, which wakes // this loop through the mailbox; only QUIC is polled here. if (comptime ninep_io.quic_enabled) { if (l.quic) |*listener| { fds[n] = listener.poll(); src[n] = .ninep_quic; n += 1; } } } // The wait is the 9P connections' turn with the core. if (n == 0) { pardes.turn.rest(); nap(if (timeout_ms == 0) 16 else timeout_ms); pardes.turn.wake(); if (timeout_ms != 0) if (s.test_clock) |*virtual| { virtual.* = @max(virtual.*, s.core.nextWake() orelse virtual.*); }; return; } var timeout: c_int = if (timeout_ms == 0) -1 else @intCast(@min(timeout_ms, std.math.maxInt(c_int))); if (s.nextWake(now)) |due| timeout = if (timeout < 0) due else @min(timeout, due); if (completed or compiled or s.check_files or regridded or s.ninep_pending) timeout = 0; pardes.turn.rest(); const ready = libc.poll(&fds, @intCast(n), timeout); pardes.turn.wake(); if (ready > 0) s.dispatch(fds[0..n], src[0..n]); // The core's wake ran out; the pump advances it from `now`. if (timeout_ms != 0 and monotonicMs() -| now >= timeout_ms) if (s.test_clock) |*virtual| { virtual.* = @max(virtual.*, s.core.nextWake() orelse virtual.*); }; _ = s.drainCompletions(true); s.takeExits(); s.expire(monotonicMs()); s.harvest(); s.reloadWatched(); _ = s.reconcile(); } fn dispatch(s: *Session, fds: []const libc.pollfd, src: []const Source) void { for (fds, src) |pfd, source| switch (source) { .listener => if (pfd.revents != 0) s.accept(), .completion => {}, .client => |i| { const c = &s.clients[i]; if (c.fd < 0) continue; if (pfd.revents & poll_out != 0) s.flush(c); if (c.fd < 0) continue; if (pfd.revents & poll_in != 0) { s.receive(c, i); } else if (pfd.revents & (poll_hup | poll_err | poll_nval) != 0) { s.close(c, .peer); } }, .pty => |pane| { if (s.ptys[pane].fd < 0) continue; if (pfd.revents & poll_out != 0) s.flushPty(pane); if (s.ptys[pane].fd < 0) continue; if (pfd.revents & poll_in != 0) { s.readPty(pane); } else if (pfd.revents & (poll_hup | poll_err | poll_nval) != 0) { s.paneEof(pane); } }, .inotify => if (pfd.revents & poll_in != 0) s.drainInotify(), .ninep_quic => {}, }; } fn nextWake(s: *const Session, now: i64) ?c_int { if (now == 0) return null; // no clock; see `monotonicMs` var due: ?i64 = null; for (s.retired_shells) |shell| if (shell.pid != 0) { due = now + 10; break; }; for (s.ptys) |pty| if (pty.fd < 0 and pty.pid != 0) { due = now + 10; break; }; for (&s.clients) |*c| { if (c.fd < 0 or c.attached) continue; const at = c.accepted_ms + @as(i64, s.greet_deadline_ms); due = if (due) |d| @min(d, at) else at; } if (s.listener >= 0 and s.accept_paused_ms > now) due = if (due) |d| @min(d, s.accept_paused_ms) else s.accept_paused_ms; if (s.ninep) |l| { if (l.nextDue()) |ms| { const at9 = now + ms; due = if (due) |d| @min(d, at9) else at9; } } // Unattended messages go on the clock: look again in a while. if (s.core.unattended) for (s.core.panes) |slot| if (slot) |pane| if (pane.msg_len > 0) { due = if (due) |d| @min(d, now + 100) else now + 100; break; }; const at = due orelse return null; return @intCast(@max(0, @min(at - now, std.math.maxInt(c_int)))); } fn expire(s: *Session, now: i64) void { if (now == 0) return; // no clock: enforce nothing rather than everything for (&s.clients) |*c| { if (c.fd < 0 or c.attached) continue; if (now - c.accepted_ms >= s.greet_deadline_ms) s.close(c, .silent); } } fn accept(s: *Session) void { for (0..max_clients + 1) |_| { const fd = libc.accept(s.listener, null, null); if (fd < 0) { switch (libc.errno(fd)) { .AGAIN, .INTR, .CONNABORTED => return, else => { s.accept_paused_ms = monotonicMs() + accept_pause_ms; return; }, } } ninep_io.setCloexec(fd); setNonblock(fd); if (comptime darwin) { const on: c_int = 1; _ = libc.setsockopt(fd, libc.SOL.SOCKET, libc.SO.NOSIGPIPE, &on, @sizeOf(c_int)); } const slot = for (&s.clients, 0..) |*c, i| { if (c.fd < 0) break i; } else { s.refuseFd(fd, .full); _ = libc.close(fd); continue; }; s.clients[slot] = .{ .fd = fd, .accepted_ms = monotonicMs() }; } } fn receive(s: *Session, c: *Client, slot: u8) void { var buf: [read_chunk]u8 = undefined; const got = libc.read(c.fd, &buf, buf.len); if (got == 0) return s.close(c, .peer); // clean EOF: the frontend left if (got < 0) return switch (libc.errno(got)) { .INTR, .AGAIN => {}, else => s.close(c, .read), }; c.in.appendSlice(s.gpa, buf[0..@intCast(got)]) catch return s.close(c, .oom); s.account(); if (c.fd < 0) return; // it was this one s.consume(c, slot); } fn consume(s: *Session, c: *Client, slot: u8) void { var off: usize = 0; while (true) { const found = wire.framed(c.in.items[off..]) catch return s.close(c, .protocol); const msg = found orelse break; s.apply(c, slot, msg.tag, msg.payload) catch return s.close(c, .protocol); if (c.fd < 0) return; // apply closed it, buffers and all off += msg.total; } if (off == 0) return; if (off == c.in.items.len) { c.in.clearRetainingCapacity(); return retire(s.gpa, &c.in); } std.mem.copyForwards(u8, c.in.items, c.in.items[off..]); c.in.items.len -= off; } fn apply(s: *Session, c: *Client, slot: u8, tag: u8, payload: []const u8) wire.Error!void { if (tag == @intFromEnum(wire.ClientTag.hello)) { if (c.attached) return error.BadValue; const claimed = try wire.helloVersion(payload); if (claimed != wire.version) { log.debug("frontend speaks protocol {d}, this session speaks {d}", .{ claimed, wire.version }); return s.refuse(c, .version); } } switch (try wire.decodeClient(tag, payload)) { .hello => |h| { if (s.core.quit) return s.refuse(c, .quitting); c.cols = h.cols; c.rows = h.rows; c.attached = true; c.need_full = true; c.greet = true; }, .bye => s.close(c, .bye), .event => |ev| { if (!c.attached) return error.BadValue; switch (ev) { .resize => |r| { c.cols = r.cols; c.rows = r.rows; if (r.row_metrics) |metrics| s.core.row_metrics = metrics; }, else => { s.origin = slot; s.core.update(ev); }, } }, } } fn reconcile(s: *Session) bool { var cols: u16 = 0; var rows: u16 = 0; for (&s.clients) |*c| { if (!c.attached) continue; cols = if (cols == 0) c.cols else @min(cols, c.cols); rows = if (rows == 0) c.rows else @min(rows, c.rows); } var regridded = false; // Against the core's size too: a root ctl `size` may have set it // while nobody was attached. if (cols != 0 and (cols != s.cols or rows != s.rows or cols != s.core.screen_w or rows != s.core.screen_h)) { s.cols = cols; s.rows = rows; for (&s.clients) |*c| c.need_full = true; s.core.update(.{ .resize = .{ .cols = cols, .rows = rows } }); regridded = true; } for (&s.clients, 0..) |*c, i| { if (!c.greet) continue; c.greet = false; s.send(c, .{ .welcome = .{ .slot = @intCast(i), .cols = s.cols, .rows = s.rows } }); } return regridded; } /// Compiles the post chain's files (a finished compile wakes the loop /// through the mailbox) and sends each attached frontend the chain it /// has not seen: on attach, and when it or a file's SPIR-V changes. True /// when a compile was taken in, so the loop draws what it said. /// ponytail: a terminal frontend is sent it too and drops it; a hello /// saying which frontend it is would spare those bytes. fn sendPost(s: *Session) bool { const files = for (s.core.settings.post.list()) |entry| { if (entry.scene == null) break true; } else false; const compiled = s.shaders.sync(s.gpa, s.io, s.core, if (files) s.inotify() else s.inotify_fd, &s.watches, .{ .ctx = &s.mailbox, .call = wakeForExit }); const key = @as(u64, s.shaders.revision) << 8 | @intFromEnum(s.core.settings.shader_animation); var chain: wire.Post = .{ .animation = s.core.settings.shader_animation, .len = s.core.settings.post.len }; _ = s.shaders.view(&s.core.settings.post, &chain.passes); for (&s.clients) |*c| { if (!c.attached or c.greet or c.post_sent == key) continue; c.post_sent = key; s.send(c, .{ .post = chain }); } return compiled; } fn send(s: *Session, c: *Client, msg: wire.ServerMsg) void { const want = wire.serverBound(msg); s.scratch.ensureTotalCapacity(s.gpa, want) catch return s.close(c, .oom); const bytes = wire.encodeServer(s.scratch.allocatedSlice()[0..want], msg) catch |err| { log.debug("message {t} not encodable: {t}", .{ msg, err }); return; }; s.queue(c, bytes); } fn queue(s: *Session, c: *Client, bytes: []const u8) void { if (c.out.items.len > out_backlog) return s.close(c, .backlog); s.account(); if (c.fd < 0) return; // the fattest peer was this one c.out.appendSlice(s.gpa, bytes) catch return s.close(c, .oom); s.flush(c); } fn account(s: *Session) void { var total: usize = 0; var worst: ?*Client = null; var worst_bytes: usize = 0; for (&s.clients) |*c| { if (c.fd < 0) continue; const held = c.in.items.len + c.out.items.len; total += held; if (held > worst_bytes) { worst_bytes = held; worst = c; } } if (total <= session_backlog) return; if (worst) |c| s.close(c, .backlog); } fn flush(s: *Session, c: *Client) void { var off: usize = 0; while (off < c.out.items.len) { const n = libc.send(c.fd, c.out.items.ptr + off, c.out.items.len - off, nosignal); if (n < 0) switch (libc.errno(n)) { .INTR => continue, .AGAIN => break, else => return s.close(c, .write), }; if (n == 0) break; off += @intCast(n); } if (off == 0) return; if (off == c.out.items.len) { c.out.clearRetainingCapacity(); return retire(s.gpa, &c.out); } std.mem.copyForwards(u8, c.out.items, c.out.items[off..]); c.out.items.len -= off; } fn refuse(s: *Session, c: *Client, why: wire.Refusal) void { s.refuseFd(c.fd, why); s.close(c, .refused); } fn refuseFd(_: *Session, fd: c_int, why: wire.Refusal) void { var buf: [wire.header_len + 1]u8 = undefined; const bytes = wire.encodeServer(&buf, .{ .refuse = why }) catch unreachable; var off: usize = 0; while (off < bytes.len) { const n = libc.send(fd, bytes.ptr + off, bytes.len - off, nosignal); if (n < 0 and libc.errno(n) == .INTR) continue; if (n <= 0) break; // it left before hearing why; nothing to do off += @intCast(n); } } fn close(s: *Session, c: *Client, why: Closed) void { if (c.fd < 0) return; log.debug("frontend detached: {t}", .{why}); _ = libc.close(c.fd); c.in.deinit(s.gpa); c.out.deinit(s.gpa); c.mirror.deinit(s.gpa); const gone = s.slotOf(c); if (s.origin) |i| if (i == gone) { s.origin = null; }; c.* = .{}; } }; /// SIGTERM, SIGINT or SIGHUP: quit as Exit does, through the loop, so its /// sockets are unlinked on the way out rather than left for the next /// session to trip on. The handler only sets the flag and wakes the loop. var quit_signalled = std.atomic.Value(bool).init(false); var quit_wake: c_int = -1; fn onQuitSignal(_: std.posix.SIG) callconv(.c) void { quit_signalled.store(true, .release); if (quit_wake >= 0) _ = libc.send(quit_wake, "q", 1, nosignal); } /// The screen of a detached session no frontend has sized. pub const detached_cols = 160; pub const detached_rows = 50; pub fn run(init: std.process.Init, opts: pardes.Options, name: []const u8) !void { const gpa = init.gpa; const allocs = pardes.memory.init(gpa); defer pardes.memory.deinit(); var options = opts; // No frontend yet, so no size to take: one roomy enough for a script's // panes (80x24 left pane/new refused after a few); `size` sets another, // and an attaching frontend its own. options.cols = detached_cols; options.rows = detached_rows; options.image_allocator = allocs.image; options.pdf_allocator = allocs.pdf; options.tree_sitter_allocator = allocs.tree_sitter; options.frame_allocator = allocs.frame; pardes.image.start(init.io, allocs.image); if (comptime pardes.pdf_enabled) pardes.pdf.start(allocs.pdf); pardes.syntax.start(allocs.tree_sitter); defer { pardes.image.stop(); if (comptime pardes.pdf_enabled) pardes.pdf.stop(); pardes.syntax.stop(); } const core = if (options.load_path) |lp| blk: { const bytes = try filesystem.readFile(gpa, lp); defer gpa.free(bytes); break :blk try pardes.dump.initFromDump(allocs.pardes, options, bytes); } else try pardes.Pardes.init(allocs.pardes, options); var session: Session = .{ .gpa = gpa, .worker_gpa = allocs.lsp, .io = init.io, .core = core, .cols = options.cols, .rows = options.rows, .prompt_rcs = host_io.Shell.prepare(), .test_clock = host_io.testClock(), }; defer session.core.deinit(); defer session.deinit(); try session.initAsync(); if (!session.listen(name)) { try std.Io.File.stderr().writeStreamingAll(init.io, "pardes: could not bind a detached session socket\n"); return error.NoSocket; } session.ninep = ninep_io.listen(init.io, gpa, core, opts.ninep_name, name, opts.ninep_tcp, opts.ninep_quic); if (session.ninep == null) return error.ListenFailed; session.ninep.?.wake_ctx = &session; session.ninep.?.wake = Session.wakeNinep; const h = session.host(); core.host = h; while (core.nextEffect()) |effect| core.perform(effect); session.in_loop = true; quit_wake = session.mailbox.wake[1]; const on_quit: std.posix.Sigaction = .{ .handler = .{ .handler = onQuitSignal }, .mask = std.posix.sigemptyset(), .flags = 0 }; for ([_]std.posix.SIG{ .HUP, .INT, .TERM }) |sig| std.posix.sigaction(sig, &on_quit, null); defer quit_wake = -1; while (!session.core.quit and !quit_signalled.load(.acquire)) { pardes.turn.restoreSettled(); try session.core.pump(h); if (session.core.quit) break; if (session.core.takeRestore()) |path| restore: { const bytes = filesystem.readRestore(gpa, path, session.core.settings.dump_dir.get()) catch |err| { session.core.reportError(session.core.active, "Restore", err); break :restore; }; defer gpa.free(bytes); session.restore(bytes, path) catch |err| session.core.reportError(session.core.active, "Restore", err); } } } test "detached queued results preserve current requests and are discarded before Restore" { const gpa = std.testing.allocator; var s: Session = .{ .gpa = gpa, .worker_gpa = gpa, .io = std.testing.io, .core = try pardes.Pardes.init(gpa, .{ .tty_only = true, .cols = 40, .rows = 12 }), .cols = 40, .rows = 12, }; defer s.core.deinit(); defer s.deinit(); try s.initAsync(); while (s.core.nextEffect()) |_| {} _ = try s.core.setTestFile("saved body\n"); s.core.lspRequest(0, .status, ""); const old_id = s.core.lsp_wait.?.id; s.core.lspRequest(0, .status, ""); const current_id = s.core.lsp_wait.?.id; s.lsp_task = .{ .id = current_id, .future = .{ .any_future = null, .result = {} } }; s.mailbox.post(.{ .lsp = .{ .id = old_id, .rows = try gpa.dupe(u8, "old result\n") } }); try std.testing.expect(s.drainCompletions(true)); try std.testing.expectEqual(current_id, s.lsp_task.?.id); try std.testing.expectEqual(current_id, s.core.lsp_wait.?.id); try dump.dumpState(s.core); const saved = try gpa.dupe(u8, s.core.dump_out.?); defer gpa.free(saved); while (s.core.nextEffect()) |_| {} s.mailbox.post(.{ .lsp = .{ .id = current_id, .rows = try gpa.dupe(u8, "queued before restore\n") } }); const outputs = try gpa.alloc([]u8, 1); outputs[0] = try gpa.dupe(u8, "old filter output\n"); try std.testing.expect(s.pipe_tasks.add(.{ .id = 77, .future = .{ .any_future = null, .result = {} } })); s.mailbox.post(.{ .pipe = .{ .id = 77, .success = true, .outputs = outputs } }); Session.lspStatus(&s, "old status"); try s.restore(saved, "test"); try std.testing.expect(s.lsp_task == null); try std.testing.expectEqual(@as(usize, 0), s.pipe_tasks.len); try std.testing.expect(!s.drainCompletions(true)); try std.testing.expectEqualStrings("saved body\n", s.core.panes[0].?.file.?.content); s.core.lspRequest(0, .status, ""); const restored_id = s.core.lsp_wait.?.id; try std.testing.expect(restored_id > current_id); s.lsp_task = .{ .id = restored_id, .future = .{ .any_future = null, .result = {} } }; s.mailbox.post(.{ .lsp = .{ .id = current_id, .rows = try gpa.dupe(u8, "late old result\n") } }); _ = s.drainCompletions(true); try std.testing.expectEqual(restored_id, s.lsp_task.?.id); try std.testing.expectEqual(restored_id, s.core.lsp_wait.?.id); s.mailbox.post(.{ .lsp = .{ .id = restored_id, .rows = &.{} } }); const host = s.host(); host.vtable.wait_input.?(host.ctx, 1); try std.testing.expect(s.lsp_task == null); try std.testing.expect(s.core.lsp_wait == null); } test "a detached session's terminal pane takes its shell with it when it closes" { const gpa = std.testing.allocator; var s: Session = .{ .gpa = gpa, .worker_gpa = gpa, .io = std.testing.io, .core = try pardes.Pardes.init(gpa, .{ .tty_only = true, .cols = 40, .rows = 12 }), .cols = 40, .rows = 12, }; defer s.core.deinit(); defer s.deinit(); // A stand-in shell: it waits for a signal, and a hangup ends it. const pid = libc.fork(); if (pid == 0) while (true) { _ = libc.poll(&[0]libc.pollfd{}, 0, -1); }; try std.testing.expect(pid > 0); var fds: [2]c_int = undefined; try std.testing.expectEqual(@as(c_int, 0), libc.pipe(&fds)); defer _ = libc.close(fds[1]); s.ptys[3] = .{ .fd = fds[0], .pid = pid }; const host = s.host(); host.vtable.close_pty.?(host.ctx, 3); try std.testing.expect(s.ptys[3].fd < 0); const deadline = monotonicMs() + 2000; while (monotonicMs() < deadline) { s.harvest(); const running = s.ptys[3].pid != 0 or for (s.retired_shells) |shell| { if (shell.pid == pid) break true; } else false; if (!running) break; _ = libc.poll(&[0]libc.pollfd{}, 0, 5); } try std.testing.expect(s.ptys[3].pid == 0); for (s.retired_shells) |shell| try std.testing.expect(shell.pid != pid); } test "detached worker setup failure completes requests without changing document bytes" { const gpa = std.testing.allocator; var failing = std.testing.FailingAllocator.init(gpa, .{ .fail_index = 0 }); var s: Session = .{ .gpa = gpa, .worker_gpa = failing.allocator(), .io = std.testing.io, .core = try pardes.Pardes.init(gpa, .{ .tty_only = true, .cols = 40, .rows = 12 }), .cols = 40, .rows = 12, }; defer s.core.deinit(); defer s.deinit(); try s.initAsync(); while (s.core.nextEffect()) |_| {} const pane = try s.core.setTestFile("one\n"); s.core.host = s.host(); s.core.lspRequest(0, .status, ""); while (s.core.nextEffect()) |effect| s.core.perform(effect); try std.testing.expect(s.core.lsp_wait == null); try std.testing.expect(s.lsp_task == null); pane.body.cur_col = 2; pane.body.vsel = .{ .active = true, .row = 0, .col = 0, .explicit = true }; s.core.update(.{ .key = .{ .cp = '|' } }); s.core.update(.{ .key = .{ .cp = 't', .text = "tr a-z A-Z" } }); s.core.update(.{ .key = .{ .cp = pardes.Key.enter } }); try std.testing.expect(s.core.pipe.wait != null); while (s.core.nextEffect()) |effect| s.core.perform(effect); try std.testing.expect(s.core.pipe.wait == null); try std.testing.expectEqual(@as(usize, 0), s.pipe_tasks.len); try std.testing.expectEqualStrings("one\n", pane.file.?.content); try std.testing.expect(failing.has_induced_failure); } pub const setCloexec = ninep_io.setCloexec; pub fn setNonblock(fd: c_int) void { const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0)); if (flags < 0) return; var o: libc.O = @bitCast(@as(u32, @bitCast(flags))); o.NONBLOCK = true; _ = libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(o))))); } pub const nosignal: u32 = if (darwin) 0 else libc.MSG.NOSIGNAL; pub const poll_in: i16 = @intCast(libc.POLL.IN); pub const poll_out: i16 = @intCast(libc.POLL.OUT); pub const poll_hup: i16 = @intCast(libc.POLL.HUP); pub const poll_err: i16 = @intCast(libc.POLL.ERR); pub const poll_nval: i16 = @intCast(libc.POLL.NVAL); fn retire(gpa: std.mem.Allocator, list: *std.ArrayListUnmanaged(u8)) void { if (list.items.len != 0 or list.capacity <= idle_retain) return; list.clearAndFree(gpa); } pub fn monotonicMs() i64 { var ts: libc.timespec = undefined; if (libc.clock_gettime(.MONOTONIC, &ts) != 0) return 0; return @as(i64, ts.sec) * std.time.ms_per_s + @divTrunc(ts.nsec, std.time.ns_per_ms); } fn nap(ms: u32) void { var ts: libc.timespec = .{ .sec = @intCast(ms / 1000), .nsec = @intCast((ms % 1000) * std.time.ns_per_ms), }; _ = libc.nanosleep(&ts, null); } pub fn socketPath(buf: *[sun_path_len]u8, dir: []const u8, name: []const u8) ?[:0]const u8 { if (name.len == 0) return null; if (std.mem.indexOfAny(u8, name, "/\x00") != null) return null; return std.fmt.bufPrintSentinel(buf, "{s}/" ++ prefix ++ "{s}.sock", .{ dir, name }, 0) catch null; } const prefix = "pardes-detached-"; pub const path_max = sun_path_len; pub fn sessionPath(buf: *[path_max]u8, name: []const u8) ?[:0]const u8 { if (comptime !supported) return null; var dir_buf: [sun_path_len:0]u8 = undefined; const dir = ninep_io.socketDir(&dir_buf) orelse return null; return socketPath(buf, dir, name); } pub fn vetted(path: [:0]const u8) bool { if (comptime !supported) return false; var dir_buf: [sun_path_len:0]u8 = undefined; const dir = ninep_io.socketDir(&dir_buf) orelse return false; if (!ours(ninep_io.statNoFollow(dir) orelse return false, s_ifdir)) return false; return ours(ninep_io.statNoFollow(path) orelse return false, s_ifsock); } const s_ifmt: u32 = 0o170000; const s_ifdir: u32 = 0o040000; const s_ifsock: u32 = 0o140000; fn ours(st: ninep_io.FileFacts, kind: u32) bool { if (st.mode & s_ifmt != kind) return false; if (st.uid != libc.getuid()) return false; return st.mode & 0o077 == 0; } fn alive(path: [:0]const u8) bool { var addr: libc.sockaddr.un = .{ .path = @splat(0) }; if (path.len + 1 > addr.path.len) return true; @memcpy(addr.path[0 .. path.len + 1], path[0 .. path.len + 1]); const fd = libc.socket(libc.AF.UNIX, libc.SOCK.STREAM, 0); if (fd < 0) return true; defer _ = libc.close(fd); ninep_io.setCloexec(fd); setNonblock(fd); const rc = libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))); if (rc == 0) return true; return libc.errno(rc) != .CONNREFUSED; } fn sweep(dir: [:0]const u8) void { const d = libc.opendir(dir) orelse return; defer _ = libc.closedir(d); while (libc.readdir(d)) |ent| { const name = std.mem.sliceTo(&ent.name, 0); if (!std.mem.startsWith(u8, name, prefix) or !std.mem.endsWith(u8, name, ".sock")) continue; var path_buf: [sun_path_len:0]u8 = undefined; const path = std.fmt.bufPrintSentinel(&path_buf, "{s}/{s}", .{ dir, name }, 0) catch continue; if (!alive(path)) _ = libc.unlink(path); } }