diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-06 18:11:36 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-09-07 13:59:12 -0300 |
| commit | 60367d8fe23f6af98ec28e3cf6c2094dfe332df0 (patch) | |
| tree | 310fc734173cf771881f4691c71909135fadde97 /src/detached/server.zig | |
| parent | fa82cac885cb4738fe36d1e49b4749b5a3e31a4a (diff) | |
| download | pardes-60367d8fe23f6af98ec28e3cf6c2094dfe332df0.tar.gz pardes-60367d8fe23f6af98ec28e3cf6c2094dfe332df0.zip | |
Refactor panes and filesystem; replace FUSE with 9P
Consolidate pane, layout, memory and host code. Serve 9P by default over Unix sockets, with runtime mounts and optional TCP/QUIC transports. Remove FUSE and obsolete proof-of-concept examples.
Fix highlighting and terminal-history performance, expand differential and stress-test infrastructure, sort navigation results while preserving the next occurrence, add syntax-colored Braille minimaps, remove SPC-k, and document 9P interaction as a repository skill.
Diffstat (limited to 'src/detached/server.zig')
| -rw-r--r-- | src/detached/server.zig | 1829 |
1 files changed, 468 insertions, 1361 deletions
diff --git a/src/detached/server.zig b/src/detached/server.zig index db93e2e4..8b49bcb2 100644 --- a/src/detached/server.zig +++ b/src/detached/server.zig @@ -1,617 +1,306 @@ -//! THE DETACHED CORE: one `Pardes` instance in a process with no terminal, -//! serving N frontends over one unix socket. -//! -//! THIS SIDE OWNS THE CORE, AND EVERYTHING UNDER IT. `Session` is a `host.Host` -//! implementation that performs the machine-local half of a host itself — it -//! forks the pane shells, writes the files, watches the paths — and whose -//! `pull_wait_input` is one `poll(2)` over the listener, every attached -//! frontend, every pane's pty master and the inotify descriptor. A frontend owns -//! a screen and a keyboard and nothing else (client.zig). So the `Pardes` is -//! here, `update` is called from here, the shells are forked from here, and the -//! same screen is on every attached frontend at once — `screen -x`, not N -//! sessions. -//! -//! WHAT THIS SIDE SERVES ITSELF, WHICH IS NOW ALL OF IT. A unix socket means -//! the core and its frontends are on the SAME machine, so there is no question -//! of whose process table, whose disk or whose inotify descriptor a call is -//! about — and given that, the process that must hold them is the long-lived -//! one. A shell forked by a frontend dies with that frontend, and a session -//! whose whole promise is outliving the frontend attached to it cannot keep its -//! panes that way. So this file forks the pane shells (`host_io.forkShell`), -//! writes the files (`host_io.writeFileBytes`), marks the directories -//! (file_watch.zig) and drains the pty masters in its own `poll(2)`. THE PANE -//! SHELLS OUTLIVE EVERY FRONTEND: attach, detach, kill the terminal, attach -//! from another one, and the build that was running in pane 3 is still running -//! and has been scrolling into the core the whole time. -//! -//! Only what this vtable leaves null falls through to the core's own -//! `host.Fallback` — and host.zig says in as many words that a zero-method host -//! is a complete pardes. What a frontend can still do BETTER is exactly what -//! needs the human's own display, and nothing else: put a yank on the clipboard -//! in front of them, take a paste off it, open a link in their browser. Three -//! messages, which is why the routing table below is as short as it is. -//! -//! 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. Two rules over four messages: -//! * 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. -//! * ORIGIN, ELSE PRIMARY — `read_clipboard` (the one `pull_` on the wire), -//! `open_link`, and `detach`. Each answers a thing a HUMAN just did, and the -//! answer belongs to that human: the paste must come from the keyboard that -//! asked for it, a link must open in front of the person who clicked it, and -//! a `Detach` typed in one frontend must send THAT frontend away and leave -//! the others painting. `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` — the lowest -//! attached slot, i.e. the oldest surviving attachment, a rule that is -//! stable while frontends come and go and needs no election — covers an -//! effect that no input caused at all. -//! -//! `detach` is the odd one and is worth naming as such: it is not an EFFECT the -//! session performs on the world, it is SESSION CONTROL — one frontend asking to -//! stop being a frontend. That is why wire.zig gives it 0x05, in the -//! 0x01..0x0f session range beside `quit`, rather than a number in the 0x10.. -//! range where every tag is one `push_` method that reaches a disk, a clipboard -//! or a browser. And it is why this side does nothing but send it: see `detach`. -//! There is no third rule, and the class of message it used to serve is gone: -//! `spawn`, `pty_write`, `pty_resize`, `write_file`, `write_dump`, `watch_file`, -//! `watch_theme` and `dump_themes` were routed to ONE frontend precisely -//! because each has one real resource behind it, and every one of them is now -//! performed HERE, once, by the process that owns the resource. Two frontends -//! can no longer fork two shells for pane 3 or race each other writing one -//! path, because neither of them writes anything. -//! -//! FAIRNESS, and why no client — and no shell — can stall the core or another -//! client. The property the bullets below add up to is worth stating as one -//! sentence, because it is what a detached session is FOR: there is no path on -//! which this process blocks indefinitely. Every descriptor it holds is -//! non-blocking, the single `poll(2)` is the only place it sleeps, and every -//! queue that could grow without bound has a ceiling with a stated answer for -//! reaching it. A daemon nobody is looking at cannot be made to stop looking -//! after the shells nobody else is keeping. -//! * ONE `poll(2)` per pump covers the listener, all `max_clients` frontends, -//! all `pardes.MAX_PANES` pty masters and the inotify descriptor: -//! `poll_slots` descriptors, one syscall, no thread per client and none per -//! pty. Putting the shells in the poll set the clients were already in is -//! what lets a daemon own sixteen of them and stay single-threaded. -//! * one read per pty per round, which is `receive`'s rule for clients -//! applied to shells: a `yes` in pane 1 gets one turn and the loop moves on -//! to the other panes, the frontends and the frame. -//! * EVERY descriptor is non-blocking, sockets and pty masters alike, and a -//! pane owes its bytes the same way a client does. A blocking write to a -//! master was the one hole this file's own comment used to argue was safe — -//! "the peer on a pty is a shell this process forked, not a stranger who can -//! stop reading on purpose" — and that was wrong, because the peer is -//! whatever program the human ran in that pane. `sleep 3600` plus a paste -//! larger than the pty's input buffer parked the WHOLE daemon inside -//! `write(2)`: no frame to any frontend, fifteen other masters unread, no -//! `accept`, no `expire`, no inotify drain. So a pane has an out-queue and a -//! POLLOUT, on the descriptor that was already in the set. See `ptyWrite`. -//! * 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 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 host_api = @import("../host.zig"); const wire = @import("wire.zig"); -/// The machine-local half of a host — fork a shell onto a pty, put bytes on a -/// disk — shared verbatim with the tty shell, and the sharing is the point: -/// `spawn` and `writeFile` below are the same two operations tty.zig performs, -/// and having them in one file is what keeps a daemon's pane and a terminal's -/// pane the same pane. See host_io.zig's header for why the daemon is the side -/// that performs them. const host_io = @import("../host_io.zig"); +const selection_pipe = @import("../selection_pipe.zig"); -/// ...and the inotify half, likewise shared: `applyEffect` is the mark-then- -/// reconcile transaction the tty and sdl shells run, and this session runs the -/// identical one. All that differs is who waits on the descriptor — a thread -/// there, `waitInput`'s poll set here. const file_watch = @import("../file_watch.zig"); -/// acme's control filesystem, which a detached session had no way to serve -/// until now: `push_fs_reply` was the one machine-local effect this host left -/// null, so a script could drive a tty or an SDL session and not a daemon — -/// the configuration whose whole promise is outliving the terminal. -/// -/// It costs less here than it does in the desktop shells. They need a thread -/// blocked on `poll()` to notice a request and wake their loop -/// (`fs_service.wake`); this process already owns a `poll(2)` over everything -/// else it waits on, so `/dev/fuse` is one more descriptor in that set and -/// there is no thread at all. `Source.fuse` is the arm; `pollFrame` is the -/// drain, at the same point in the frame that tty.zig drains. -const fuse = @import("../fuse.zig"); -const fs_service = @import("../fs_service.zig"); -/// ...and the SECOND transport onto that same tree, on a unix socket of its -/// own. `Source.ninep_listener` and `Source.ninep` are its arms; `pollFrame` -/// is its drain, beside the mount's, at the same point in the frame. -const fs9_service = @import("../fs9_service.zig"); +const ninep_io = @import("../9p_io.zig"); -/// Host-lifetime storage for the OSC 133 rc files a forked shell sources, held -/// by `Session` because a Session is exactly one host's lifetime. -const shell_bin = @import("../shell_bin.zig"); - -/// `shellCwd` for a pane's shell, `ttyTaken` for a pane the core is about to -/// type a command line into, `readFile` for `run`'s `--load`. const look = @import("../look.zig"); -/// The "saved <path>" / "dumped themes <path>" message row, stamped the way -/// every other host stamps it — one clock format across every frontend. -const message = @import("../message.zig"); - -/// `Dump themes` writes the reference set out as .zon, into the core's own -/// `opts.config_dir`. -const user_config = @import("../user_config.zig"); +const message = pardes.Pardes.Message; -// TIOCSWINSZ: absent from std.c.T on darwin — _IOW('t', 103, winsize). The -// same constant the tty, gui and macos shells spell, for the same reason. const TIOCSWINSZ: c_int = @bitCast(@as(u32, if (@hasDecl(posix.T, "IOCSWINSZ")) posix.T.IOCSWINSZ else 0x80087467)); -/// 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 const darwin = ninep_io.darwin; -/// `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; +const supported = ninep_io.supported; -/// 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; +const sun_path_len = ninep_io.sun_path_len; -/// `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 bounds a backlog of -/// the three things that are still on the wire — a welcome, a clipboard mirror, -/// a link to open — 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 taken off one pane's pty per poll round. 64 KiB is what every other -/// host's pty reader uses (`readPty` in tty.zig, gui.zig and macos.zig), and it -/// sits on `readPty`'s own frame rather than the loop's. Nothing is copied out -/// of it: `Event.output` borrows the buffer for one `update` call, so a daemon -/// serving a shell that is printing a build log asks the allocator for nothing. const pty_chunk = 64 * 1024; -/// Descriptors in the ONE poll this process runs: the listener, every frontend, -/// every pane's pty master, the inotify descriptor behind every watch, -/// `/dev/fuse`, the 9P listener, and every 9P connection. That is -/// 1 + 32 + 16 + 1 + 1 + 1 + 4 = 56 on a full house, and one syscall covers -/// all of them. -const poll_slots = 1 + max_clients + pardes.MAX_PANES + 2 + 1 + fs9_service.max_conns; +const poll_slots = 2 + max_clients + pardes.MAX_PANES + 2 + 2 + @as(usize, @intFromBool(ninep_io.quic_enabled)) + ninep_io.max_conns; -/// How many times `reloadWatched` will honour `file_watch.reloadChanged`'s -/// request for another pass within one round. See `reloadWatched`. const reload_retries = 4; -/// 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: a client's `out` tops out at `out_backlog` plus the one -/// oversized message allowed through whole, so `out_backlog` alone permits -/// 32 * (1 + 1.6) MiB, about 83 MiB of a daemon nobody is looking at. -/// -/// DERIVED, and the derivation IS the fix. This was the literal `4 << 20`, -/// which was by coincidence the exact value of tty.zig's `max_paste_bytes` — -/// and `in` grows to hold one WHOLE message, so a frontend assembling the very -/// paste wire.zig names as one of the two messages that set `max_payload` -/// crossed the table's ceiling while still receiving it. The session then -/// closed the only frontend it had, mid-paste, with `.backlog`, which is the -/// diagnostic for a peer that STOPPED reading. The documented maximum paste -/// could not complete. Two whole `max_payload`s is the smallest number that is -/// headroom rather than another coincidence: one peer may legitimately be -/// assembling a message of the largest size `framed` will accept while the rest -/// of the table holds frames, and past 32 MiB the fattest peer is the peer that -/// stopped draining. Neither the mirrors nor the pane queues are in this -/// number: a mirror is this session's own bookkeeping for a client it chose to -/// serve, and a pane is bounded per pane by `pty_backlog` because it is not a -/// peer and cannot be closed to reclaim anything. const session_backlog = 2 * @as(usize, wire.max_payload); -/// Bytes of un-drained INPUT one pane's shell may owe before more is refused. -/// -/// A pane is not a client, so the answer cannot be `out_backlog`'s: a client -/// that stops draining is closed, and the thing at the other end of a pty is a -/// program the human is running. This refuses the write and says so on the -/// pane's message row instead, which is the only honest answer left — dropping -/// input silently loses half a command line, and killing a shell to reclaim a -/// megabyte destroys work. -/// -/// Checked BEFORE the append, exactly as `queue` checks `out_backlog`, and that -/// is what makes 1 MiB enough: any single write lands whole, so a maximum paste -/// into an empty queue is never truncated. What gets refused is MORE input typed -/// at a program that has stopped reading its input at all — `sleep 3600`, a -/// stopped job, anything blocked on its own output. const pty_backlog = 1 << 20; -/// How long a connection has to say `hello`, and the ONE number both ends of -/// this transport time the handshake against. `pub` because a frontend that -/// waited longer than the session is willing to hold its slot would report a -/// timeout for a slot that had already been taken back, and a frontend that -/// waited less would give up on a session that was still going to answer — two -/// halves of one deadline, and two literals is how they drift apart. -/// -/// The `Session` field it initialises is a field and not this constant for -/// exactly one reason: the test for expiry would otherwise have to sleep five -/// seconds. See `Session.greet_deadline_ms`. pub const greet_deadline_default_ms: u32 = 5_000; -/// 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: forked by THIS process, drained by its poll set, reaped by -/// it. -/// -/// There is no owner here and nothing is owed, and the absence is the whole -/// change. This struct used to record which frontend had been asked to fork a -/// pane and re-ask the next arrival when that frontend left, because the pty -/// lived in the frontend that forked it; and a `--detach`, whose panes always -/// exist before its socket does, had nobody to ask at all and had to remember -/// the request instead. Both were one problem, and forking here dissolves both: -/// a startup layout's shells are forked during `run`'s pre-loop drain with -/// nobody attached, and they are still those same shells when the tenth -/// frontend attaches an hour later. const Pty = struct { - /// The pty master, non-negative exactly when this pane has a live shell. - /// While it is here it is in the poll set (`waitInput`), NON-BLOCKING like - /// every other descriptor this file holds — `spawn` flips it, because - /// `forkpty` hands it back blocking and tty.zig's streaming reader wants it - /// that way. fd: c_int = -1, - /// Kept past the fork for `look.shellCwd` and `look.ttyTaken`, both of which - /// ask /proc about this pid rather than about the descriptor. pid: posix.pid_t = 0, - /// Bytes owed to this shell's stdin, drained by POLLOUT and bounded by - /// `pty_backlog`. The same shape as `Client.out`, for the same reason: the - /// thing on the far side may not be reading, and this process must not wait - /// to find out. See `ptyWrite`. + kill_at: i64 = 0, out: std.ArrayListUnmanaged(u8) = .empty, }; -/// What one descriptor in `waitInput`'s poll set is. A tagged union rather than -/// the bare slot index this loop used to carry alongside its `pollfd`s, because -/// the set now holds four different kinds of thing and a `u8` cannot say which. +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, - /// The acme filesystem's descriptor. Its arm does one thing only: notice - /// that the connection has gone. Being in the set is otherwise the whole - /// point, because a readable `/dev/fuse` must END THE SLEEP so that - /// `pollFrame` — which runs after `pull_wait_input` returns — reaches the - /// drain. - /// - /// The drain is not done HERE, and the reason is not re-entrancy: both - /// `pull_wait_input` and `push_poll_frame` are called from inside - /// `Pardes.pump`, so either would re-enter. It is that `dispatch` is - /// mid-iteration over the `fds[0..n]`/`src[0..n]` SNAPSHOT `waitInput` - /// built, and a filesystem request reaches `core.perform`, which drains the - /// whole effect ring — including a `push_spawn`, whose `Session.spawn` - /// closes a pane's master and forks a new one. A later `.pty` entry in the - /// same pass would then apply the old descriptor's `revents` to a brand new - /// one, and the `fd < 0` guard cannot see it because the fd is valid, - /// merely different. Same hazard as the client slots, same reason. - fuse, - /// The 9P listener, and one arm per connection on it. Same shape and same - /// hazard as `fuse` above, and stated separately only because the hazard - /// is WORSE here: a 9P `Twrite` to `ctl` reaches `core.perform` exactly as - /// a FUSE write does, so it can `push_spawn` and close a pane's master - /// under a later `.pty` entry in this same pass. So `dispatch` moves BYTES - /// — accept, read into `push`, drain `output` — and never serves a - /// request; `pollFrame` does the serving, after the snapshot is done with. ninep_listener, - /// A connection's index in `Listener.conns`. + ninep_quic, ninep: u8, }; pub const Session = struct { gpa: std.mem.Allocator, - /// The Io every filesystem read this host performs goes through: - /// file_watch.zig's reload of a changed pane, and `user_config.dumpThemes`. - /// Required and not optional — a session that owns the disk work cannot be - /// handed a null disk. + worker_gpa: std.mem.Allocator, io: std.Io, 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, - /// Each pane's shell. Forked here, drained by the poll set, reaped by - /// `harvest`. See `Pty`. ptys: [pardes.MAX_PANES]Pty = @splat(.{}), - /// The OSC 133 rc files a forked shell sources, staged once for the life of - /// this host exactly as tty.zig stages them for the life of a terminal: - /// `shell_bin.resolve` hands a child pointers into these buffers and the - /// child holds them until it execs, so they must not live in a stack frame. - /// The default is the empty one, which `resolve` reads as "this shell gets - /// no prompt marks"; `run` supplies a staged one. - prompt_rcs: shell_bin.PromptRcs = .{}, - /// The one inotify descriptor behind every watch this session holds, and the - /// last member of the poll set. Opened lazily — see `inotify`. + 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, - /// acme's control filesystem, or null when `--fs` was not asked for or the - /// mount failed. Owned here rather than by `run` so that `deinit` unmounts - /// on every path out, including the error ones. - fs: ?*fuse.Fs = null, - /// The last drain stopped at `max_batch` with requests still in the kernel. - /// Same role as `check_files`: nothing else will wake us, because no - /// acknowledgement has gone back, so the next round must not sleep. - fs_pending: bool = false, - /// The same tree on a unix socket, or null when `--fs9` was not asked for - /// or the bind failed. Owned here, like `fs`, so that `deinit` closes the - /// socket and unlinks its path on every way out. - ninep: ?*fs9_service.Listener = null, - /// A 9P drain stopped with work still owed. `fs_pending`'s twin, and it - /// needs its own field rather than sharing that one: they are cleared by - /// different drains, and or'ing them into one flag would make a busy 9P - /// script keep the FUSE drain's `pending` set forever. + ninep: ?*ninep_io.Listener = null, ninep_pending: bool = false, - /// WHICH TRANSPORT THE REQUEST BEING SERVED CAME FROM, for the whole of - /// one `fs_service.drain` and never outside one. - /// - /// A filesystem reply reaches its transport through `push_fs_reply`, which - /// is a HOST method: the core answers a `fs_req` and does not know, and - /// must not know, that this process has two filesystems on one tree. This - /// field is the routing origin docs/9p.typ's layering table says a second - /// listener costs, and it is one pointer set by the drain that already - /// knows the answer. - /// - /// Null outside a drain, and `fsReply` then falls back to the mount. That - /// fallback is for a reply the core produced with no request outstanding, - /// which `Fs.reply` drops on a slot lookup; it is NOT a routing guess, and - /// a 9P reply cannot reach it — every 9P request is answered inside the - /// `step` that made it, which is inside the drain that set this. - fs_origin: ?fs_service.Transport = null, - /// Which directory mark belongs to which pane, and the generation the core - /// has already accepted from each. file_watch.zig owns the shape and the - /// transaction; this host owns only the descriptor and the wake. watches: file_watch.Table = @splat(null), - /// A reconcile pass is due: a watched directory had an edge, or the last - /// pass asked for another one. Consumed at the end of `waitInput`, which is - /// where a core change still makes the current frame — `Pardes.pump` renders - /// after `pull_wait_input` returns. check_files: bool = false, - /// False during `run`'s pre-loop effect drain. Read in exactly one place: - /// `Pardes.loadThemeFile` animates an interactive theme change and must not - /// animate a startup one, and a `ThemeFile` in a boot layout is a startup - /// one. tty.zig spells the same distinction `threads_ok`. in_loop: bool = false, - /// 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, and - /// the number itself is `greet_deadline_default_ms`, which the frontend half - /// of this transport reads too. greet_deadline_ms: u32 = greet_deadline_default_ms, - /// 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 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); + } - pub fn deinit(s: *Session) void { - // First, and before the pane shells: a script blocked on `event` is - // holding a kernel request, and `Fs.deinit` answers everything still - // parked and ABORTS THE CONNECTION before unmounting, so that reader - // wakes with ENODEV while its own shell is still alive to run its exit - // path. After `closePty` it would be woken by a hangup instead, with - // nothing left to exit into. - // - // Not, as this comment first claimed, because the harvest takes time: - // `harvest` is `waitpid(WNOHANG)` in a loop and blocks for nothing. The - // one step here that CAN take milliseconds is this one, because - // `Fs.deinit` forks `fusermount3` and waits for it untimed — which is - // also why the `.quit` owed to every frontend now queues behind a - // subprocess. Worth knowing; not worth reordering, because the shells - // matter more than the milliseconds. - if (s.fs) |f| { - f.deinit(); - s.fs = null; + 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)); + } } - // ...and the 9P socket, for the mount's reason above: a script blocked - // on `event` over 9P is woken by the EOF its own connection closing - // produces, while its shell is still alive to run its exit path. The - // orphaned fids are not pumped — the core they would report releases to - // is going with them. + return batch.len != 0 or batch.status != null; + } + + pub fn restore(s: *Session, bytes: []const u8) !void { + const replacement = try s.core.restore(bytes); + 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); + s.check_files = false; + if (s.ninep) |listener| listener.reset(s.core); + 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 { + if (s.mailbox.wake[0] >= 0) pardes.lsp.setStatusSink(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; } - // 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); - // The pane shells go with the SESSION and not with a frontend, which is - // this file's whole change. Closing a master is what hangs its shell up; - // `harvest` collects whatever has already exited, and the process is - // about to leave, so anything slower than that is the kernel's job. for (0..s.ptys.len) |pane| s.closePty(@intCast(pane)); - s.harvest(); - // Every mark dies with the descriptor, so there is nothing to unmark. + 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; } - // Unlinks the two rc files staged for this host's shells. s.prompt_rcs.deinit(); } - /// 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; + 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; - 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. + ninep_io.setCloexec(fd); 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); @@ -623,12 +312,7 @@ pub const Session = struct { 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; @@ -643,15 +327,13 @@ pub const Session = struct { 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 { + pub fn host(s: *Session) host_io.Host { return .{ .ctx = s, .vtable = &vtable }; } @@ -659,60 +341,97 @@ pub const Session = struct { return @ptrCast(@alignCast(ctx.?)); } - /// Eighteen methods, and NOT the fullest host in the tree — that claim - /// stood here, was believed, and was copied into docs/detached.md before an - /// audit counted the others. The tty and SDL shells fill TWENTY each - /// (everything but `pull_gpio_toggle` and `push_detach`) and macOS fifteen, - /// so this host is the only one that implements `push_detach` and otherwise - /// the least complete of the three desktop hosts. What is true is narrower - /// and is the point anyway: it performs every MACHINE-LOCAL effect there is, - /// and the four of host.zig's twenty-two it leaves null are null because - /// there is nothing here for them to do. Two of those four are real losses a - /// person can notice — no `pull_lsp` and no `pull_pipe`, because both want - /// the worker pool this deliberately single-threaded loop does not have. The - /// other two are not losses at all: `push_post_present` marks the moment a - /// frame reached a screen and this process has no screen, and - /// `pull_gpio_toggle` wants pads. - /// - /// `push_fs_reply` was the third real loss until this commit. It was null - /// because the daemon mounted no /dev/fuse, and the consequence was that a - /// detached session — the configuration whose whole promise is outliving the - /// terminal — was the one configuration no script could drive. It mounts one - /// now; see `fs` and `pollFrame`. - /// - /// `push_detach` is the one entry here that is not an effect. See `detach`. - const vtable: host_api.Host.VTable = .{ - .pull_wait_input = waitInput, - .push_present = present, - .push_poll_frame = pollFrame, - .push_spawn = spawn, - .push_pty_write = ptyWrite, - .push_pty_resize = ptyResize, - .push_pty_signal = ptySignal, - .pull_tty_taken = ttyTaken, - .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, - .push_detach = detach, - .push_fs_reply = fsReply, + const vtable: host_io.Host.VTable = .{ + .wait_input = waitInput, + .present = present, + .poll_frame = pollFrame, + .spawn = spawn, + .pty_write = ptyWrite, + .pty_resize = ptyResize, + .pty_signal = ptySignal, + .tty_taken = ttyTaken, + .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, }; - // ---- routing ---------------------------------------------------------- + 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.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) }); + } - /// 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]; @@ -721,8 +440,6 @@ pub const Session = struct { return s.primary(); } - /// 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))); } @@ -731,87 +448,36 @@ pub const Session = struct { for (&s.clients) |*c| if (c.attached) s.send(c, msg); } - // ---- the host methods: pseudo-terminals ------------------------------- - 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 - // The in-process host's reaping rule and its reason, verbatim from - // tty.zig `spawn`: the core reuses pane ids and there is no close - // effect, so a deleted pane's shell lives in its slot until a respawn - // lands here. s.closePty(pane); - var cwd_buf: [256:0]u8 = undefined; - var cwd_z: ?[*:0]const u8 = null; - if (cwd.len > 0 and cwd.len < cwd_buf.len) { - @memcpy(cwd_buf[0..cwd.len], cwd); - cwd_buf[cwd.len] = 0; - cwd_z = @ptrCast(&cwd_buf); - } + 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_z, - // The SESSION grid — which `reconcile` already made the smallest - // common one across everyone attached, and which survives every - // frontend leaving, so a shell forked into an empty session is - // still sized like the pane the core reflowed. + cwd, s.core.screen_h, s.core.screen_w, - s.fs, - ); - // A `forkpty` that failed left `master` holding a number this process - // does not own. The shells get away with not checking because they hand - // the descriptor to a reader task that simply ends; this one would go - // into `poll(2)`, come back POLLNVAL, and be closed out from under - // whoever really owns it. - if (child.pid < 0) return; + s.ninep, + ) catch |err| return s.core.reportError(pane, "shell", err); s.ptys[pane] = .{ .fd = child.file.handle, .pid = child.pid }; - // ...and the master joins the rule every other descriptor in this file - // obeys. `forkpty` hands it back BLOCKING, and host_io.zig leaves it that - // way because tty.zig streams it from a thread that wants a blocking - // read; a poll loop wants the opposite, and one blocking `write(2)` here - // is the whole session parked. Only this side of the pty is affected — - // the master and the slave are separate open file descriptions, so the - // shell's own stdin stays exactly as `forkpty` made it. setNonblock(child.file.handle); - // The pane's starting directory, for the tags. `pollFrame` keeps it - // current after a `cd`; this is the one before the first frame. - var lbuf: [1024]u8 = undefined; - if (look.shellCwd(child.pid, &lbuf)) |wd| s.core.setCwd(pane, wd); + var lbuf: [pardes.memory.limits.host_path_cap + 1]u8 = undefined; + if (host_io.shellCwd(child.pid, &lbuf)) |wd| s.core.setCwd(pane, wd); } - /// Keystrokes and pastes into the shell, QUEUED and never blocked on. - /// - /// This was one blocking `host_io.writeFd`, and the comment defending it - /// argued that "the peer is a shell this process forked rather than a - /// stranger who can stop reading on purpose". The peer is whatever program - /// the human ran in the pane: `sleep 3600`, a job stopped with ^Z, anything - /// blocked writing its own output. Any of those plus a paste larger than the - /// pty's input buffer — four kilobytes, and a frontend is entitled to send a - /// four-MEGABYTE paste — put this single-threaded process to sleep inside - /// `write(2)` with the whole session behind it: no frame to any frontend, - /// fifteen other masters unread, no `accept`, no `expire`, no inotify drain. - /// - /// So a pane owes bytes the way a client does, and the answer is the shape - /// this file already had for exactly this problem. What differs is what - /// happens when the queue will not drain: a client that stops reading is - /// CLOSED, and a pane cannot be, because closing it kills a program the - /// human is running. See `pty_backlog` — the write is refused and said out - /// loud on the pane's own message row. 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]; - // A pane with no shell swallows what is typed at it, which is exactly - // what the core does with a null method. if (pt.fd < 0) return; - // BEFORE the append, which is `queue`'s rule and gives `queue`'s - // guarantee: one write always lands whole, so the biggest paste anyone - // can send is never truncated on arrival, and what is refused is the - // NEXT one typed at a program that has read nothing. if (pt.out.items.len > pty_backlog) { var mbuf: [256]u8 = undefined; const text = std.fmt.bufPrint( @@ -822,13 +488,8 @@ pub const Session = struct { return s.core.setMessage(pane, text); } pt.out.appendSlice(s.gpa, bytes) catch { - // Out of memory for a keystroke. The shell is fine and the session - // is fine; this one write is not, and saying so is all there is. return s.core.setMessage(pane, "input refused: out of memory"); }; - // Try immediately. On an idle pty this empties the queue in one write and - // the descriptor never asks for a POLLOUT at all, which keeps a session - // of keystrokes exactly as cheap as it was. s.flushPty(pane); } @@ -841,56 +502,32 @@ pub const Session = struct { _ = posix.system.ioctl(fd, TIOCSWINSZ, @intFromPtr(&ws)); } - /// `pty/ctl`'s `sig` — and the host where it matters most, because these - /// shells outlive every frontend: a script that signals a build in a - /// detached session is signalling a process nobody has a terminal on. - /// `fd < 0` is a pane with no shell, the same silence `ptyWrite` gives it. 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]; if (pt.fd < 0) return; - look.signalTty(pt.pid, pt.fd, sig); + host_io.signalTty(pt.pid, pt.fd, sig); } - /// Is this pane's tty still the prompt we forked, or has a program taken it? - /// - /// Answerable at all only because the pty is HERE. While a pane's shell - /// lived in a frontend this method had to stay null, and a null one means - /// the core types every `Exec` at the shell — into vim, into a pager, into - /// an agent waiting on stdin. Lazy by construction (host.zig): it runs where - /// the core is about to type a command line, so the /proc walk costs an - /// ordinary frame nothing. 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 look.ttyTaken(pt.pid, pt.fd); + return host_io.ttyTaken(pt.pid, pt.fd); } - // ---- the host methods: the filesystem --------------------------------- - fn writeFile(ctx: ?*anyopaque, pane: u8, path: []const u8, bytes: []const u8) void { const s = of(ctx); - host_io.writeFileBytes(path, bytes) catch |err| + filesystem.write(s.core, path, bytes) catch |err| return s.core.saveFailed(pane, "save", err); - // Our own write is about to come back as an inotify edge: restamp from - // the bytes we just put there so the reconcile reads as "no change". - // Only when this IS the pane's watched file — a `Save <elsewhere>` must - // not silence a real change to the file the pane has open. Six lines - // shared with tty.zig `writeFile` over the same `file_watch.Table`, - // which is what makes a save in a detached pane behave like a save in a - // terminal one. 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 => {}, }; }; - // ...and say so on the pane's message row. AFTER the write, not beside - // it: the early return above is a save that did not happen and must not - // be reported as one. var mbuf: [256]u8 = undefined; s.core.setMessage(pane, message.stamp(&mbuf, "saved", path)); } @@ -899,22 +536,13 @@ pub const Session = struct { const s = of(ctx); var pbuf: [1024:0]u8 = undefined; const path = pardes.dump.outPath(&pbuf) orelse return; - host_io.writeFileBytes(path, bytes) catch |err| return s.core.reportError(0, "dump", err); - // Where it landed, which is what puts `Restore <path>` in the topbar - // (pardes.zig `write_dump`). A dump of a detached session now lands in - // the same directory a terminal session's does, rather than in whatever - // directory the frontend that happened to be primary was started from. + filesystem.write(s.core, path, bytes) catch |err| return s.core.reportError(0, "dump", err); s.core.setLastDump(path); } - fn watchFile(ctx: ?*anyopaque, pane: u8, _: []const u8, on: bool) void { + fn watchFile(ctx: ?*anyopaque, pane: u8, _: []const u8, on: bool, mode: pardes.WatchMode) void { const s = of(ctx); - // The path argument is unused because `applyEffect` takes it off the - // core's own pane, together with the serial and the generation that make - // the reconcile safe. That is the one thing a frontend could not do — it - // had no core — and it is why the frontend's copy of this method needed a - // second table of pathnames to go with the watch table. - if (file_watch.applyEffect(s.core, s.io, s.gpa, s.inotify(), &s.watches, pane, on)) + if (file_watch.applyEffect(s.core, s.io, s.inotify(), &s.watches, pane, on, mode)) s.check_files = true; } @@ -926,10 +554,8 @@ pub const Session = struct { fn dumpThemes(ctx: ?*anyopaque, pane: u8) void { const s = of(ctx); - // A session started without one has nowhere to put them; the core's - // options are the only place that answer lives. const config_dir = s.core.opts.config_dir orelse return; - const out_dir = user_config.dumpThemes(s.io, s.gpa, config_dir, pardes.themes) catch |err| { + 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; }; @@ -938,182 +564,59 @@ pub const Session = struct { s.core.setMessage(pane, message.stamp(&mbuf, "dumped themes", out_dir)); } - // ---- the host methods: the desktop ------------------------------------ - 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 host methods: session control -------------------------------- - - /// `Detach` in an attached frontend: that frontend leaves, the session and - /// every other frontend carry on. tmux's `detach-client`. - /// - /// A `send` and NOTHING ELSE, and each of the three things it does not do is - /// deliberate. It does not quit — the whole point is that the session - /// survives, and a detach that took the daemon with it would be `quit` under - /// another name. It does not touch the core — no pane closes, no shell dies, - /// no frame changes; the grid is retaken by `reconcile` from the frontends - /// that remain, on the ordinary path, because a frontend leaving is already a - /// case this file handles. And it does not close the connection: the frontend - /// closes its own socket when it reads the message, and the peer-hangup path - /// then frees the slot exactly as it does for a frontend somebody killed. - /// Closing from this side would race the frontend's own teardown for no gain. - /// - /// It is therefore the one vtable entry here that is not an effect on the - /// world but SESSION CONTROL — one frontend asking to stop being a frontend - /// — which is why wire.zig numbers it 0x05, in the session range beside - /// `quit`, rather than in 0x10.. where every tag reaches a disk, a clipboard - /// or a browser. The module header's routing table says the same. - /// - /// ORIGIN, ELSE PRIMARY, for `read_clipboard`'s and `open_link`'s reason: it - /// answers something one particular human just typed, so it has to reach that - /// human's screen and not somebody else's — sending a detach to the wrong - /// frontend takes away a session from a person who did not ask. Nobody - /// attached at all is a no-op, and correctly so: there is no frontend to - /// detach, and the core has nothing to record about one. fn detach(ctx: ?*anyopaque) void { const s = of(ctx); if (s.origins()) |c| s.send(c, .detach); } - /// The core's answer to one filesystem request, handed straight back to the - /// transport holding it. `bytes` was resolved by `pardes.fsPayload` inside - /// `perform` and is borrowed only for this call, so a body read is a window - /// onto the pane's live text and copies nothing. `.again` needs no case: - /// both transports read the status and re-park the request themselves. - /// - /// WHICH transport is `fs_origin`, set by the drain that asked. This used - /// to be `s.fs` unconditionally, which was right while a mount was the only - /// answer there was and became a silent misroute the moment `--fs9` gave - /// the session a second one: `Fs.reply` looks the tag up in ITS park table, - /// finds nothing, and returns — so a 9P `Tattach` was answered into the - /// void and its client waited forever. Measured against plan9port's `9p`. - fn fsReply(ctx: ?*anyopaque, reply: *const pardes.acmefs.Reply, bytes: []const u8) void { - const s = of(ctx); - if (s.fs_origin) |t| return t.reply(reply, bytes); - if (s.fs) |f| f.reply(reply, bytes); - } - - // ---- pane shells ------------------------------------------------------ - - /// Each pane's live cwd, for the tags. One readlink of /proc per pane that - /// has a shell, per frame, which is what tty.zig's `pollFrame` costs — and - /// why the far more expensive question, whether a program has taken the - /// pane's tty, is a pull asked at the `Exec` that cares (`ttyTaken`) - /// instead of polled here. - /// - /// This could not exist before. A detached session's shells lived in a - /// frontend, and a frontend has no core to report a cwd TO, so a `cd` in a - /// detached pane never reached its tag no matter how many frontends were - /// watching. The pids are here now, so it does. fn pollFrame(ctx: ?*anyopaque) void { const s = of(ctx); - // acme's filesystem first in the pass, for the reason tty.zig gives at - // its own call site: an edit a script just made through `body` belongs - // in the surface this frame composes, not the next one. The flag is - // read by `waitInput`, which must not sleep while the kernel still has - // requests we have not acknowledged. - if (s.fs) |f| { - s.fs_origin = f.transport(); - s.fs_pending = fs_service.drain(s.fs_origin.?, s.core).pending; - s.fs_origin = null; - } - // ...and the 9P connections, in the same breath and for the same - // reason. One connection is one `Transport`, so this is - // `fs_service.drain` per connection per frame, with `fs_origin` naming - // the one being served so its replies come back to it (see `fsReply`). - // HERE rather than in `dispatch` for `Source.ninep`'s reason: a - // `Twrite` to `ctl` can `push_spawn`, and `dispatch` is mid-iteration - // over a descriptor snapshot when it runs. - if (s.ninep) |l| { - // Before the drain, so a slot held by silence is taken back on the - // same frame it expires rather than one drain later. - l.expire(); - var pending = false; - for (0..fs9_service.max_conns) |i| { - const t = l.transport(@intCast(i)) orelse continue; - s.fs_origin = t; - const d = fs_service.drain(t, s.core); - s.fs_origin = null; - if (l.settle(@intCast(i), d)) pending = true; - } - s.ninep_pending = pending; - } + if (s.ninep) |l| s.ninep_pending = if (ninep_io.quic_enabled and l.quic != null) l.tick(s.core).pending else l.drain(s.core).pending; for (&s.ptys, 0..) |*pt, pane| { if (pt.fd < 0) continue; - var lbuf: [1024]u8 = undefined; - if (look.shellCwd(pt.pid, &lbuf)) |cwd| s.core.setCwd(pane, cwd); + var lbuf: [pardes.memory.limits.host_path_cap + 1]u8 = undefined; + if (host_io.shellCwd(pt.pid, &lbuf)) |cwd| s.core.setCwd(pane, cwd); } } - /// Drop a pane's shell: out of the poll set, out of the process. Closing the - /// master is what hangs the shell up — which is true only because - /// `host_io.forkShell` puts FD_CLOEXEC on it, so no LATER pane's shell is - /// still holding a copy open. The pid is left to `harvest`, because a - /// `waitpid` here would return 0 for a shell that has not noticed the hangup - /// yet and that answer is worth nothing. - /// - /// The core is NOT told. Its two callers are `spawn` — a respawn, where the - /// core is the thing that asked — and `deinit`, where there is no core left - /// to tell. The path that does tell it is `paneEof`. fn closePty(s: *Session, pane: u8) void { const pt = &s.ptys[pane]; if (pt.fd < 0) return; _ = libc.close(pt.fd); - // Before the reset, or the queue's allocation goes with the slot: what - // is in it is input a program that is not reading never took, and there - // is nobody left to hand it to. + pt.fd = -1; pt.out.deinit(s.gpa); - pt.* = .{}; + 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; + }; } - /// Push what the kernel will take of what this pane owes its shell, and - /// leave the rest for a POLLOUT. `flush`'s body, on a pty instead of a - /// socket, down to the `retire` that hands a drained megabyte back. - /// - /// The one difference is what an error means. A failed write to a SOCKET - /// closes a client; a failed write to a master means the slave side is gone, - /// which is the same event as a read of 0. It is NOT the same moment, - /// though, and that is why the error arm reads the pane before it ends it: - /// linux's `n_tty_write` returns EIO the instant the slave has no open - /// descriptors left, while `n_tty_read` on that same master still hands back - /// what the shell wrote before it went — so the write fails while the last - /// line is still retrievable, and ending the pane first would throw it away. - /// That is the very thing the dispatch's `.pty` branch protects against when - /// it takes POLLIN before POLLHUP, and it has to hold here too, because two - /// paths reach this arm with no read of their own in between: the dispatch - /// runs POLLOUT before POLLIN, and `ptyWrite` calls this during `perform`, - /// after this round's `readPty` has already been and gone. fn flushPty(s: *Session, pane: u8) void { const pt = &s.ptys[pane]; var off: usize = 0; @@ -1121,22 +624,13 @@ pub const Session = struct { 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, - // The pty's input buffer is full: the rest waits for POLLOUT, - // and this is the case the whole change exists for. .AGAIN => break, else => { - // `readPty` either takes that last chunk or reaches the end - // itself and has already ended the pane; the guard is what - // stops the second `paneEof` from being a double-end. s.readPty(pane); if (s.ptys[pane].fd >= 0) s.paneEof(pane); return; }, }; - // No progress and no error. host_io.zig's `writeFd` says why this is - // a `break` and never a retry: looping on a zero-byte write is a - // spin, and a spin in here is the whole session at 100% of a core - // with no syscall for a signal to interrupt. if (n == 0) break; off += @intCast(n); } @@ -1149,119 +643,51 @@ pub const Session = struct { pt.out.items.len -= off; } - /// One read per readable pty per round — `receive`'s rule for clients, - /// applied to shells: a `yes` in pane 1 gets one turn and the loop moves on - /// to the other panes, the frontends and the frame. - /// - /// Nothing is copied. `Event.output` borrows the buffer for the length of - /// one `update` call, which is the same borrow window every other host gives - /// a pty chunk — tty.zig frees its duplicate the line after the update — - /// except that this one never allocated a duplicate to free. A daemon - /// serving sixteen shells printing build logs asks the allocator for - /// nothing. 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)) { - // A master that said POLLIN and then had nothing is not an error; - // the next round asks again. .INTR, .AGAIN => {}, - // EIO is how linux reports the slave side going away, which is the - // ordinary end of a shell rather than a fault. else => s.paneEof(pane), }; s.core.update(.{ .output = .{ .pane = pane, .bytes = buf[0..@intCast(got)] } }); } - /// The shell in `pane` is gone. The descriptor leaves the poll set BEFORE - /// the core is told, because an `eof` is what makes the core offer a respawn - /// and a respawn into a slot still holding the old fd would leak it. fn paneEof(s: *Session, pane: u8) void { s.closePty(pane); s.harvest(); s.core.update(.{ .eof = .{ .pane = pane } }); } - /// Collect every child that has exited. - /// - /// `waitpid(-1)` and not a pid list, because the only children this process - /// LEAVES UNREAPED are pane shells (`host_io.forkShell`) — so "any exited - /// child" and "an exited pane shell" are the same set — and because the pids - /// a list would hold are exactly the ones it cannot help with: a respawn - /// closes a master, and the shell that gets the hangup exits some - /// milliseconds later with its slot already reused by a different shell. - /// Anything else this process forks — `fusermount3`, from `Fs.mount`, - /// `sweepStale` and `Fs.deinit` — is reaped by its own spawner with a - /// pid-specific blocking wait before control returns here, so the set this - /// sees is still only shells. A future worker that forks and does not wait - /// would break that, and this is the sentence it has to come back and edit. - /// - /// Nothing here waits, so a session whose shells are all running pays one - /// syscall that returns 0. Called once per poll round and again wherever a - /// shell is dropped, which is what keeps a daemon that runs for a week and - /// spawns a thousand shells free of zombies — the one bookkeeping cost a - /// long-lived process pays that a frontend, which exits, never did. - fn harvest(_: *Session) void { - while (true) { - // 0: there are children and none has exited. -1: no children at all. - if (libc.waitpid(-1, null, libc.W.NOHANG) <= 0) return; - } + 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); } - // ---- watched files ---------------------------------------------------- + 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; + } + } - /// The one inotify descriptor behind every watch, opened on first use. - /// - /// Lazy for two reasons pointing the same way: a `Session` is built as a - /// struct literal (client.zig's test harness is one) and so has no init hook - /// to open it in, and a session whose panes are all shells never watches a - /// path and has no use for one. -1 on anything but linux and on a failed - /// `inotify_init1`, which file_watch.zig reads as "mark nothing" — the core - /// then keeps its own record of what was asked and simply never gets a - /// reload, which is what a host with no watcher has always done. - /// - /// `polled` is true because this descriptor is drained from `poll`, not - /// from a thread parked in a wait (tty.zig `watchFiles`): `drainInotify` - /// must be able to stop. On linux that is IN_NONBLOCK; on macos a kqueue - /// needs nothing, since the timeout argument to `kevent(2)` decides. 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; } - /// A directory this session marked had an edge. The CONTENTS are discarded - /// on purpose, exactly as tty.zig's watcher thread discards them: a record - /// names a mark and a filename, and reconciling every mark against the - /// generation the core accepted is both cheaper and safer than deciding from - /// the record which pane it meant. - /// - /// DRAINED TO EMPTY, in a loop, and one read was a real cost rather than the - /// coalescing this comment used to claim. The descriptor is level-triggered, - /// so a queue left partly full makes `poll` return ready again immediately — - /// and each of those rounds is a whole `pump`: `reloadChanged` over all 17 - /// slots, every watched text pane re-read from disk and re-hashed, a render, - /// a present. A `git checkout` can queue the kernel's whole 16384 events; at - /// roughly 128 records per 4 KiB that was ~128 spin rounds and some two - /// thousand whole-file reads for one command, at 100% of a core, while every - /// frontend got a frame per round it could not use. The fd is IN_NONBLOCK - /// (`inotify`), so the loop ends on EAGAIN. fn drainInotify(s: *Session) void { if (file_watch.drain(s.inotify_fd)) s.check_files = true; } - /// Reconcile every marked pane and the theme file. Called at the END of - /// `waitInput`, which is what puts a reload in THIS frame: `Pardes.pump` - /// renders after `pull_wait_input` returns. - /// - /// The loop is `reloadChanged`'s contract. It asks for another pass when a - /// PDF's pathname changed between the stat before MuPDF reopened it and the - /// stat after — a save that landed mid-reconcile, where committing either - /// identity would lose a generation. tty.zig posts that request back into - /// its event queue; this loop has no queue, so it is retried here and - /// BOUNDED, because a file being rewritten in a loop must not hold the core. - /// What is left over is picked up by the next directory edge. fn reloadWatched(s: *Session) void { if (!s.check_files) return; s.check_files = false; @@ -1270,16 +696,10 @@ pub const Session = struct { } } - // ---- 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); } @@ -1289,8 +709,6 @@ pub const Session = struct { 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 @@ -1307,61 +725,31 @@ pub const Session = struct { 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, and the only place it - /// waits on ANYTHING: one `poll(2)` over the listener, every attached - /// frontend, every pane's pty master and the inotify descriptor. No thread - /// per client, no thread per shell, no watcher thread, and nothing here - /// blocks on a single peer. - /// - /// That the shells are in this set and not on threads of their own is what - /// lets a daemon own sixteen of them and stay a single-threaded state - /// machine — and it costs the shells nothing, because a pty master is - /// pollable and a pane's output has nowhere to go but the core this loop is - /// driving anyway. 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`). + const completed = s.drainCompletions(true); for (&s.clients) |*c| if (c.fd >= 0) s.flush(c); - // The session grid, retaken BEFORE the sleep as well as after it. A - // client can leave OUTSIDE this function — a `set_clipboard` broadcast - // whose write failed during `perform` closes it — and the minimum across - // attached frontends would then stay sized for a frontend that is gone - // until some descriptor happened to become readable, which on an idle - // session is never. That is what this call buys, and `regridded` is what - // it costs: a round that has just told the core to reflow must not then - // sleep on it, because the frame carrying that reflow is the one - // `Pardes.pump` composes the moment this returns. const regridded = s.reconcile(); const now = monotonicMs(); var fds: [poll_slots]libc.pollfd = undefined; var src: [poll_slots]Source = 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. + 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 }; @@ -1378,68 +766,39 @@ pub const Session = struct { src[n] = .{ .client = @intCast(i) }; n += 1; } - // The pane shells, and note what is NOT conditional on a frontend: a - // session with nobody attached still polls these, still reads them and - // still feeds the core. That is the difference between a detach that - // pauses your build and a detach that does not. for (&s.ptys, 0..) |*pt, pane| { if (pt.fd < 0) continue; fds[n] = .{ .fd = pt.fd, - // POLLOUT only while this pane owes its shell bytes, which is - // the same rule and the same reason as a client's: asking for it - // unconditionally makes every idle pty a ready descriptor and - // turns the poll into a spin. .events = if (pt.out.items.len != 0) poll_in | poll_out else poll_in, .revents = 0, }; src[n] = .{ .pty = @intCast(pane) }; n += 1; } - // Opened only once something asked to be watched, so an unwatched - // session simply has one fewer descriptor here (see `inotify`). if (s.inotify_fd >= 0) { fds[n] = .{ .fd = s.inotify_fd, .events = poll_in, .revents = 0 }; src[n] = .inotify; n += 1; } - // ...and acme's filesystem, when there is one. Its arm in `dispatch` - // does nothing: this descriptor is here to END THE SLEEP, so that the - // `pollFrame` after `pull_wait_input` returns reaches the drain. The - // desktop shells buy the same wake with a thread; one poll slot is - // cheaper and cannot race the loop. - // ...and NOT once it is dead. `fuse_dev_poll` answers `EPOLLERR` as soon - // as the connection is gone, POSIX reports `POLLERR` whatever the events - // mask asked for, and this arm cannot consume it — so an external - // `fusermount3 -u`, a sysfs abort, or systemd taking `/run/user/$UID` - // away at final logout (exactly when a detached session is supposed to - // keep running) would make `poll(2)` return instantly, forever, and burn - // a whole core for the life of the daemon. Measured at 100% of one CPU - // before this guard. fuse.zig's own poll thread has carried the - // equivalent check all along, which is why the desktop shells never - // showed it and this loop did. - if (s.fs) |f| if (f.fd >= 0 and !f.dead) { - fds[n] = .{ .fd = f.fd, .events = poll_in, .revents = 0 }; - src[n] = .fuse; - n += 1; - }; - // ...and the 9P socket, when there is one: the listener, plus one - // descriptor per live connection. POLLOUT only while a connection owes - // bytes, which is the same rule and the same reason as a client's and - // a pty's — asking for it unconditionally makes every idle socket a - // ready descriptor and turns the poll into a spin. if (s.ninep) |l| { - // `accepting`, not `fd >= 0`: a listener paused after an EMFILE - // must leave the set, or the backlog it could not drain reports - // ready on every poll and spins the core. `nextDue` carries the - // moment it comes back. if (l.accepting()) { - fds[n] = .{ .fd = l.fd, .events = poll_in, .revents = 0 }; - src[n] = .ninep_listener; - n += 1; + for ([_]c_int{ l.fd, l.tcp_fd }) |fd| { + if (fd < 0) continue; + fds[n] = .{ .fd = fd, .events = poll_in, .revents = 0 }; + src[n] = .ninep_listener; + n += 1; + } + } + if (comptime ninep_io.quic_enabled) { + if (l.quic) |*listener| { + fds[n] = listener.poll(); + src[n] = .ninep_quic; + n += 1; + } } - for (0..fs9_service.max_conns) |i| { - if (!l.live(@intCast(i))) continue; + for (0..ninep_io.max_conns) |i| { + if (l.conns[i].fd < 0) continue; fds[n] = .{ .fd = l.conns[i].fd, .events = if (l.owes(@intCast(i))) poll_in | poll_out else poll_in, @@ -1449,108 +808,38 @@ pub const Session = struct { n += 1; } } - // A session with no listener, no clients, no shells and no watches 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. - // - // Nothing is owed on this path. `check_files` is only ever set by a - // watch, and a watch means the inotify descriptor is in the set; - // `reconcile` posts a resize only when a client is ATTACHED, which means - // its socket is in the set; and `fs_pending` is only ever set by a drain, - // which runs only when `fs` is live, which puts `/dev/fuse` in the set. - // So `n == 0` implies `!regridded` and `!fs_pending` too. `--fs9` is - // the ONE exception, and it is why `ninep_pending` is named on the - // `timeout = 0` line below rather than here: a hung-up 9P connection - // still owing the core its orphaned fids has no descriptor at all, so - // it can be the only work left with `n == 0`. It is also finite — one - // release per fid — so the 16 ms nap that path takes costs it a couple - // of frames and never a stall. 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); - // ...and two things are due on nothing at all rather than on a - // descriptor: a reconcile pass a watch effect asked for (`watchFile` ran - // during `perform`, outside this function) and a regrid this round has - // already performed. Both are consumed before this function returns, so - // the round must not sleep before reaching them. - if (s.check_files or regridded or s.fs_pending or s.ninep_pending) timeout = 0; + if (completed or s.check_files or regridded or s.ninep_pending) timeout = 0; const ready = libc.poll(&fds, @intCast(n), timeout); - // A timeout is an ordinary frame boundary and EINTR is a signal we do not - // handle here. Neither skips anything below any more: what used to be an - // early `return` here is why a client closed without any descriptor being - // readable — which is every `expire` — left the session grid sized for a - // frontend that had gone, until the next readable event, on an idle - // session possibly hours later. if (ready > 0) s.dispatch(fds[0..n], src[0..n]); - // AFTER the dispatch, and that ordering is itself a fix. `expire` frees a - // client slot and `accept` — which runs INSIDE the dispatch — fills the - // lowest free one, so an expire that ran first could hand a slot to a new - // connection within this same round and the dispatch would then apply the - // OLD connection's `revents` to the new descriptor: a POLLHUP from the - // peer that left, closing the peer that just arrived. The dispatch's - // `c.fd < 0` guard cannot see that, because the fd is perfectly valid — - // it is simply a different fd. Expiring after means a freed slot is - // refilled no earlier than the next round, which builds a fresh `fds` for - // it. It fixes a smaller thing for free, too: a connection whose `hello` - // arrived in THIS round is attached before its deadline is judged, - // instead of being taken back with its handshake still unread. + _ = s.drainCompletions(true); s.expire(monotonicMs()); - // Unconditional, and not only where a shell is noticed to have died: a - // shell whose master `spawn` closed on a respawn exits after that close, - // with no descriptor left for anyone to see it on. See `harvest`. s.harvest(); - // Both before this function returns, so a file that changed on disk and a - // frontend that left during this round are in the frame `Pardes.pump` - // composes next rather than the one after it. s.reloadWatched(); _ = s.reconcile(); } - /// One pass over the descriptors `poll` reported ready. Split out of - /// `waitInput` for one reason: everything that must happen AFTER it — - /// `expire`, `harvest`, `reloadWatched`, `reconcile` — is then stated once, - /// in one order, where no early return can skip it. An early return past - /// that list is exactly what findings 5 and 6 were. 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]; - // A slot closed earlier in this same pass (its peer hung up, a - // decode failed) must not be touched through a stale revents. 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) { - // POLLIN wins when both are set: a peer that wrote and then - // closed has bytes still worth reading. s.close(c, .peer); } }, .pty => |pane| { if (s.ptys[pane].fd < 0) continue; - // What this pane still owes its shell, which is the whole of - // finding 1's drain: `ptyWrite` queued it and stopped at EAGAIN - // rather than sleeping, and this is where the rest goes. if (pfd.revents & poll_out != 0) s.flushPty(pane); - // `flushPty` ends the pane when the slave side has gone. if (s.ptys[pane].fd < 0) continue; - // The same precedence as a client's, and it matters more here: a - // shell that printed its last line and exited reports - // POLLIN|POLLHUP together, and taking the hangup first would - // throw that line away. `readPty` reaches the EOF by reading 0. if (pfd.revents & poll_in != 0) { s.readPty(pane); } else if (pfd.revents & (poll_hup | poll_err | poll_nval) != 0) { @@ -1558,23 +847,10 @@ pub const Session = struct { } }, .inotify => if (pfd.revents & poll_in != 0) s.drainInotify(), - // The only revents worth a word: a dead connection must leave the - // set, or the `POLLERR` it reports on every future poll spins the - // loop. `Fs.next` would set `dead` on its first failed read anyway; - // saying it here costs nothing and saves the one spinning round. - .fuse => if (pfd.revents & (poll_hup | poll_err | poll_nval) != 0) { - if (s.fs) |f| f.dead = true; - }, .ninep_listener => if (pfd.revents != 0) { if (s.ninep) |l| l.accept(); }, - // BYTES ONLY. The requests those bytes decode into are served by - // `pollFrame`, after this snapshot is done with — see - // `Source.ninep`. POLLOUT before POLLIN so a reply the last frame - // could not finish writing goes before more work arrives, and - // POLLIN before the hangup because a script that wrote a `Tclunk` - // and closed has bytes still worth reading; the `live` guard is the - // client slots' `c.fd < 0`, for its reason. + .ninep_quic => {}, .ninep => |i| if (s.ninep) |l| { if (!l.live(i)) continue; if (pfd.revents & poll_out != 0) l.flush(i); @@ -1588,13 +864,17 @@ pub const Session = struct { }; } - /// 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.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); @@ -1602,8 +882,6 @@ pub const Session = struct { } 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; - // ...and the 9P listener's own greet deadline, for the reason its - // `expire` gives: four slots is a cheaper denial than thirty-two. if (s.ninep) |l| { if (l.nextDue()) |ms| { const at9 = now + ms; @@ -1614,10 +892,6 @@ pub const Session = struct { 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| { @@ -1626,55 +900,27 @@ pub const Session = struct { } } - /// 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); + ninep_io.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; @@ -1683,10 +929,6 @@ pub const Session = struct { } } - /// 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); @@ -1696,9 +938,6 @@ pub const Session = struct { 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); @@ -1709,9 +948,6 @@ pub const Session = struct { 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; @@ -1726,21 +962,7 @@ pub const Session = struct { } fn apply(s: *Session, c: *Client, slot: u8, tag: u8, payload: []const u8) wire.Error!void { - // `Hello.version` BEFORE the payload is decoded, which is the whole - // point of wire.zig putting it first at a fixed offset: a mismatch has to - // stay diagnosable when the rest of the layout is the part that changed. - // Checking it inside the `.hello` arm defeated exactly that guarantee — - // `decodeClient` refuses a cols/rows this build does not like and refuses - // trailing bytes, so a v2 hello with one extra field came back as - // `.protocol` and the `refuse .version` the frontend needs to say - // something useful was never sent. `wire.helloVersion` reads the one - // field without decoding the rest, and lives in the file that owns the - // layout. if (tag == @intFromEnum(wire.ClientTag.hello)) { - // A second hello on one connection is not a resize; it is a peer - // that is not speaking this protocol. Judged here rather than in the - // arm below so that a repeat hello is a protocol error whatever - // version it claims. if (c.attached) return error.BadValue; const claimed = try wire.helloVersion(payload); if (claimed != wire.version) { @@ -1755,20 +977,12 @@ pub const Session = struct { 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; @@ -1782,14 +996,6 @@ pub const Session = struct { } } - /// 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. - /// - /// True when the CORE was told to reflow, which is the one thing a caller - /// has to react to: the frame carrying that reflow is the next one - /// `Pardes.pump` composes, so a `waitInput` that hears true must not go to - /// sleep before returning. See its `regridded`. fn reconcile(s: *Session) bool { var cols: u16 = 0; var rows: u16 = 0; @@ -1798,17 +1004,10 @@ pub const Session = struct { 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. var regridded = false; 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 } }); regridded = true; @@ -1821,18 +1020,10 @@ pub const Session = struct { return regridded; } - // ---- 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`, and with - // every effect that carried a whole file gone from this protocol the - // only payload that can still get there is a yank of more than - // 16 MiB. The session keeps it — `setClipboard` put it in the core's - // own clipboard before this was ever queued — and the frontends' - // desktop clipboards do not get it, out loud rather than silently. log.debug("message {t} not encodable: {t}", .{ msg, err }); return; }; @@ -1840,36 +1031,13 @@ pub const Session = struct { } 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, plus a frame apiece. 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. - /// - /// `items.len` and NOT `capacity`, which was half of the bug in - /// `session_backlog`'s history. An ArrayList grows geometrically, so a - /// client's `in.capacity` crossed a 4 MiB ceiling while it was still - /// assembling a paste of roughly 2.8 MiB — the peer was punished for the - /// allocator's rounding rather than for anything it held. What this is - /// asking is "how much is a peer making this session hold RIGHT NOW", and - /// that is `items.len`; capacity above it is transient by construction, - /// because `retire` hands back anything over `idle_retain` the moment a - /// buffer empties. fn account(s: *Session) void { var total: usize = 0; var worst: ?*Client = null; @@ -1893,8 +1061,6 @@ pub const Session = struct { 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), }; @@ -1910,18 +1076,11 @@ pub const Session = struct { 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; @@ -1934,10 +1093,6 @@ pub const Session = struct { } } - /// Free one slot. A frontend dying takes NOTHING with it: not the core, not - /// the listener, not another frontend's frames, and — since this file - /// forks — not its panes' shells either. 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}); @@ -1946,49 +1101,23 @@ pub const Session = struct { 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 that is the whole of it. A frontend used to take its panes' - // shells with it and leave them owed to whoever attached next, because - // the ptys were in its process; a pane that survived a detach looked - // alive, produced nothing and swallowed everything typed into it. The - // shells are here now, so a frontend leaving is a screen going away and - // nothing else. 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, and -/// it is no longer half a promise: the startup spawns are already queued and are -/// PERFORMED here, on this process's own process table. So a session binds its -/// socket with every pane's shell already forked and already in the poll set, -/// and the first frontend to attach — whether that is a second later or the -/// next morning — is sent a frame of shells that have been printing into the -/// core since before it existed. Nothing is remembered for a later frontend, -/// because nothing is owed to one. 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(); + const allocs = pardes.memory.init(gpa); + defer pardes.memory.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); @@ -1999,86 +1128,144 @@ pub fn run(init: std.process.Init, opts: pardes.Options, name: []const u8) !void } const core = if (options.load_path) |lp| blk: { - const bytes = try look.readFile(gpa, lp); + const bytes = try filesystem.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, + .worker_gpa = allocs.lsp, .io = init.io, .core = core, .cols = options.cols, .rows = options.rows, - // Staged before the first fork and owned by the Session for exactly as - // long as it can fork: `shell_bin.resolve` hands a child pointers into - // these buffers, and the child holds them until it execs. - // - // `prepareForFork` and not `PromptRcs.init` alone: this host forks bash - // through the same `resolve` its siblings do and was the one that never - // silenced Apple's zsh-migration banner, so every pane in a detached - // session on macOS opened with it printed across the top. It also had - // no `adoptSystemPath`, which a daemon needs more than anyone — it is - // the host most likely to be started by launchd. - .prompt_rcs = shell_bin.prepareForFork(), + .prompt_rcs = host_io.Shell.prepare(), }; + defer session.core.deinit(); defer session.deinit(); + try session.initAsync(); 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; } - // Before the host is installed and before the startup drain, so a script - // that races the daemon's launch finds a tree whose panes already exist. - // Null on every failure — no fuse3, no `user_allow_other`, a kernel without - // FUSE — and a failure must cost the operator their scripting, never their - // session. `fs_service.start` has already said so on pane 0's message row. - // - // NO `fs_service.wake`. That call exists to start a thread that blocks on - // `poll()` and pokes a loop the thread does not otherwise share; this - // process polls `/dev/fuse` itself, in the same syscall as everything else. - // See `Source.fuse`. - session.fs = fs_service.start(gpa, core); - // ...and the same tree on a unix socket, independently: `--fs` and `--fs9` - // are two transports and neither is the other's prerequisite, so a daemon - // may serve one, both or neither. Null on every failure, for - // `fs_service.start`'s reason — a transport that will not bind must cost - // the operator their scripting, never their session — and NOT a - // `return error` the way the frontend socket above is, because a session - // whose frontend socket did not bind is one nobody can ever find, while - // this one is merely one nobody can script over 9P. - // - // A bare `--fs9` is named by the SESSION rather than by the pid: the - // operator typed that name to find the daemon again, and having to look up - // a pid to reach its filesystem would undo it. `--fs9=<name>` wins. - if (opts.fs9) |named| session.ninep = fs9_service.open(gpa, named, name); + session.ninep = ninep_io.listen(gpa, opts.ninep_name, name, opts.ninep_tcp, opts.ninep_quic); + if (session.ninep == null) return error.ListenFailed; + core.fs.socket_path = session.ninep.?.path(); + core.fs.tcp_address = session.ninep.?.tcp_address; + core.fs.quic_address = session.ninep.?.quic_address; const h = session.host(); core.host = h; while (core.nextEffect()) |effect| core.perform(effect); - // Past the startup drain: a `ThemeFile` reload from here on is a human's - // and animates. See `in_loop`. session.in_loop = true; - while (!core.quit) try core.pump(h); + while (!session.core.quit) { + try session.core.pump(h); + if (session.core.quit) break; + if (session.core.takeRestore()) |path| restore: { + const bytes = filesystem.readFile(gpa, path) catch |err| { + session.core.reportError(session.core.active, "Restore", err); + break :restore; + }; + defer gpa.free(bytes); + session.restore(bytes) catch |err| session.core.reportError(session.core.active, "Restore", err); + } + } } -// --------------------------------------------------------------------------- -// the socket, nested.zig's way -// --------------------------------------------------------------------------- +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 s.core.dumpState(); + 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); + 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 "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.cur_col = 2; + pane.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); +} -/// 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; +pub const setCloexec = ninep_io.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; @@ -2087,11 +1274,6 @@ pub fn setNonblock(fd: c_int) void { _ = 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); @@ -2100,35 +1282,17 @@ 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. -/// -/// `pub` for the same reason `setNonblock`, `nosignal` and the `poll_*` -/// constants are: this file owns the transport's conventions and BOTH ends of -/// it, and the clock a handshake is timed against is one of them. client.zig -/// times its wait for a `welcome` on this and against -/// `greet_deadline_default_ms`, so the two ends cannot disagree about how long -/// the handshake is allowed to take. 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); } -/// 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), @@ -2137,14 +1301,7 @@ fn nap(ms: u32) void { _ = 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; @@ -2152,78 +1309,33 @@ pub fn socketPath(buf: *[sun_path_len]u8, dir: []const u8, name: []const u8) ?[: 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; + const dir = ninep_io.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 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; -/// 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 { +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; } -/// 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; @@ -2231,18 +1343,13 @@ fn alive(path: [:0]const u8) bool { const fd = libc.socket(libc.AF.UNIX, libc.SOCK.STREAM, 0); if (fd < 0) return true; defer _ = libc.close(fd); - nested.setCloexec(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; } -/// Unlink the sockets of detached sessions that are gone — our own litter, -/// which the bare `Attach`'s "whichever session is there" would otherwise count -/// as a session (client.zig `resolve`). `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); |
