diff options
Diffstat (limited to 'src/detached/server.zig')
| -rw-r--r-- | src/detached/server.zig | 1271 |
1 files changed, 1271 insertions, 0 deletions
diff --git a/src/detached/server.zig b/src/detached/server.zig new file mode 100644 index 00000000..6d13a4ae --- /dev/null +++ b/src/detached/server.zig @@ -0,0 +1,1271 @@ +//! THE DETACHED CORE: one `Pardes` instance in a process with no terminal, +//! serving N frontends over one unix socket. +//! +//! THIS SIDE OWNS THE CORE. `Session` is a `host.Host` implementation whose +//! methods encode wire messages instead of doing IO, and whose +//! `pull_wait_input` is a `poll(2)` over the listener and every attached +//! frontend. The frontends own terminals and nothing else (client.zig). So the +//! `Pardes` is here, `update` is called from here, and the same screen is on +//! every attached frontend at once — `screen -x`, not N sessions. +//! +//! WHAT THIS SIDE SERVES ITSELF. Every method this vtable leaves null falls +//! through to the core's own `host.Fallback`: the embedded source filesystem, +//! the in-process clipboard, silent ptys. host.zig says in as many words that a +//! zero-method host is a complete pardes, and that is exactly what a session +//! with nothing attached is. Everything a real frontend can do BETTER — fork a +//! shell on a real tty, put bytes on a real disk, reach a real desktop +//! clipboard — is asked of a frontend, and the routing table below says which. +//! +//! ROUTING, and it is not "push means broadcast". A push reaches every HOST +//! (host.zig's rule, which `Fanout.isPull` enforces); this is ONE host that +//! happens to be backed by several frontends, and how it spreads a call inside +//! itself is its own business. Three rules, one per kind of side effect: +//! * BROADCAST — the frame, and `set_clipboard`. Every screen must show the +//! same thing, and a yank in a shared session is a session-wide fact that +//! every attached desktop is entitled to. +//! * PRIMARY ONLY — `spawn`, `pty_write`, `pty_resize`, `write_file`, +//! `write_dump`, `watch_file`, `watch_theme`, `dump_themes`. Each of these +//! has ONE real resource behind it, and doing it twice is not doing it +//! twice as well: two frontends forking a shell for pane 3 gives the pane +//! two shells, and two frontends writing one path race each other. Primary +//! is the lowest attached slot, i.e. the oldest surviving attachment — a +//! rule that is stable while frontends come and go and needs no election. +//! A pane's shell therefore lives in the frontend that forked it: when that +//! frontend leaves, its panes stop producing output and the session's text, +//! files and layout carry on. That is a real limit and it is stated here +//! rather than papered over, because migrating a live pty between processes +//! is a different feature. +//! * ORIGIN, ELSE PRIMARY — `read_clipboard` (the one `pull_` on the wire) +//! and `open_link`. Both answer a thing a HUMAN just did, and the answer +//! belongs on that human's machine: the paste must come from the keyboard +//! that asked for it, and a link must open in front of the person who +//! clicked it. `origin` is the frontend whose event was applied most +//! recently. Effects drain after a whole batch of events (pardes.zig +//! `pump`), so in the rare case where two frontends type in the same +//! millisecond the second one wins; the fallback to primary covers an +//! effect that no input caused at all. +//! +//! FAIRNESS, and why no client can stall the core or another client: +//! * every descriptor is non-blocking, and there is no thread per client. One +//! `poll(2)` per pump covers the listener and all `max_clients` frontends. +//! * FRAMES ARE NOT QUEUED. A client with bytes still owed to the kernel is +//! SKIPPED for this frame and its mirror is left alone, so the next frame +//! it does get is a diff against what it actually has. A slow frontend +//! therefore sees fewer, larger frames instead of a growing queue, and +//! coalescing costs no byte surgery at all. +//! * what is left in a client's out-queue is control messages, and it is +//! capped (`out_backlog`). The cap is checked BEFORE an append, so a single +//! oversized message still goes out whole and what gets refused is a client +//! that has stopped draining: it is closed. Its session and its peers are +//! untouched, and it may reattach and be sent a full frame. +//! * `max_clients` is a REFUSAL, not a queue — the same shape and the same +//! number as fuse.zig's park table, and for the same reason: the listener +//! is always accepted from even when the table is full, because a +//! level-triggered `poll` on a backlog nobody accepts returns ready +//! forever and spins a core. Bounded per round all the same (`accept`), and +//! a connection that never says `hello` loses its slot +//! (`greet_deadline_ms`) — a slot held by silence is the same denial as a +//! queue, arrived at from the other end. +//! * the TABLE is accounted, not just each client (`session_backlog`), and a +//! drained client gives its buffers back (`idle_retain`): 32 slots each +//! holding one 4 MiB paste is 128 MiB of a daemon nobody is looking at. +//! +//! THE SOCKET follows nested.zig's conventions exactly, and they ARE +//! nested.zig's: `socketDir`, `ensureSocketDir`, `statNoFollow` and +//! `setCloexec` are imported from it rather than copied, because one directory +//! vetted by two predicates is how the two go out of step. `$XDG_RUNTIME_DIR` +//! else `~/.local/state/pardes` created 0700 and vetted (never /tmp), +//! `chmod 0600` before `listen(2)`, CLOEXEC on the listener and on every +//! accepted connection. The NAME differs on purpose: +//! `pardes-detached-<name>.sock` rather than `pardes-<pid>.sock`, so that +//! nested.zig's sweeper — which only recognises all-digit pids — never unlinks +//! a live detached session, and so that a person can say `--detach=work` +//! instead of learning a pid. +//! +//! WHO MAY BIND A NAME, and this side is not allowed to guess. `bind(2)` on a +//! unix socket is an atomic exclusive create, so it decides: a name whose +//! socket ANSWERS is a live session and `listen` refuses rather than taking it +//! (an unconditional unlink-before-bind is how a second `--detach=work` used +//! to steal the socket out from under every frontend attached to the first). +//! The only file this process unlinks is one it proved dead — a connect that +//! was REFUSED — and `alive` is the single place that judgement is made, for +//! `listen` and for the sweep both. +//! +//! ...and both ends do the vetting. `vetted` is the frontend's half: a socket +//! at a path anyone could plant receives every keystroke that frontend +//! collects, so the client checks the directory and the socket before it +//! connects, exactly as this side checks them before it binds. +const std = @import("std"); +const libc = std.c; +const pardes = @import("../pardes.zig"); +const host_api = @import("../host.zig"); +const wire = @import("wire.zig"); + +/// Diagnostics for whoever is running the daemon. Every one of these is a +/// `debug`, and the level is not a judgement about how bad the thing is: +/// main.zig's logFn drops this scope entirely unless PARDES_LOG is set, so what +/// decides whether a human sees it is that variable and not the level. Reaching +/// for `warn` instead would change exactly one thing — a TEST binary does not +/// go through logFn, and its stderr is the build runner's failure signal. +const log = std.log.scoped(.detached); + +/// nested.zig owns the socket conventions this file shares — the directory, +/// its vetting, the stat that will not follow a symlink, CLOEXEC — and its +/// module comment carries the reasoning for each. Imported and not copied: +/// see the module header. +const nested = @import("../nested.zig"); + +/// `pub` for client.zig, which needs the same platform answer for the same +/// reason: SIGPIPE is per-write on linux and per-socket on darwin. +pub const darwin = nested.darwin; + +/// Same two ingredients as nested.zig needs, minus the ancestor walk: unix +/// sockets and a per-user runtime directory. Anywhere else there is no detached +/// session and `listen` says so. +const supported = nested.supported; + +/// `sun_path` is 108 bytes on linux and 104 on darwin, taken from the struct so +/// that the buffers, the fit checks and the memcpy cannot disagree with the +/// kernel or with each other. +const sun_path_len = nested.sun_path_len; + +/// How many frontends may be attached at once. The number and the shape are +/// fuse.zig's park table: 32 slots, and overflow is a refusal rather than a +/// queue. A session with 32 frontends on it is not a session, it is a mistake, +/// and the 33rd gets told so instead of waiting in a backlog nobody drains. +pub const max_clients = 32; + +/// Bytes of un-drained CONTROL messages a client may owe before it is closed. +/// Frames are not in here (see the module header), so this is a backlog of +/// ACTIONS — spawns, clipboard mirrors, file writes — and a frontend that has +/// not taken 1 MiB of those has stopped reading its socket. Checked before an +/// append rather than after, so one oversized message is never the thing that +/// trips it. +const out_backlog = 1 << 20; + +/// One read per client per poll round (see `receive`). 16 KiB is two orders of +/// magnitude past a keystroke and small enough to sit on the loop's stack; a +/// 4 MiB paste arrives across several rounds, which is the point. +const read_chunk = 16 * 1024; + +/// Bytes of client traffic — every in-queue and out-queue together — this +/// session may hold before it starts closing the peers holding it. +/// `out_backlog` bounds ONE slot and this bounds the table, which is not the +/// same ceiling: 32 clients each a byte under their own cap is 32 MiB of a +/// daemon nobody is looking at. 4 MiB is one whole paste in flight plus every +/// frame queue a real session builds, and past it the fattest peer is the peer +/// that stopped reading. The mirrors are NOT in this number: a mirror is this +/// session's own bookkeeping for a client it chose to serve, not something a +/// peer can grow. +const session_backlog = 4 << 20; + +/// What a DRAINED client is allowed to keep. `in` grows to hold one whole +/// message, so a single 4 MiB paste otherwise leaves 4 MiB resident in that +/// slot for the life of the session — 128 MiB across a full table, for +/// something that happened once. Anything above one `read_chunk` is handed +/// back the moment the buffer empties, and the next message pays one +/// allocation for it; below that it is kept, so a session of keystrokes never +/// asks the allocator at all. +const idle_retain = read_chunk; + +/// How long the listener is left out of the poll set after an `accept` that +/// failed for a reason that persists (EMFILE above all). See `accept`: the +/// alternative was sleeping 100 ms inside the core. +const accept_pause_ms = 100; + +/// Why a client's connection ended. Only ever logged (`PARDES_LOG=1`), and +/// spelled out because "connection closed" is the one diagnostic that has never +/// helped anybody. +const Closed = enum { bye, peer, protocol, backlog, silent, write, read, oom, refused, quitting }; + +const Client = struct { + fd: c_int = -1, + /// The `hello` landed and was accepted. Before that the connection exists + /// but votes on nothing and is sent no frames: its geometry is unknown. + attached: bool = false, + /// A `welcome` is owed, and is sent once this round's geometry has settled + /// so the number in it is the one the next frame will use. + greet: bool = false, + /// This frontend's own window, as its last `hello`/`resize` said. One vote + /// in `reconcile`'s minimum, never the session's grid by itself. + cols: u16 = 0, + rows: u16 = 0, + /// Bytes read and not yet a whole message. + in: std.ArrayListUnmanaged(u8) = .empty, + /// Bytes owed to the kernel. + out: std.ArrayListUnmanaged(u8) = .empty, + /// What this client's grid holds, so the next frame can be a diff. Advanced + /// only when a frame is actually queued for it, which is what makes a + /// skipped frame correct rather than lost. + mirror: std.ArrayListUnmanaged(pardes.Cell) = .empty, + /// The next frame must be full: freshly attached, or the session geometry + /// moved under it. + need_full: bool = true, + /// Monotonic milliseconds at `accept`, and the only thing an un-greeted + /// connection is timed against. See `Session.greet_deadline_ms`. + accepted_ms: i64 = 0, +}; + +/// A pane's shell: which frontend was asked to fork it, and where. +/// +/// WHY THE SESSION REMEMBERS THIS. A spawn is the one primary-only call that +/// has to survive having no frontend to serve it. Every boot layout creates its +/// panes before the socket exists, so a `--detach` performs its startup spawns +/// with nobody attached — and dropping them meant a session that opened with +/// panes whose shells had never been forked, forever, in silence. So a spawn +/// with no primary is OWED, and asked of whoever attaches next. +/// +/// It is also what makes `pty_write` reach the right process. A pane's pty +/// lives in the frontend that forked it, which is not always the primary: A +/// attaches and forks the shells, B attaches, A leaves — the panes are re-owed +/// to B, and then C attaching into A's freed slot becomes primary while the +/// ptys are in B. Routing a pane's bytes by its OWNER rather than by the +/// primary is the difference between typing into a shell and typing into +/// nothing. +const Shell = struct { + /// The slot that was asked to fork this pane's shell, or null when nobody + /// has been. + owner: ?u8 = null, + /// A spawn owed to whoever attaches next: either it was never asked, or the + /// frontend holding it left and took the pty with it. + owed: bool = false, + /// Copied, because `push_spawn`'s `cwd` borrows the core's memory for the + /// length of that one call and this outlives it by definition. + cwd: std.ArrayListUnmanaged(u8) = .empty, +}; + +pub const Session = struct { + gpa: std.mem.Allocator, + core: *pardes.Pardes, + /// -1 when nothing is bound: an unsupported platform, or a bind that + /// failed. A session with no listener is a session nobody can attach to, + /// which still runs. + listener: c_int = -1, + /// The bound path, kept so teardown unlinks exactly what was created and + /// nothing else — guarded on the fd, like nested.zig's `unlisten`. + path_buf: [sun_path_len]u8 = undefined, + path_len: usize = 0, + clients: [max_clients]Client = @splat(.{}), + /// The session grid: the smallest common one across attached frontends. + /// Seeded from the core's own startup size so the first attach of an + /// identically sized frontend posts no resize at all. + cols: u16, + rows: u16, + /// Whose input was applied last, for the two calls that must go back to one + /// particular frontend. See the module header. + origin: ?u8 = null, + /// One encode buffer, reused. Grown to whatever the largest message so far + /// needed rather than sized from `wire.max_payload`, which would be 16 MiB + /// of resident memory for a session whose frames are six kilobytes. + scratch: std.ArrayListUnmanaged(u8) = .empty, + /// Where each pane's shell lives, and which spawns are still owed. See + /// `Shell`. + shells: [pardes.MAX_PANES]Shell = @splat(.{}), + /// How long a connection may stay silent before the session takes its slot + /// back. `Client.open` writes its `hello` in the same call that connects, + /// so a peer that has said nothing for five seconds is not a frontend that + /// was slow, and thirty-two of them used to fill the table and lock every + /// real frontend out with a `refuse .full`. + /// + /// A field rather than a constant for exactly one reason: the test for that + /// would otherwise have to sleep five seconds. Nothing else changes it. + greet_deadline_ms: u32 = 5_000, + /// Monotonic milliseconds until which the LISTENER is left out of the poll + /// set, because an `accept` failed for a reason that persists. See `accept`. + accept_paused_ms: i64 = 0, + + // ---- lifetime --------------------------------------------------------- + + pub fn deinit(s: *Session) void { + // Tell everyone the session is over before the socket disappears, so a + // frontend exits on a `quit` rather than on a read error whose meaning + // it has to guess. Best effort by construction: these descriptors are + // non-blocking, so a frontend that is not reading gets the EOF instead + // — which is a case it has to handle regardless. + 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 (&s.shells) |*sh| sh.cwd.deinit(s.gpa); + } + + /// Bind and listen. False when there is no socket, and a session without + /// one is simply one nobody can attach to — the same posture nested.zig + /// takes, and for the same reason: a failed bind must not cost a launch. + 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 = nested.socketDir(&dir_buf) orelse return false; + if (!nested.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; + nested.setCloexec(fd); + // `bind` IS the exclusive create — it fails with EADDRINUSE the moment + // the path exists — so it, and nothing else, decides who owns a name. + // There is no unlink before it: unlinking unconditionally is how a + // second `pardes --detach=work` took the socket away from a live + // session, leaving every frontend attached to a file no new frontend + // could reach. + if (libc.bind(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) != 0) { + // The one case that is not a collision: a session killed rather + // than quit ran no teardown, so its file outlived it. `alive` is + // the only thing that may say so, and it says so only about a + // connect that was REFUSED. + 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; + } + } + // Owner-only, and BEFORE listen(2), which is the first moment anyone + // could connect. The directory is already private; this is the second + // wall, and this socket carries keystrokes into a live editor. + _ = libc.chmod(path, 0o600); + // A backlog of max_clients: past that the kernel refuses the connect + // itself, which is the same answer `accept` would give. + 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; + // Guarded on the fd, so a bind that FAILED cannot unlink a path this + // process never created. + 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_api.Host { + return .{ .ctx = s, .vtable = &vtable }; + } + + fn of(ctx: ?*anyopaque) *Session { + return @ptrCast(@alignCast(ctx.?)); + } + + /// Thirteen methods, and the seven that are missing are missing on purpose + /// — see wire.zig's header for each one's reason. `push_poll_frame` and + /// `push_post_present` carry no information a frame does not; the four + /// synchronous or dispatched pulls and `push_fs_reply` belong to whoever + /// owns the core, which is this process. + const vtable: host_api.Host.VTable = .{ + .pull_wait_input = waitInput, + .push_present = present, + .push_spawn = spawn, + .push_pty_write = ptyWrite, + .push_pty_resize = ptyResize, + .push_write_file = writeFile, + .push_write_dump = writeDump, + .push_watch_file = watchFile, + .push_watch_theme = watchTheme, + .push_dump_themes = dumpThemes, + .push_set_clipboard = setClipboard, + .pull_read_clipboard = readClipboard, + .push_open_link = openLink, + }; + + // ---- routing ---------------------------------------------------------- + + /// The oldest surviving attachment. No election and no state: slots are + /// filled lowest-first, so the lowest attached one is the oldest that is + /// still here. + fn primary(s: *Session) ?*Client { + for (&s.clients) |*c| if (c.attached) return c; + return null; + } + + /// ...and the frontend whose input we are answering, when there is one. + fn origins(s: *Session) ?*Client { + if (s.origin) |i| { + const c = &s.clients[i]; + if (c.attached) return c; + } + return s.primary(); + } + + /// The frontend holding pane `pane`'s pty, which is NOT the primary in + /// general — see `Shell`. Null when no frontend holds it, which is a pane + /// with no child: the core's own answer to that is silence, and so is this. + fn holder(s: *Session, pane: u8) ?*Client { + if (pane >= s.shells.len) return null; + const owner = s.shells[pane].owner orelse return null; + const c = &s.clients[owner]; + return if (c.attached) c else null; + } + + /// Which slot this client is. From the pointer because every caller here + /// holds a `*Client` and not its index. + 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); + } + + // ---- the host methods ------------------------------------------------- + + fn spawn(ctx: ?*anyopaque, pane: u8, cwd: []const u8) void { + const s = of(ctx); + if (pane >= s.shells.len) return; // the core indexes its own panes + const sh = &s.shells[pane]; + // Kept whether or not there is somebody to ask, because the pane now + // exists either way and the cwd is the only thing that cannot be + // reconstructed later. + sh.cwd.clearRetainingCapacity(); + sh.cwd.appendSlice(s.gpa, cwd) catch {}; + if (s.primary()) |c| { + sh.owner = s.slotOf(c); + sh.owed = false; + return s.send(c, .{ .spawn = .{ .pane = pane, .cwd = cwd } }); + } + // Nobody can fork a shell right now — a startup layout, or every + // frontend gone. NOT dropped: `reconcile` asks the next arrival. + sh.owner = null; + sh.owed = true; + } + + fn ptyWrite(ctx: ?*anyopaque, pane: u8, bytes: []const u8) void { + const s = of(ctx); + // To the frontend that forked this pane's shell, not to the primary: + // the pty is in that process and nowhere else. A pane whose holder is + // gone is silent, which is what the core does with a null method. + if (s.holder(pane)) |c| s.send(c, .{ .pty_write = .{ .pane = pane, .bytes = bytes } }); + } + + fn ptyResize(ctx: ?*anyopaque, pane: u8, cols: u16, rows: u16) void { + const s = of(ctx); + if (s.holder(pane)) |c| s.send(c, .{ .pty_resize = .{ .pane = pane, .cols = cols, .rows = rows } }); + } + + fn writeFile(ctx: ?*anyopaque, pane: u8, path: []const u8, bytes: []const u8) void { + const s = of(ctx); + if (s.primary()) |c| return s.send(c, .{ .write_file = .{ .pane = pane, .path = path, .bytes = bytes } }); + // Nobody attached, and a Put must not evaporate. This method being + // non-null means the core did NOT reach for its own filesystem, so the + // obligation a null method would have discharged is discharged here by + // hand — the same shape `readClipboard` below has, and host.zig's rule + // that a zero-method host is a complete pardes. + s.core.fallback.writeFile(path, bytes); + } + + fn writeDump(ctx: ?*anyopaque, bytes: []const u8) void { + const s = of(ctx); + if (s.primary()) |c| return s.send(c, .{ .write_dump = bytes }); + // ...and the same for a Dump, including the part that makes the bytes + // reachable again: a real host reports where it landed, which is what + // puts `Restore <path>` in the topbar (pardes.zig `write_dump`). + s.core.fallback.writeFile(pardes.fallback_dump_path, bytes); + s.core.setLastDump(pardes.fallback_dump_path); + } + + fn watchFile(ctx: ?*anyopaque, pane: u8, path: []const u8, on: bool) void { + const s = of(ctx); + if (s.primary()) |c| return s.send(c, .{ .watch_file = .{ .pane = pane, .path = path, .on = on } }); + // No frontend to watch a path, so the core's own record of what was + // asked is the whole of what a watch means here — exactly what a null + // method leaves behind. + if (pane < s.core.fallback.watched.len) s.core.fallback.watched[pane] = on; + } + + fn watchTheme(ctx: ?*anyopaque, generation: u32, on: bool) void { + const s = of(ctx); + // Dropped with nobody attached, and that is the whole of it: the core's + // own `theme_file` effect does nothing for a null method either, so + // there is no obligation left over. Same for `dump_themes` below. + if (s.primary()) |c| s.send(c, .{ .watch_theme = .{ .generation = generation, .on = on } }); + } + + fn dumpThemes(ctx: ?*anyopaque, pane: u8) void { + const s = of(ctx); + if (s.primary()) |c| s.send(c, .{ .dump_themes = .{ .pane = pane } }); + } + + fn setClipboard(ctx: ?*anyopaque, text: []const u8) void { + const s = of(ctx); + // Mirrored into the core's own clipboard ALWAYS, not only when nobody + // is attached: `readClipboard` answers from it when there is no + // frontend, and a frontend can leave between the yank and the paste. A + // yank that a detached session then pasted as the previous yank is the + // bug this one line is. + s.core.fallback.setClipboard(text); + s.broadcast(.{ .set_clipboard = text }); + } + + /// The one `pull_` that crosses the wire, and it stays a pull for exactly + /// the reason host.zig gives: two frontends answering would paste the + /// clipboard twice for one Ctrl-V. + fn readClipboard(ctx: ?*anyopaque) void { + const s = of(ctx); + if (s.origins()) |c| return s.send(c, .read_clipboard); + // Nobody attached. This method being non-null means the core will NOT + // reach for its own fallback, so an unanswered request would leave + // `clip_pending` armed forever — host.zig's note that a null method + // answers immediately is the obligation being met here by hand. + 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 }); + // No desktop in reach, so the link goes where a host with no browser + // puts it: the core's record of the last one asked for, which is what + // `Fallback.setLink` is and what the acme filesystem reads back. + s.core.fallback.setLink(url); + } + + // ---- the frame -------------------------------------------------------- + + fn present(ctx: ?*anyopaque, surface: *const pardes.Surface) void { + const s = of(ctx); + for (&s.clients) |*c| { + if (!c.attached) continue; + // A client that has not drained what it already owes does not get + // this frame, and its mirror is deliberately left where it is: the + // next frame it does get is a diff against what it really has. A + // slow frontend gets fewer, larger frames rather than a queue. + 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); + s.scratch.ensureTotalCapacity(s.gpa, want) catch return s.close(c, .oom); + // Nothing comparable on the far side is the LATE JOINER and the RESIZE + // in one test: either way the whole grid has to be described. + 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.encodeFrame( + s.scratch.allocatedSlice()[0..want], + surface.cols, + surface.rows, + cursor, + cells, + prev, + ) catch |err| { + // A frame this protocol cannot carry is a grid past `max_cols` / + // `max_rows`, or a cursor the core placed outside its own surface. + // Dropping the frame keeps the session alive with a stale screen, + // which is strictly better than dropping the frontend — a frontend + // REFUSES such a frame and hangs up — and the log says which. + 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 + // The mirror advances only now, and only because the bytes are on the + // wire or in the kernel's buffer for it. + c.mirror.resize(s.gpa, cells.len) catch return s.close(c, .oom); + @memcpy(c.mirror.items, cells); + c.need_full = false; + } + + // ---- the loop --------------------------------------------------------- + + /// The only place this process sleeps, which is what `pull_wait_input`'s + /// comment in host.zig requires of whoever serves it. One `poll(2)` covers + /// the listener and every attached frontend; there is no thread per client + /// and nothing here blocks on a single peer. + fn waitInput(ctx: ?*anyopaque, timeout_ms: u32) void { + const s = of(ctx); + // Push what the kernel will take before sleeping: a client that becomes + // writable while we are inside poll(2) would otherwise be a frame late, + // and a frame late is a frame skipped (see `present`). + for (&s.clients) |*c| if (c.fd >= 0) s.flush(c); + + const now = monotonicMs(); + var fds: [max_clients + 1]libc.pollfd = undefined; + var slots: [max_clients + 1]u8 = undefined; + var n: usize = 0; + // The listener is left OUT of the set while accepting is paused, which + // is how an EMFILE is waited out without the core sleeping (see + // `accept`). Every frontend already attached goes on being served. + 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 }; + slots[n] = 0; + 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, + }; + slots[n] = @intCast(i); + n += 1; + } + // A detached session with no listener and no clients has no event + // source at all. Returning immediately would spin the outer + // `while (!core.quit)` at full speed, so sleep the interval the core + // offered and, when it offered none, a frame's worth. + if (n == 0) return nap(if (timeout_ms == 0) 16 else timeout_ms); + // Zero is the core's word for "sleep until something happens" (see + // pardes.zig `pump`: it passes a frame interval only while an animation + // is running). poll spells that -1. + var timeout: c_int = if (timeout_ms == 0) -1 else @intCast(@min(timeout_ms, std.math.maxInt(c_int))); + // Two things here are due on a CLOCK rather than on a descriptor: a + // handshake that has to expire, and a paused listener that has to come + // back. An indefinite poll would sit through both — and thirty-two + // peers that connect and then say nothing, with the session otherwise + // idle, IS the denial `greet_deadline_ms` exists to answer — so the + // wait is clamped to whichever is due first. + if (s.nextWake(now)) |due| timeout = if (timeout < 0) due else @min(timeout, due); + const ready = libc.poll(&fds, @intCast(n), timeout); + // Expired unconditionally: a slot held by silence comes back on a + // timeout exactly as it does on a wakeup, and a poll that returned + // nothing is the ordinary way this deadline is reached. + s.expire(monotonicMs()); + // A timeout is an ordinary frame boundary and EINTR is a signal we do + // not handle here; both simply come back next pump. + if (ready <= 0) return; + + var k: usize = 0; + if (watching_listener) { + if (fds[0].revents != 0) s.accept(); + k = 1; + } + while (k < n) : (k += 1) { + const c = &s.clients[slots[k]]; + // A slot closed earlier in this same pass (its peer hung up, a + // decode failed, its handshake expired) must not be touched + // through a stale revents. + if (c.fd < 0) continue; + if (fds[k].revents & poll_out != 0) s.flush(c); + if (c.fd < 0) continue; + if (fds[k].revents & poll_in != 0) { + s.receive(c, slots[k]); + } else if (fds[k].revents & (poll_hup | poll_err | poll_nval) != 0) { + // POLLIN wins when both are set: a peer that wrote and then + // closed has bytes still worth reading. + s.close(c, .peer); + } + } + s.reconcile(); + } + + /// Milliseconds until the next deadline that is kept by the CLOCK rather + /// than by a descriptor, or null when there is none. Floored at zero, so a + /// deadline already past polls once without blocking instead of blocking + /// forever on a negative timeout. + fn nextWake(s: *const Session, now: i64) ?c_int { + if (now == 0) return null; // no clock; see `monotonicMs` + var due: ?i64 = null; + 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; + const at = due orelse return null; + return @intCast(@max(0, @min(at - now, std.math.maxInt(c_int)))); + } + + /// Take the slots of connections that never said `hello` back. A connection + /// that holds a slot in silence denies a real frontend exactly as a queue + /// would, and `Client.open` writes its hello in the same call that + /// connects, so there is nothing legitimate to wait for. + 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); + } + } + + /// Always accept, even with a full table: the tempting alternative — stop + /// accepting and let the kernel hold the surplus — is a spin, because + /// `poll` is level triggered and an unaccepted backlog reports ready + /// forever. fuse.zig's park table learned that as a deadlock; here it is + /// 100% of a core. + /// + /// BOUNDED all the same. `max_clients + 1` is enough to fill an empty table + /// and refuse one more, and past that the surplus waits in the backlog for + /// the next round — one pump later, with every frontend drawn in between. + /// The `while (true)` this replaces let a peer dialling in a loop hold the + /// core inside `accept` for as long as it kept dialling, and the core is + /// what draws every other frontend's screen. + 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)) { + // The backlog is empty, which is this loop's ordinary exit. + .AGAIN, .INTR, .CONNABORTED => return, + // Anything else — EMFILE above all — persists until some + // other descriptor is freed, and `poll` is LEVEL + // triggered: coming straight back means poll reports the + // listener ready again immediately and the core spins at + // 100% until the condition clears. The old answer was a + // 100 ms nanosleep, which parks the CORE — every attached + // frontend stops being drawn for a tenth of a second + // because a descriptor ran out. So the LISTENER is dropped + // from the poll set for that beat instead, and the session + // goes on serving the frontends it has. + else => { + s.accept_paused_ms = monotonicMs() + accept_pause_ms; + return; + }, + } + } + nested.setCloexec(fd); + setNonblock(fd); + if (comptime darwin) { + // linux says MSG_NOSIGNAL per write; darwin says it once per + // socket. Either way a frontend that dies mid-frame must not + // take the session down with SIGPIPE. + 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 { + // Refused, and told why, on a connection accepted purely so + // that the listener stays quiet. + s.refuseFd(fd, .full); + _ = libc.close(fd); + continue; + }; + s.clients[slot] = .{ .fd = fd, .accepted_ms = monotonicMs() }; + } + } + + /// One read per client per round. A frontend that never stops talking gets + /// one turn and then the loop moves on to the others and to the frame — + /// which is fuse.zig's `retry` rule (one attempt per parked request per + /// frame) applied to sockets. + 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); + // The table's own ceiling, checked where the table grows: a peer that + // sends the first half of a 16 MiB message and stops is holding memory + // no per-message check can see. See `session_backlog`. + 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; + // The decoded Event BORROWS these bytes, so the buffer is not + // compacted until every message already in it has been applied — + // the same borrow window the tty host gives a pty chunk. + 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 { + var scratch: wire.Scratch = .{}; + switch (try wire.decodeClient(tag, payload, &scratch)) { + .hello => |h| { + // A second hello on one connection is not a resize; it is a + // peer that is not speaking this protocol. + if (c.attached) return error.BadValue; + if (h.version != wire.version) { + log.debug("frontend speaks protocol {d}, this session speaks {d}", .{ h.version, wire.version }); + return s.refuse(c, .version); + } + if (s.core.quit) return s.refuse(c, .quitting); + c.cols = h.cols; + c.rows = h.rows; + c.attached = true; + c.need_full = true; + // Greeted after `reconcile`, so the geometry in the welcome is + // the one this client's first frame will actually use. + c.greet = true; + }, + .bye => s.close(c, .bye), + .event => |ev| { + // Input before a handshake has no geometry behind it and no + // version agreement either. + if (!c.attached) return error.BadValue; + switch (ev) { + // A frontend's resize is about ITS window. The core only + // ever sees the smallest common grid, which `reconcile` + // posts once per round when it moves — forwarding this raw + // would let whichever frontend resized last win. + .resize => |r| { + c.cols = r.cols; + c.rows = r.rows; + }, + else => { + s.origin = slot; + s.core.update(ev); + }, + } + }, + } + } + + /// Settle the session grid and greet whoever arrived, once per poll round + /// rather than once per message: three frontends attaching in the same + /// round are one resize, not three reflows of every pane. + fn reconcile(s: *Session) void { + 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); + } + // Nobody attached: keep the grid we had. A detached session is not a + // session of no size, it is one nobody is looking at, and reflowing + // every pane to nothing for zero readers is work with no reader. + if (cols != 0 and (cols != s.cols or rows != s.rows)) { + s.cols = cols; + s.rows = rows; + // Every mirror is now the wrong shape. `encodeFrame` reaches the + // same conclusion from the cell count alone, but saying it here is + // what makes a reshape with the SAME cell count (80x24 -> 48x40) + // safe too. + for (&s.clients) |*c| c.need_full = true; + s.core.update(.{ .resize = .{ .cols = cols, .rows = rows } }); + } + 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 } }); + } + s.flushOwed(); + } + + /// Hand every owed spawn to the frontend that can serve it. Runs at the end + /// of a poll round, so a frontend that has just been greeted is asked for + /// its panes' shells in the same round it arrived — and a session that was + /// started with panes and no frontend (which is every `--detach`) is a + /// session whose panes get their shells from the first attach rather than + /// never. See `Shell`. + fn flushOwed(s: *Session) void { + const c = s.primary() orelse return; + const slot = s.slotOf(c); + for (&s.shells, 0..) |*sh, pane| { + if (!sh.owed) continue; + sh.owed = false; + sh.owner = slot; + s.send(c, .{ .spawn = .{ .pane = @intCast(pane), .cwd = sh.cwd.items } }); + // The send closed it, and `close` put its panes back on the owed + // list; the ones this loop has not reached are still owed anyway. + if (c.fd < 0) return; + } + } + + // ---- bytes ------------------------------------------------------------ + + 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| { + // The only reachable case is a payload past `max_payload`: a save + // of a pane holding more text than this protocol carries. The + // session keeps it (the core's own filesystem already has it) and + // the frontend's copy does not happen — said out loud rather than + // silently. + log.debug("message {t} not encodable: {t}", .{ msg, err }); + return; + }; + s.queue(c, bytes); + } + + fn queue(s: *Session, c: *Client, bytes: []const u8) void { + // BEFORE the append, so one oversized message always goes out whole and + // what this refuses is a client that has stopped draining. + if (c.out.items.len > out_backlog) return s.close(c, .backlog); + // ...and the table as a whole, which `out_backlog` does not bound: 32 + // slots one byte under it each is 32 MiB. See `session_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); + // Try immediately: on a local socket this empties the queue in one + // write, and `present` skips a client whose queue is not empty. + s.flush(c); + } + + /// Close the peer holding the most of the table when the table as a whole + /// is over `session_backlog`. One peer per call, and the fattest one, + /// because this is only ever asked when the total is already over and the + /// peer holding the most of it is the peer that stopped reading. The next + /// append asks again, so a second offender is closed a message later rather + /// than in a loop that could empty the table on one bad frame. + 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.capacity + c.out.capacity; + 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, + // The kernel's buffer is full: the rest waits for POLLOUT, and + // this client is skipped for frames until it drains. + .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; + } + + /// Say why, then hang up. The refusal is written with a plain blocking + /// write on a socket nobody has sent anything on yet: it is six bytes, and + /// queueing it would mean keeping a slot for a connection being rejected. + fn refuse(s: *Session, c: *Client, why: wire.Refusal) void { + s.refuseFd(c.fd, why); + s.close(c, .refused); + } + + /// Writes only. The descriptor belongs to the caller — `refuse` hands it to + /// `close`, and the full-table path in `accept` closes it itself — because + /// closing here as well is a double close, and the number is reusable the + /// instant the first one lands. + 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); + } + } + + /// Free one slot. A frontend dying takes NOTHING with it: not the core, not + /// the listener, not another frontend's frames. Its buffers go back and the + /// slot is reusable on the next connect. + 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); + // Which slot this is, so a departing frontend cannot leave `origin` + // pointing at it and send the next `read_clipboard` to a stranger. + if (s.origin) |i| if (i == gone) { + s.origin = null; + }; + // ...and its panes' shells died with the process that forked them. They + // go back on the owed list, so the frontend that replaces this one is + // asked to fork them again in the directory they were forked in: the + // alternative — which is what this did — is a pane that looks alive, + // produces nothing, and swallows everything typed into it. Migrating a + // live pty between processes is the other answer and is a different + // feature; a fresh shell is the one this transport can keep. + for (&s.shells) |*sh| { + const owner = sh.owner orelse continue; + if (owner != gone) continue; + sh.owner = null; + sh.owed = true; + } + c.* = .{}; + } +}; + +// --------------------------------------------------------------------------- +// the process +// --------------------------------------------------------------------------- + +/// `pardes --detach[=<name>]`: one core, no terminal, a socket. The loop is the +/// core's own `pump`, exactly as the tty and gui shells run it — this frontend +/// simply has no window of its own. +/// +/// The pre-loop effect drain is here for the same reason tty.zig has one: the +/// startup spawns are already queued, and they have to be PERFORMED before the +/// loop rather than left in the queue. They reach no frontend — there is none +/// yet — and are remembered instead, then asked of the first attach; `Shell` +/// says why that is the only shape that works for a session whose panes exist +/// before its socket does. +pub fn run(init: std.process.Init, opts: pardes.Options, name: []const u8) !void { + const gpa = init.gpa; + const allocs = pardes.allocators.init(gpa); + defer pardes.allocators.deinit(); + var options = opts; + options.image_allocator = allocs.image; + options.pdf_allocator = allocs.pdf; + options.tree_sitter_allocator = allocs.tree_sitter; + options.frame_allocator = allocs.frame; + + // The core's own subsystems, not host work: a detached session syntax + // highlights and decodes images exactly like an attached one. + 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 @import("../look.zig").readFile(gpa, lp); + defer gpa.free(bytes); + break :blk try pardes.Pardes.initFromDump(allocs.pardes, options, bytes); + } else try pardes.Pardes.init(allocs.pardes, options); + defer core.deinit(); + + var session: Session = .{ .gpa = gpa, .core = core, .cols = options.cols, .rows = options.rows }; + defer session.deinit(); + if (!session.listen(name)) { + // Loud, and on stderr rather than through the log: a `--detach` whose + // socket did not bind is a session nobody will ever find, and exiting + // is the only honest answer. + try std.Io.File.stderr().writeStreamingAll(init.io, "pardes: could not bind a detached session socket\n"); + return error.NoSocket; + } + + const h = session.host(); + core.host = h; + while (core.nextEffect()) |effect| core.perform(effect); + while (!core.quit) try core.pump(h); +} + +// --------------------------------------------------------------------------- +// the socket, nested.zig's way +// --------------------------------------------------------------------------- + +/// Re-exported so the frontend half of this transport (client.zig) has ONE +/// import for the socket conventions, and so that the file which owns the +/// convention is the file it asks. The definition and its reasoning are +/// nested.zig's. +pub const setCloexec = nested.setCloexec; + +/// Every descriptor in this transport is non-blocking, on both sides: the core +/// must never park on a peer (`waitInput`), and a frontend must never park on +/// the session (client.zig `wait`). `pub` for that second caller. +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))))); +} + +/// A dead peer must never kill this process, and that is as true of a frontend +/// whose session ended as of a session whose frontend died — so client.zig +/// takes this one too. linux says it per write, darwin once per socket (see +/// `accept`); the `if (darwin)` is what keeps `MSG.NOSIGNAL`, which darwin's +/// headers do not have, out of that build. +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); + +/// Give a drained buffer's memory back, and only a big one's: see +/// `idle_retain`. Called where a queue empties rather than on a timer, because +/// that is the one moment the capacity is provably unused. +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); +} + +/// Monotonic milliseconds, the clock macos.zig's fling already times with and +/// for its reason: MONOTONIC and not REALTIME, because a handshake that expired +/// because NTP stepped the wall clock backwards is a bug nobody reproduces. +/// +/// Zero on failure, and every caller treats zero as "no clock" and enforces no +/// deadline at all — a session that cannot read a clock keeps every slot rather +/// than dropping every slot. +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); +} + +/// Sleep, for the one case that has no descriptor to wait on (see `waitInput`). +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); +} + +/// `<dir>/pardes-detached-<name>.sock`. The prefix differs from nested.zig's +/// `pardes-<pid>.sock` on purpose: that file's sweeper unlinks the socket of any +/// name whose digits name a dead pid, and a session called `work` must never +/// look like one. The buffer is sun_path-sized, so a name that does not fit is +/// no address at all rather than a truncated one pointing somewhere else. +pub fn socketPath(buf: *[sun_path_len]u8, dir: []const u8, name: []const u8) ?[:0]const u8 { + // A name is one path component and nothing clever: a `/` would put the + // socket somewhere else entirely, and a NUL would truncate the address. + 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-"; + +/// The path a FRONTEND connects to for a session called `name`. Derived here +/// rather than in client.zig because this file owns the convention, and the +/// side that binds and the side that connects must not be able to disagree +/// about it. `path_max` is the buffer a caller has to supply. +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 = nested.socketDir(&dir_buf) orelse return null; + return socketPath(buf, dir, name); +} + +/// The FRONTEND's half of the vetting this file does before it binds, and the +/// reason it is here rather than in client.zig: one convention, one predicate, +/// one file that owns both. +/// +/// Until this, the server refused a directory anyone else could write and a +/// socket anyone else could talk to, and the client connected to whatever it +/// found at the path it derived — which is the asymmetry this module's header +/// condemns in as many words. A socket planted at a path a frontend derives +/// from `$XDG_RUNTIME_DIR` receives every keystroke that frontend collects, and +/// answers with frames of its choosing. +/// +/// Checked and then connected, in that order, which is a TOCTOU only for +/// somebody who can already write the directory — and the directory is the +/// first thing this refuses. +pub fn vetted(path: [:0]const u8) bool { + if (comptime !supported) return false; + var dir_buf: [sun_path_len:0]u8 = undefined; + const dir = nested.socketDir(&dir_buf) orelse return false; + if (!ours(nested.statNoFollow(dir) orelse return false, s_ifdir)) return false; + return ours(nested.statNoFollow(path) orelse return false, s_ifsock); +} + +const s_ifmt: u32 = 0o170000; +const s_ifdir: u32 = 0o040000; +const s_ifsock: u32 = 0o140000; + +/// Is this a `kind` we own, with nothing granted to group or other? The three +/// questions `nested.ensureSocketDir` asks of the directory, asked of the +/// SOCKET too: the two walls are the directory's mode and the file's, and a +/// frontend that checks only one of them has checked neither. +fn ours(st: nested.DirFacts, kind: u32) bool { + if (st.mode & s_ifmt != kind) return false; + if (st.uid != libc.getuid()) return false; + return st.mode & 0o077 == 0; +} + +/// Is something LISTENING at `path`? The one place this file decides whether a +/// socket file is a corpse, asked by `listen` before it takes a name over and +/// by `sweep` before it unlinks anything. +/// +/// nested.zig can ask `kill(0)` because its filenames carry a pid; a detached +/// session is named by a PERSON, so the question is put to the socket: a +/// connect to a bound path with no listener is refused (ECONNREFUSED), and that +/// refusal is the ONLY evidence of death this accepts. Everything else is life, +/// including the case a blocking connect used to turn into a hang — a live +/// session busy inside the core has a full backlog and answers EAGAIN, which is +/// why this socket is NON-BLOCKING. EPERM, a socket() that failed and a path +/// that no longer fits are all "not proven dead" too, and leave the file alone. +/// +/// THE WINDOW THIS CANNOT SEE, stated because it is real: a session between its +/// own `bind` and its `listen(2)` also answers ECONNREFUSED and is alive. It is +/// two syscalls wide, it is only ever entered by another `pardes --detach` +/// starting in the same instant, and what the loser loses is a NAME (its +/// `listen` fails and it says so) rather than a session. Closing it needs a +/// lock file per session, which is a second thing to leak. +/// +/// The successful-connect case costs the live session one slot for one round: +/// closing this descriptor immediately turns the pending connection into an +/// EOF, which `receive` reads as a frontend that left. +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); + nested.setCloexec(fd); + setNonblock(fd); + const rc = libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))); + if (rc == 0) return true; + return libc.errno(rc) != .CONNREFUSED; +} + +/// Unlink the sockets of detached sessions that are gone — our own litter, +/// which `--attach`'s "the one session there is" would otherwise count as a +/// session (tty.zig `sessionName`). `alive` is the whole of the judgement. +/// +/// Bounded: one readdir of a directory only we write to, one connect each. +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); + } +} |
