diff options
| author | Gabriel Schneider <[email protected]> | 2026-08-26 13:27:46 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-08-27 09:47:39 -0300 |
| commit | 11f380f6d7222f2cad93c2cdf13701ea1f903d47 (patch) | |
| tree | 803194ee5853a6b4cda93f90a95e28d1f02e69ae /src/detached | |
| parent | fbc194068687e49a8490c85c9f1257a2f2bb9079 (diff) | |
| download | pardes-11f380f6d7222f2cad93c2cdf13701ea1f903d47.tar.gz pardes-11f380f6d7222f2cad93c2cdf13701ea1f903d47.zip | |
One core behind N frontends, the board's own runner moved in, and every board cap on one screen
## The wire is the effect stream, not a new protocol
`pardes --detach` leaves a core running with no terminal; `pardes --attach` is a frontend that owns
a terminal and a socket and nothing else. N frontends on one core all look at the same screen —
`screen -x`, not N sessions.
The codec (`src/detached/wire.zig`) carries exactly one `Event` or one `Host.VTable` call per
message. That is not a coincidence and it is why there is no third vocabulary to keep in step: the
core's IO seam was already a struct of function pointers with plain-data arguments, so a socket is
a legal implementation of it. `nested.zig`'s socket could not be reused — it carries a builtin
command line, and a command line cannot carry a frame.
ARCHITECTURE-NEUTRAL on purpose, not as decoration. The frontend on the far end may be
riscv32-freestanding on the ESP32-P4 while the core is x86_64 Linux, so every field is an explicit
little-endian fixed width and no message is a blit of a native struct. A protocol that only works
between two builds of the same compiler would have thrown away the one frontend that motivated it.
## The board comes in; its toolchain stays out
`src/p4.zig` becomes `src/esp32p4.zig`, and the pardes half of `../05-zig-p4` — the vaxis-over-
serial runner, the UART editor terminal, the keystroke rescue ring, the on-die test suite — moves
into `src/esp32p4/`. `build.zig.zon` gains `.zig_p4 = .{ .path = "../05-zig-p4" }`, so
`zig build -Dplatform=esp32p4 -Desp32p4-firmware` builds, flashes, monitors and self-tests the
board from this repo's `build.zig`.
The DIVISION is the point. What moved is what only pardes wants: the runner that drives a pardes
core over a serial line. What stayed is everything a second project would also want — the HAL, the
register/radio/oracle layers, the linker script, `_start`. `zig_p4` declares no dependencies of its
own and its `build()` early-returns when it is not the root package, so this costs the package
graph exactly zero packages and the editor's own builds nothing at all.
## limits.zig: nine forgettable places become one budget
Nine `platform == .esp32p4` capacity tests lived in nine files. They were never nine decisions —
they are ONE decision, how much memory this build may spend, taken nine times where no reader could
see the total. `src/limits.zig` puts the whole budget on one screen with every cap named against
what it is measured against, derived from two booleans.
The payoff is testability on a machine that is not the board: the caps are ordinary comptime values,
so a host build can be compiled against the board's numbers and the parking, eviction and clamping
paths a 240 KiB core takes get exercised by the normal test suite instead of only over a UART.
## A bare `zig build`
`zig build` with no arguments now builds the tty and GUI binaries and installs them into
`~/.local/bin`, and says so once on stdout with the flag that overrides it. The old default built
one binary into `zig-out` — a path nothing on a `PATH` ever looks at, which made "build it" and
"use it" two different commands for no reason.
Diffstat (limited to 'src/detached')
| -rw-r--r-- | src/detached/client.zig | 866 | ||||
| -rw-r--r-- | src/detached/server.zig | 1271 | ||||
| -rw-r--r-- | src/detached/wire.zig | 1755 |
3 files changed, 3892 insertions, 0 deletions
diff --git a/src/detached/client.zig b/src/detached/client.zig new file mode 100644 index 00000000..8d118782 --- /dev/null +++ b/src/detached/client.zig @@ -0,0 +1,866 @@ +//! THE FRONTEND SIDE of a detached session: a socket, a grid, and no core. +//! +//! THIS SIDE DOES NOT OWN A `Pardes`. That is the one thing to be clear about, +//! because the shape invites the mistake: there is no `update`, no `postEvent` +//! and no `Host` in this file. The core lives in the detached process +//! (server.zig), which is also where every `Host.VTable` call originates. A +//! frontend's whole job is two sentences long — send the input it collects, draw +//! the frames it is sent — and this module is exactly that and nothing more. +//! +//! `Client` is NOT a renderer either. It owns `grid` and `cursor`: the cells the +//! session is showing, kept current by applying frames as they arrive. Whoever +//! owns a terminal (or a window, or an ESP32-P4 panel) reads those and paints. +//! The split is deliberate — pardes already has terminal frontends, and a second +//! one living in here would be a second convention for a job that has one. +//! +//! THE ATTACH IS NOT A BLOCKING HANDSHAKE. `open` connects and writes the +//! `hello`; the `welcome` (or the `refuse`) arrives through the ordinary +//! `wait`/`next` loop like everything else. A frontend therefore has one loop +//! and one place it sleeps, instead of a startup path that can hang for two +//! seconds before its terminal is even in raw mode. `attached()` says whether +//! the session has greeted us; `refusal` says why it did not. +//! +//! WHAT ARRIVES, and what a frontend is expected to do with it. `next` hands +//! back one decoded `wire.ServerMsg` at a time: +//! * `welcome` and `frame` have already been applied to `slot`/`cols`/`rows`/ +//! `grid`/`cursor` by the time you see them. They are returned so a frontend +//! knows the screen moved. +//! * `spawn`, `pty_write`, `pty_resize`, `write_file`, `write_dump`, +//! `watch_file`, `watch_theme`, `dump_themes` are the session asking this +//! frontend for the real host work it has and a daemon does not: fork a +//! shell on a real tty, put bytes on a real disk, watch a path. A frontend +//! with none of that ignores them, exactly as a null vtable method does — +//! and only ONE attached frontend is ever asked (server.zig's routing). +//! * `set_clipboard`, `open_link` and `read_clipboard` are the desktop. The +//! answer to `read_clipboard` is not a reply message: it is an ordinary +//! `Event.paste` sent back through `send`, which is the same asynchronous +//! shape `pull_read_clipboard` already has in-process. +//! * `refuse` is followed by the session closing the connection, and `quit` +//! means the session itself has ended. +//! +//! GEOMETRY. `cols`/`rows` are the SESSION's grid, which with several frontends +//! attached is the smallest common one and can be smaller than this frontend's +//! window. Tell the session about the window with `resize`; do not assume the +//! next frame will agree with it. +//! +//! BORROWED BYTES. Every slice in a returned `ServerMsg` points into this +//! client's receive buffer and is valid until the next call to `next` or +//! `wait`. A frontend that needs a path or a payload for longer copies it — the +//! same rule the core's own `Event.output` bytes have. +const std = @import("std"); +const libc = std.c; +const pardes = @import("../pardes.zig"); +const server = @import("server.zig"); +const wire = @import("wire.zig"); + +const read_chunk = 16 * 1024; + +pub const Error = error{ + /// No `$XDG_RUNTIME_DIR` and no `$HOME`, or a name that is not one path + /// component — there is no socket path to try. + NoSessionPath, + /// Nothing is listening there: the name is wrong, or that session ended. + /// Its socket file, if it is still lying about, is unlinked by the sweep the + /// next detached session runs. + NoSession, + /// The socket is there and this frontend will not talk to it: the directory + /// or the socket is not a private one of ours (server.zig `vetted`). A + /// planted socket at a derivable path collects every keystroke typed into + /// the frontend that trusts it, so this is refused rather than reported as + /// "no session" — the two need different answers from a human. + NotPrivate, + /// The session hung up, or the connection failed under us. Every read and + /// write path funnels here: a frontend's answer to all of them is the same + /// (report and exit), so distinguishing them would be a distinction nobody + /// acts on. + Closed, + /// The session said something before it said `welcome`. + Ungreeted, +}; + +pub const Client = struct { + gpa: std.mem.Allocator, + fd: c_int = -1, + /// Which client slot the session gave this connection. Diagnostics only, + /// and it exists so both sides print the same number. + slot: u8 = 0, + /// Set by a `refuse`, and the reason a frontend prints before exiting. + refusal: ?wire.Refusal = null, + /// The SESSION's grid, not this frontend's window. Zero until the welcome. + cols: u16 = 0, + rows: u16 = 0, + /// `cols * rows` cells: what the session is showing right now. + grid: std.ArrayListUnmanaged(pardes.Cell) = .empty, + cursor: ?wire.Cursor = null, + in: std.ArrayListUnmanaged(u8) = .empty, + out: std.ArrayListUnmanaged(u8) = .empty, + /// Bytes of `in` belonging to the message `next` returned last. Compacted at + /// the top of the next call, which is exactly what makes that message's + /// borrowed slices valid until then and no longer. + held: usize = 0, + + /// Connect to the session called `name` and say hello, telling it this + /// frontend's window. Does not wait: the greeting arrives through `next`. + pub fn open(gpa: std.mem.Allocator, name: []const u8, cols: u16, rows: u16) (Error || wire.Error)!Client { + var path_buf: [server.path_max]u8 = undefined; + const path = server.sessionPath(&path_buf, name) orelse return error.NoSessionPath; + // Both ends vet, and this is this end's half: the session refuses a + // directory or a socket anyone else can reach before it binds, and until + // this a frontend connected to whatever it found at the path it derived. + // See server.zig `vetted` for what is asked and why it is asked there. + if (!server.vetted(path)) return error.NotPrivate; + 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 error.NoSession; + // CLOEXEC before anything can fork, and a frontend DOES fork: the pane + // shells it is asked to spawn are its own children, and one of them + // holding this socket would keep the session believing a frontend is + // attached long after this process left. + server.setCloexec(fd); + if (comptime server.darwin) { + // linux says MSG_NOSIGNAL per write and darwin says it once per + // socket, which is what `nosignal` being 0 on darwin MEANS — so + // this call is the whole of that platform's protection and the + // comment on `nosignal` used to say `open` made it without `open` + // making it. A session that ends mid-write must not take the + // frontend down with SIGPIPE. server.zig `accept` is the mirror. + const on: c_int = 1; + _ = libc.setsockopt(fd, libc.SOL.SOCKET, libc.SO.NOSIGPIPE, &on, @sizeOf(c_int)); + } + // Still blocking for the connect itself, which on AF_UNIX either lands + // in the listener's backlog immediately or is refused; there is no + // in-progress state to poll for. + if (libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) != 0) { + _ = libc.close(fd); + return error.NoSession; + } + server.setNonblock(fd); + var c: Client = .{ .gpa = gpa, .fd = fd }; + errdefer c.deinit(); + try c.send(.{ .hello = .{ .cols = cols, .rows = rows } }); + return c; + } + + /// Has the session greeted us? Until it has, `grid` is empty and nothing + /// has been drawn. + pub fn attached(c: *const Client) bool { + return c.cols != 0; + } + + pub fn deinit(c: *Client) void { + if (c.fd >= 0) _ = libc.close(c.fd); + c.fd = -1; + c.grid.deinit(c.gpa); + c.in.deinit(c.gpa); + c.out.deinit(c.gpa); + } + + /// Leave without ending the session. The `bye` is a courtesy — the session + /// handles a frontend that simply dies, and proving it does is what that + /// test is for — but it turns "the peer vanished" into "the peer left" in + /// the session's log, which is worth seven bytes. + pub fn detach(c: *Client) void { + c.send(.bye) catch {}; + c.deinit(); + } + + /// One message on its way to the core. Everything a frontend collects goes + /// through here: keys, the mouse, pty output from the shells it forked, a + /// paste answering a `read_clipboard`. + pub fn send(c: *Client, msg: wire.ClientMsg) (Error || wire.Error)!void { + const want = wire.clientBound(msg); + c.out.ensureUnusedCapacity(c.gpa, want) catch return error.Closed; + const at = c.out.items.len; + // Encoded straight into the queue's tail rather than through a scratch + // buffer: a paste is four megabytes and copying it twice is two copies. + c.out.items.len += want; + const bytes = wire.encodeClient(c.out.items[at..], msg) catch |err| { + // AND THE QUEUE GOES BACK. `want` bytes of it are uninitialised + // right now, and leaving them there — which is what a bare `try` + // did — puts that much stack-shaped garbage on the socket at the + // next flush: the session decodes it as a message, refuses it, and + // drops a frontend whose only mistake was a message this protocol + // cannot carry (an `Event.command` past 64 KiB is `Overlong`, and a + // caller that mis-sized the queue is `NoSpace`). + c.out.items.len = at; + return err; + }; + c.out.items.len = at + bytes.len; + return c.flush(); + } + + /// Tell the session this frontend's window changed. Not a promise about the + /// next frame: with other frontends attached the session grid is the + /// smallest common one. + pub fn resize(c: *Client, cols: u16, rows: u16) (Error || wire.Error)!void { + return c.send(.{ .event = .{ .resize = .{ .cols = cols, .rows = rows } } }); + } + + /// Wait up to `timeout_ms` for the session to say something, and push + /// whatever we still owe it. Zero blocks. This is the frontend's one + /// sleeping place FOR THE SOCKET; the terminal it draws on is polled by + /// whoever owns that, which is why this takes a timeout rather than a second + /// descriptor. + pub fn wait(c: *Client, timeout_ms: u32) (Error || wire.Error)!void { + if (c.fd < 0) return error.Closed; + try c.flush(); + var fds: [1]libc.pollfd = .{.{ + .fd = c.fd, + .events = if (c.out.items.len != 0) poll_in | poll_out else poll_in, + .revents = 0, + }}; + const timeout: c_int = if (timeout_ms == 0) -1 else @intCast(@min(timeout_ms, std.math.maxInt(c_int))); + if (libc.poll(&fds, 1, timeout) <= 0) return; // a timeout, or EINTR + if (fds[0].revents & poll_out != 0) try c.flush(); + // POLLIN wins over POLLHUP: a session that wrote a `quit` and then + // closed has bytes worth reading. + if (fds[0].revents & poll_in != 0) return c.fill(); + if (fds[0].revents & (poll_hup | poll_err | poll_nval) != 0) return error.Closed; + } + + /// The next complete message, or null when the buffer holds only part of + /// one. `welcome` and `frame` have already been applied to this client's + /// own state; every slice in the result borrows the receive buffer until the + /// next call here or to `wait`. + pub fn next(c: *Client) (Error || wire.Error)!?wire.ServerMsg { + // Retire the message returned last, now that its borrow window is over. + if (c.held != 0) { + if (c.held == c.in.items.len) { + c.in.clearRetainingCapacity(); + } else { + std.mem.copyForwards(u8, c.in.items, c.in.items[c.held..]); + c.in.items.len -= c.held; + } + c.held = 0; + } + const found = (try wire.framed(c.in.items)) orelse return null; + const msg = try wire.decodeServer(found.tag, found.payload); + c.held = found.total; + switch (msg) { + .welcome => |v| { + // Both sides check the version. This side checks it too because + // a session speaking something else may not have recognised our + // hello as one either, and a frontend must not paint a frame it + // decoded by a layout the other end does not use. + if (v.version != wire.version) return error.Ungreeted; + c.slot = v.slot; + try c.reshape(v.cols, v.rows); + }, + .refuse => |why| c.refusal = why, + .frame => |f| { + // A frame before the greeting would be the session drawing for a + // connection it never accepted. + if (!c.attached()) return error.Ungreeted; + // A geometry change and a late attach are the same case on this + // side too: reshape, and require the FULL frame the session + // promises for it. Applying a diff to a grid we just cleared + // would leave every untouched cell blank. + if (f.cols != c.cols or f.rows != c.rows) { + if (f.kind != .full) return error.BadValue; + try c.reshape(f.cols, f.rows); + } + try f.apply(c.grid.items); + c.cursor = f.cursor; + }, + // The session has ended. Left for the caller to act on, and the + // descriptor stays open so `deinit` is the only place that closes. + .quit => {}, + else => {}, + } + return msg; + } + + // ---- internals -------------------------------------------------------- + + fn reshape(c: *Client, cols: u16, rows: u16) (Error || wire.Error)!void { + c.grid.resize(c.gpa, @as(usize, cols) * @as(usize, rows)) catch return error.Closed; + // Unpainted, which a frontend draws as the terminal's own default cell. + // The full frame that follows paints over it. + @memset(c.grid.items, .{}); + c.cols = cols; + c.rows = rows; + c.cursor = null; + } + + /// Take everything the kernel is holding, not one chunk of it. The server + /// deliberately reads its clients ONE chunk per round, because it is + /// dividing a loop between thirty-two of them; a frontend has exactly one + /// peer, and reading one 16 KiB slice per poll would leave it a full frame + /// behind on every large one — and the session drops a frontend whose queue + /// it cannot drain (server.zig `out_backlog`). + fn fill(c: *Client) (Error || wire.Error)!void { + var buf: [read_chunk]u8 = undefined; + while (true) { + const got = libc.read(c.fd, &buf, buf.len); + if (got == 0) return error.Closed; + if (got < 0) return switch (libc.errno(got)) { + .INTR => continue, + // Nothing more is ready; what we have is what there was. + .AGAIN => {}, + else => error.Closed, + }; + c.in.appendSlice(c.gpa, buf[0..@intCast(got)]) catch return error.Closed; + } + } + + fn flush(c: *Client) Error!void { + if (c.fd < 0) return error.Closed; + var off: usize = 0; + while (off < c.out.items.len) { + const n = libc.send(c.fd, c.out.items.ptr + off, c.out.items.len - off, nosignal); + if (n < 0) switch (libc.errno(n)) { + .INTR => continue, + // The session is not draining us. The rest waits for POLLOUT; + // the queue is bounded in practice because a frontend's output + // is keystrokes and pty chunks, never frames. + .AGAIN => break, + else => return error.Closed, + }; + if (n == 0) break; + off += @intCast(n); + } + if (off == 0) return; + if (off == c.out.items.len) return c.out.clearRetainingCapacity(); + std.mem.copyForwards(u8, c.out.items, c.out.items[off..]); + c.out.items.len -= off; + } +}; + +// The socket primitives are server.zig's, which is the file that owns this +// transport's conventions and both ends of it — see its `setNonblock`, +// `nosignal` and `poll_*`. There was a third copy of all of them here. +const nosignal = server.nosignal; +const poll_in = server.poll_in; +const poll_out = server.poll_out; +const poll_hup = server.poll_hup; +const poll_err = server.poll_err; +const poll_nval = server.poll_nval; + +// --------------------------------------------------------------------------- +// tests +// --------------------------------------------------------------------------- +// +// One real core, one real unix socket, real frontends. The harness runs the +// server's vtable in the order `Pardes.pump` runs it, so what these exercise is +// the transport as the core actually drives it rather than a mock of it. + +const testing = std.testing; +const host_api = @import("../host.zig"); + +extern "c" fn setenv(name: [*:0]const u8, value: [*:0]const u8, overwrite: c_int) c_int; +extern "c" fn unsetenv(name: [*:0]const u8) c_int; + +/// A session on a socket of its own under `.zig-cache/tmp`, so a test never +/// collides with a real session in `$XDG_RUNTIME_DIR` and never depends on that +/// variable being set at all. It is process-wide, so it is saved and restored. +const Harness = struct { + tmp: std.testing.TmpDir, + saved: ?[:0]const u8, + saved_buf: [4096:0]u8 = undefined, + core: *pardes.Pardes, + session: server.Session, + arena: std.heap.ArenaAllocator, + name_buf: [32]u8 = undefined, + name: []const u8 = &.{}, + + fn init(h: *Harness, cols: u16, rows: u16) !void { + h.tmp = std.testing.tmpDir(.{}); + errdefer h.tmp.cleanup(); + h.saved = if (libc.getenv("XDG_RUNTIME_DIR")) |v| + try std.fmt.bufPrintSentinel(&h.saved_buf, "{s}", .{std.mem.span(v)}, 0) + else + null; + errdefer h.restoreEnv(); + var dir_buf: [4096:0]u8 = undefined; + const dir = try std.fmt.bufPrintSentinel(&dir_buf, ".zig-cache/tmp/{s}", .{h.tmp.sub_path}, 0); + // `ensureSocketDir` refuses anything with a bit granted to group or + // other, which is the whole point of it; a tmpDir arrives 0755. + try testing.expectEqual(@as(c_int, 0), libc.chmod(dir, 0o700)); + _ = setenv("XDG_RUNTIME_DIR", dir.ptr, 1); + h.core = try pardes.Pardes.init(testing.allocator, .{ .tty_only = true, .cols = cols, .rows = rows }); + errdefer h.core.deinit(); + h.arena = .init(testing.allocator); + errdefer h.arena.deinit(); + h.session = .{ .gpa = testing.allocator, .core = h.core, .cols = cols, .rows = rows }; + h.name = "s"; + try testing.expect(h.session.listen(h.name)); + // The pre-loop drain tty.zig has, for its reason: the startup spawns are + // already queued and a session must not open its socket with panes that + // have not been created. + h.core.host = h.session.host(); + while (h.core.nextEffect()) |e| h.core.perform(e); + } + + fn restoreEnv(h: *Harness) void { + if (h.saved) |v| { + _ = setenv("XDG_RUNTIME_DIR", v.ptr, 1); + } else _ = unsetenv("XDG_RUNTIME_DIR"); + } + + fn deinit(h: *Harness) void { + h.session.deinit(); + h.core.deinit(); + h.arena.deinit(); + h.restoreEnv(); + h.tmp.cleanup(); + } + + /// `Pardes.pump`, with the one substitution a single-threaded test needs: a + /// bounded wait, so a frontend that says nothing cannot hang the suite + /// where the real session would sleep until it spoke. The queued-input + /// drain `pump` does between the two is absent because this host has no + /// queue: it calls `update` directly (borrowed bytes), and `postEvent` is + /// for hosts whose worker threads post from off the loop. + fn pump(h: *Harness) !void { + const host = h.session.host(); + host.vtable.pull_wait_input.?(host.ctx, 20); + while (h.core.nextEffect()) |e| h.core.perform(e); + if (h.core.quit) return; + _ = h.arena.reset(.retain_capacity); + const surface = try h.core.render(h.arena.allocator()); + host.vtable.push_present.?(host.ctx, surface); + } + + /// Pump until this client has the message we are waiting for. Bounded, so a + /// broken transport fails a test rather than hanging the suite. + fn pumpUntil(h: *Harness, c: *Client, comptime want: std.meta.Tag(wire.ServerMsg)) !wire.ServerMsg { + for (0..64) |_| { + try h.pump(); + try c.wait(5); + while (try c.next()) |msg| if (std.meta.activeTag(msg) == want) return msg; + } + return error.NeverArrived; + } + + /// Pump until this client is sent a frame that CHANGES something. The first + /// frame after an input is not always the one carrying it — an effect the + /// input queued (a `Look` on a directory emits a spawn) lands a frame + /// later, and the frame in between legitimately says nothing. + fn pumpUntilChange(h: *Harness, c: *Client) !wire.Frame { + for (0..64) |_| { + const msg = try h.pumpUntil(c, .frame); + if (msg.frame.nruns > 0) return msg.frame; + } + return error.NothingChanged; + } + + /// Pump until this client's grid is this shape, draining everything that + /// arrives. A test waits for the STATE rather than for the n-th message + /// because a resize is announced when the session settles it, which may be + /// one empty frame after the pump that caused it. + fn pumpUntilGrid(h: *Harness, c: *Client, cols: u16, rows: u16) !void { + for (0..64) |_| { + try h.pump(); + try c.wait(5); + while (try c.next()) |_| {} + if (c.cols == cols and c.rows == rows) return; + } + return error.NeverResized; + } + + /// Pump until two frontends are showing the same screen, draining both on + /// every pass. One pump sends every attached frontend a frame, so a test + /// that drains only one of them is comparing two different instants — and + /// what a shared session promises is that they CONVERGE, which is what this + /// waits for. + fn pumpUntilSameScreen(h: *Harness, a: *Client, b: *Client) !void { + for (0..64) |_| { + try h.pump(); + try a.wait(5); + try b.wait(5); + while (try a.next()) |_| {} + while (try b.next()) |_| {} + if (sameScreen(a.grid.items, b.grid.items)) return; + } + return error.NeverConverged; + } + + /// Attach a frontend and get it greeted: `open` writes the hello into the + /// listener's backlog, one pump accepts and answers it. Nothing blocks, + /// which is the whole reason the handshake is not a blocking call. + fn attach(h: *Harness, cols: u16, rows: u16) !Client { + var c = try Client.open(testing.allocator, h.name, cols, rows); + errdefer c.deinit(); + _ = try h.pumpUntil(&c, .welcome); + return c; + } +}; + +test "detached session: a frontend attaches, is greeted, and is sent the screen" { + var h: Harness = undefined; + try h.init(60, 16); + defer h.deinit(); + + var c = try h.attach(60, 16); + defer c.deinit(); + try testing.expectEqual(@as(u8, 0), c.slot); + try testing.expectEqual(@as(u16, 60), c.cols); + try testing.expectEqual(@as(u16, 16), c.rows); + + const frame = (try h.pumpUntil(&c, .frame)).frame; + // The first frame a frontend gets must be full: it has nothing to diff + // against. + try testing.expectEqual(wire.FrameKind.full, frame.kind); + try testing.expectEqual(@as(usize, 60 * 16), c.grid.items.len); + // ...and it must be the core's own frame, cell for cell. This is the whole + // claim of the transport. + _ = h.arena.reset(.retain_capacity); + const surface = try h.core.render(h.arena.allocator()); + try expectSameScreen(surface.cells, c.grid.items); +} + +test "detached session: input from a frontend reaches the core and comes back as a diff" { + var h: Harness = undefined; + try h.init(60, 16); + defer h.deinit(); + var c = try h.attach(60, 16); + defer c.deinit(); + _ = try h.pumpUntil(&c, .frame); + + // A builtin line is the cheapest input with a guaranteed visible effect, + // and it travels the same `Event` path a keystroke does. + try c.send(.{ .event = .{ .command = "Look /" } }); + const frame = try h.pumpUntilChange(&c); + // The change arrived as a DIFF: this frontend was already in sync, so + // nothing it already had was re-sent. + try testing.expectEqual(wire.FrameKind.diff, frame.kind); + _ = h.arena.reset(.retain_capacity); + try expectSameScreen((try h.core.render(h.arena.allocator())).cells, c.grid.items); +} + +test "detached session: two frontends share one screen at the smallest common grid" { + var h: Harness = undefined; + try h.init(80, 24); + defer h.deinit(); + var a = try h.attach(80, 24); + defer a.deinit(); + _ = try h.pumpUntil(&a, .frame); + + // A second, smaller frontend. tmux's rule: the session shrinks to what both + // can show, because two people looking at different screens is the point of + // a shared session lost. + var b = try h.attach(50, 12); + defer b.deinit(); + try testing.expectEqual(@as(u8, 1), b.slot); + try testing.expectEqual(@as(u16, 50), b.cols); + try testing.expectEqual(@as(u16, 12), b.rows); + + // Both are now on the 50x12 grid — and getting there IS the proof that the + // reshape was sent as a FULL frame: `next` refuses a diff whose geometry + // does not match the grid it holds, so a client whose shape moved can only + // have been reshaped by a full one. + try h.pumpUntilGrid(&a, 50, 12); + try h.pumpUntilGrid(&b, 50, 12); + try testing.expectEqual(@as(usize, 50 * 12), a.grid.items.len); + try expectSameScreen(a.grid.items, b.grid.items); + + // ...and input from EITHER moves that one screen. Both are drained on every + // pump before they are compared: one pump sends every attached frontend a + // frame, so a test that drains only one is comparing two instants. + try b.send(.{ .event = .{ .command = "Look /" } }); + _ = try h.pumpUntilChange(&a); + try h.pumpUntilSameScreen(&a, &b); +} + +test "detached session: a frontend that dies takes nothing with it" { + var h: Harness = undefined; + try h.init(60, 16); + defer h.deinit(); + var a = try h.attach(60, 16); + defer a.deinit(); + var b = try h.attach(60, 16); + _ = try h.pumpUntil(&a, .frame); + _ = try h.pumpUntil(&b, .frame); + try testing.expect(h.session.clients[0].attached); + try testing.expect(h.session.clients[1].attached); + + // Not a `bye`: the socket goes away under the session's feet, which is what + // a frontend crashing or being killed looks like from here. + b.deinit(); + _ = try h.pumpUntil(&a, .frame); + try testing.expect(!h.session.clients[1].attached); + try testing.expectEqual(@as(c_int, -1), h.session.clients[1].fd); + // The survivor is still served and the core is still running. + try testing.expect(h.session.clients[0].attached); + try testing.expect(!h.core.quit); + try a.send(.{ .event = .{ .command = "Look /" } }); + try testing.expect((try h.pumpUntilChange(&a)).nruns > 0); + + // The freed slot takes the next frontend, and the screen comes with it. + var d = try h.attach(60, 16); + defer d.deinit(); + try testing.expectEqual(@as(u8, 1), d.slot); + try testing.expectEqual(wire.FrameKind.full, (try h.pumpUntil(&d, .frame)).frame.kind); + try expectSameScreen(a.grid.items, d.grid.items); +} + +test "detached session: a frontend speaking another protocol is refused, loudly" { + var h: Harness = undefined; + try h.init(60, 16); + defer h.deinit(); + + // A raw socket rather than a `Client`, because the whole point is a peer + // that does not agree with `wire.version` — and `open` would already have + // sent a perfectly good hello. + const fd = try rawConnect(&h); + defer _ = libc.close(fd); + var buf: [64]u8 = undefined; + const hello = try wire.encodeClient(&buf, .{ + .hello = .{ .version = wire.version + 1, .cols = 60, .rows = 16 }, + }); + try testing.expectEqual(@as(isize, @intCast(hello.len)), libc.send(fd, hello.ptr, hello.len, nosignal)); + + // Two pumps: one to accept the connection, one to read the hello and + // answer it. + try h.pump(); + try h.pump(); + var got: [64]u8 = undefined; + const n = libc.read(fd, &got, got.len); + try testing.expect(n > 0); + const f = (try wire.framed(got[0..@intCast(n)])).?; + try testing.expectEqual(wire.Refusal.version, (try wire.decodeServer(f.tag, f.payload)).refuse); + // Refused means refused: the slot went back and no frame was ever sent. + for (&h.session.clients) |*slot| try testing.expect(!slot.attached); + try testing.expect(!h.core.quit); +} + +test "detached session: a peer that sends garbage is dropped, not obeyed" { + var h: Harness = undefined; + try h.init(60, 16); + defer h.deinit(); + var c = try h.attach(60, 16); + defer c.deinit(); + _ = try h.pumpUntil(&c, .frame); + + // A well-formed frame around a tag this protocol has never defined. The + // session must close the connection rather than guess at it. + const junk = [_]u8{ 0xfe, 0x00, 0x00, 0x00, 0x00 }; + try testing.expectEqual(@as(isize, junk.len), libc.send(c.fd, &junk, junk.len, nosignal)); + for (0..8) |_| { + try h.pump(); + if (!h.session.clients[0].attached) break; + } + try testing.expect(!h.session.clients[0].attached); + // The session is untouched: a hostile frontend costs a slot, not a session. + try testing.expect(!h.core.quit); +} + +test "detached session: the session outlives every frontend and keeps its grid" { + var h: Harness = undefined; + try h.init(72, 20); + defer h.deinit(); + { + var c = try h.attach(40, 10); + defer c.detach(); + _ = try h.pumpUntil(&c, .frame); + try testing.expectEqual(@as(u16, 40), h.session.cols); + } + // Nobody attached. The grid stays where the last frontend left it rather + // than collapsing: a detached session is one nobody is looking at, not one + // of no size. + try h.pump(); + try testing.expectEqual(@as(u16, 40), h.session.cols); + try testing.expectEqual(@as(u16, 10), h.session.rows); + try testing.expect(!h.core.quit); + // ...and the next frontend takes the grid over: with one attachment the + // smallest common grid IS that frontend's, so the session follows it up to + // 90x30 rather than pinning the departed one's 40x10 forever. + var again = try h.attach(90, 30); + defer again.deinit(); + try testing.expectEqual(@as(u16, 90), again.cols); + try testing.expectEqual(@as(u16, 30), again.rows); + try h.pumpUntilGrid(&again, 90, 30); + try testing.expectEqual(@as(usize, 90 * 30), again.grid.items.len); +} + +test "detached session: the seam's own routing rules, per method" { + var h: Harness = undefined; + try h.init(60, 16); + defer h.deinit(); + var a = try h.attach(60, 16); + defer a.deinit(); + var b = try h.attach(60, 16); + defer b.deinit(); + _ = try h.pumpUntil(&a, .frame); + _ = try h.pumpUntil(&b, .frame); + + const host = h.session.host(); + // The eight `push_` methods with one real resource behind them go to the + // PRIMARY only — the oldest surviving attachment — because two frontends + // forking a shell for one pane gives that pane two shells. + host.vtable.push_spawn.?(host.ctx, 1, "/tmp"); + try expectOnly(&h, &a, &b, .spawn); + host.vtable.push_pty_write.?(host.ctx, 1, "ls\n"); + try expectOnly(&h, &a, &b, .pty_write); + host.vtable.push_write_file.?(host.ctx, 1, "/tmp/x", "body"); + try expectOnly(&h, &a, &b, .write_file); + + // ...and the ones that are facts about the SESSION go to everybody. + host.vtable.push_set_clipboard.?(host.ctx, "yank"); + try expectBoth(&h, &a, &b, .set_clipboard); + + // The one pull on the wire goes to the frontend whose input caused it, and + // it is asked ONCE — two frontends answering would paste twice for one + // Ctrl-V, which is the rule host.zig states. + h.session.origin = 1; + host.vtable.pull_read_clipboard.?(host.ctx); + try expectOnly(&h, &b, &a, .read_clipboard); + h.session.origin = 0; + host.vtable.push_open_link.?(host.ctx, "https://x"); + try expectOnly(&h, &a, &b, .open_link); +} + +test "detached session: the client table is a refusal, not a queue" { + // A small grid on purpose: these thirty-two peers never read, and a full + // frame of 80x24 each would push them into the backlog rule that the next + // test is about. 20x5 keeps every queue in the kernel's own buffer. + var h: Harness = undefined; + try h.init(20, 5); + defer h.deinit(); + + var fds: [server.max_clients]c_int = @splat(-1); + defer for (fds) |fd| if (fd >= 0) { + _ = libc.close(fd); + }; + var buf: [64]u8 = undefined; + const hello = try wire.encodeClient(&buf, .{ .hello = .{ .cols = 20, .rows = 5 } }); + for (&fds) |*fd| { + fd.* = try rawConnect(&h); + try testing.expectEqual(@as(isize, @intCast(hello.len)), libc.send(fd.*, hello.ptr, hello.len, nosignal)); + } + // The listener's backlog is `max_clients` deep, so this takes a few rounds: + // `accept` drains what is there each time it is woken. + for (0..16) |_| { + try h.pump(); + var attached: usize = 0; + for (&h.session.clients) |*slot| if (slot.attached) { + attached += 1; + }; + if (attached == server.max_clients) break; + } + for (&h.session.clients) |*slot| try testing.expect(slot.attached); + + // The thirty-third is TOLD it does not fit, and told promptly: the + // alternative — leaving it in the backlog — makes a level-triggered poll + // report the listener ready forever and spins the core. + const extra = try rawConnect(&h); + defer _ = libc.close(extra); + try testing.expectEqual(@as(isize, @intCast(hello.len)), libc.send(extra, hello.ptr, hello.len, nosignal)); + var got: [64]u8 = undefined; + const refusal = for (0..8) |_| { + try h.pump(); + const n = libc.read(extra, &got, got.len); + if (n > 0) break got[0..@intCast(n)]; + } else return error.NeverRefused; + const f = (try wire.framed(refusal)).?; + try testing.expectEqual(wire.Refusal.full, (try wire.decodeServer(f.tag, f.payload)).refuse); + // ...and the thirty-two it does serve are untouched. + for (&h.session.clients) |*slot| try testing.expect(slot.attached); + try testing.expect(!h.core.quit); +} + +test "detached session: a frontend that stops reading is dropped, not waited for" { + var h: Harness = undefined; + try h.init(40, 10); + defer h.deinit(); + var good = try h.attach(40, 10); + defer good.deinit(); + + // A peer that says hello and then never reads a byte again — a frontend + // stopped in a debugger, or one whose terminal is blocked. + const mute = try rawConnect(&h); + defer _ = libc.close(mute); + var buf: [64]u8 = undefined; + const hello = try wire.encodeClient(&buf, .{ .hello = .{ .cols = 40, .rows = 10 } }); + try testing.expectEqual(@as(isize, @intCast(hello.len)), libc.send(mute, hello.ptr, hello.len, nosignal)); + for (0..8) |_| { + try h.pump(); + if (h.session.clients[1].attached) break; + } + try testing.expect(h.session.clients[1].attached); + + // Broadcast enough control traffic to pass `out_backlog`. A clipboard + // mirror is the honest vehicle: it is a real `push_` that reaches every + // frontend and carries the yank register, so this is a session yanking a + // lot rather than a synthetic poke. + const text = try testing.allocator.alloc(u8, 256 * 1024); + defer testing.allocator.free(text); + @memset(text, 'y'); + const host = h.session.host(); + for (0..24) |_| { + if (!h.session.clients[1].attached) break; + host.vtable.push_set_clipboard.?(host.ctx, text); + try h.pump(); + try good.wait(5); + while (try good.next()) |_| {} + } + // Dropped rather than queued without bound, and rather than the core + // blocking on it. + try testing.expect(!h.session.clients[1].attached); + try testing.expectEqual(@as(c_int, -1), h.session.clients[1].fd); + // The frontend that WAS reading is still attached and still being drawn + // for, which is the whole claim: one slow peer costs its own slot. + try testing.expect(h.session.clients[0].attached); + try testing.expect(!h.core.quit); + try good.send(.{ .event = .{ .command = "Msg still here" } }); + try testing.expect((try h.pumpUntilChange(&good)).nruns > 0); +} + +/// A connected socket with nothing said on it yet, for the tests whose peer is +/// deliberately not a `Client`: one that speaks another protocol, thirty-two +/// that fill the table, one that never reads. +fn rawConnect(h: *Harness) !c_int { + var path_buf: [server.path_max]u8 = undefined; + const path = server.sessionPath(&path_buf, h.name).?; + 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 error.NoSocket; + errdefer _ = libc.close(fd); + if (libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) != 0) return error.NoSession; + return fd; +} + +/// `want` reached `to` and nothing reached `other`. +fn expectOnly( + h: *Harness, + to: *Client, + other: *Client, + comptime want: std.meta.Tag(wire.ServerMsg), +) !void { + _ = try h.pumpUntil(to, want); + try other.wait(5); + while (try other.next()) |msg| if (std.meta.activeTag(msg) == want) { + std.debug.print("{t} reached a frontend it was not routed to\n", .{want}); + return error.Misrouted; + }; +} + +fn expectBoth( + h: *Harness, + a: *Client, + b: *Client, + comptime want: std.meta.Tag(wire.ServerMsg), +) !void { + _ = try h.pumpUntil(a, want); + _ = try h.pumpUntil(b, want); +} + +/// Are these two grids showing the same thing? `visuallyEqual` and not +/// `std.meta.eql`, because `Cell.text` past `len` is scratch the decoder does +/// not invent — pardes.zig says in as many words that it must never +/// manufacture a difference. +fn sameScreen(want: []const pardes.Cell, have: []const pardes.Cell) bool { + if (want.len != have.len) return false; + for (want, have) |*x, *y| if (!x.visuallyEqual(y)) return false; + return true; +} + +fn expectSameScreen(want: []const pardes.Cell, have: []const pardes.Cell) !void { + try testing.expectEqual(want.len, have.len); + for (want, have, 0..) |*x, *y, i| if (!x.visuallyEqual(y)) { + std.debug.print("cell {d} differs\n", .{i}); + return error.CellMismatch; + }; +} diff --git a/src/detached/server.zig b/src/detached/server.zig new file mode 100644 index 00000000..6d13a4ae --- /dev/null +++ b/src/detached/server.zig @@ -0,0 +1,1271 @@ +//! THE DETACHED CORE: one `Pardes` instance in a process with no terminal, +//! serving N frontends over one unix socket. +//! +//! THIS SIDE OWNS THE CORE. `Session` is a `host.Host` implementation whose +//! methods encode wire messages instead of doing IO, and whose +//! `pull_wait_input` is a `poll(2)` over the listener and every attached +//! frontend. The frontends own terminals and nothing else (client.zig). So the +//! `Pardes` is here, `update` is called from here, and the same screen is on +//! every attached frontend at once — `screen -x`, not N sessions. +//! +//! WHAT THIS SIDE SERVES ITSELF. Every method this vtable leaves null falls +//! through to the core's own `host.Fallback`: the embedded source filesystem, +//! the in-process clipboard, silent ptys. host.zig says in as many words that a +//! zero-method host is a complete pardes, and that is exactly what a session +//! with nothing attached is. Everything a real frontend can do BETTER — fork a +//! shell on a real tty, put bytes on a real disk, reach a real desktop +//! clipboard — is asked of a frontend, and the routing table below says which. +//! +//! ROUTING, and it is not "push means broadcast". A push reaches every HOST +//! (host.zig's rule, which `Fanout.isPull` enforces); this is ONE host that +//! happens to be backed by several frontends, and how it spreads a call inside +//! itself is its own business. Three rules, one per kind of side effect: +//! * BROADCAST — the frame, and `set_clipboard`. Every screen must show the +//! same thing, and a yank in a shared session is a session-wide fact that +//! every attached desktop is entitled to. +//! * PRIMARY ONLY — `spawn`, `pty_write`, `pty_resize`, `write_file`, +//! `write_dump`, `watch_file`, `watch_theme`, `dump_themes`. Each of these +//! has ONE real resource behind it, and doing it twice is not doing it +//! twice as well: two frontends forking a shell for pane 3 gives the pane +//! two shells, and two frontends writing one path race each other. Primary +//! is the lowest attached slot, i.e. the oldest surviving attachment — a +//! rule that is stable while frontends come and go and needs no election. +//! A pane's shell therefore lives in the frontend that forked it: when that +//! frontend leaves, its panes stop producing output and the session's text, +//! files and layout carry on. That is a real limit and it is stated here +//! rather than papered over, because migrating a live pty between processes +//! is a different feature. +//! * ORIGIN, ELSE PRIMARY — `read_clipboard` (the one `pull_` on the wire) +//! and `open_link`. Both answer a thing a HUMAN just did, and the answer +//! belongs on that human's machine: the paste must come from the keyboard +//! that asked for it, and a link must open in front of the person who +//! clicked it. `origin` is the frontend whose event was applied most +//! recently. Effects drain after a whole batch of events (pardes.zig +//! `pump`), so in the rare case where two frontends type in the same +//! millisecond the second one wins; the fallback to primary covers an +//! effect that no input caused at all. +//! +//! FAIRNESS, and why no client can stall the core or another client: +//! * every descriptor is non-blocking, and there is no thread per client. One +//! `poll(2)` per pump covers the listener and all `max_clients` frontends. +//! * FRAMES ARE NOT QUEUED. A client with bytes still owed to the kernel is +//! SKIPPED for this frame and its mirror is left alone, so the next frame +//! it does get is a diff against what it actually has. A slow frontend +//! therefore sees fewer, larger frames instead of a growing queue, and +//! coalescing costs no byte surgery at all. +//! * what is left in a client's out-queue is control messages, and it is +//! capped (`out_backlog`). The cap is checked BEFORE an append, so a single +//! oversized message still goes out whole and what gets refused is a client +//! that has stopped draining: it is closed. Its session and its peers are +//! untouched, and it may reattach and be sent a full frame. +//! * `max_clients` is a REFUSAL, not a queue — the same shape and the same +//! number as fuse.zig's park table, and for the same reason: the listener +//! is always accepted from even when the table is full, because a +//! level-triggered `poll` on a backlog nobody accepts returns ready +//! forever and spins a core. Bounded per round all the same (`accept`), and +//! a connection that never says `hello` loses its slot +//! (`greet_deadline_ms`) — a slot held by silence is the same denial as a +//! queue, arrived at from the other end. +//! * the TABLE is accounted, not just each client (`session_backlog`), and a +//! drained client gives its buffers back (`idle_retain`): 32 slots each +//! holding one 4 MiB paste is 128 MiB of a daemon nobody is looking at. +//! +//! THE SOCKET follows nested.zig's conventions exactly, and they ARE +//! nested.zig's: `socketDir`, `ensureSocketDir`, `statNoFollow` and +//! `setCloexec` are imported from it rather than copied, because one directory +//! vetted by two predicates is how the two go out of step. `$XDG_RUNTIME_DIR` +//! else `~/.local/state/pardes` created 0700 and vetted (never /tmp), +//! `chmod 0600` before `listen(2)`, CLOEXEC on the listener and on every +//! accepted connection. The NAME differs on purpose: +//! `pardes-detached-<name>.sock` rather than `pardes-<pid>.sock`, so that +//! nested.zig's sweeper — which only recognises all-digit pids — never unlinks +//! a live detached session, and so that a person can say `--detach=work` +//! instead of learning a pid. +//! +//! WHO MAY BIND A NAME, and this side is not allowed to guess. `bind(2)` on a +//! unix socket is an atomic exclusive create, so it decides: a name whose +//! socket ANSWERS is a live session and `listen` refuses rather than taking it +//! (an unconditional unlink-before-bind is how a second `--detach=work` used +//! to steal the socket out from under every frontend attached to the first). +//! The only file this process unlinks is one it proved dead — a connect that +//! was REFUSED — and `alive` is the single place that judgement is made, for +//! `listen` and for the sweep both. +//! +//! ...and both ends do the vetting. `vetted` is the frontend's half: a socket +//! at a path anyone could plant receives every keystroke that frontend +//! collects, so the client checks the directory and the socket before it +//! connects, exactly as this side checks them before it binds. +const std = @import("std"); +const libc = std.c; +const pardes = @import("../pardes.zig"); +const host_api = @import("../host.zig"); +const wire = @import("wire.zig"); + +/// Diagnostics for whoever is running the daemon. Every one of these is a +/// `debug`, and the level is not a judgement about how bad the thing is: +/// main.zig's logFn drops this scope entirely unless PARDES_LOG is set, so what +/// decides whether a human sees it is that variable and not the level. Reaching +/// for `warn` instead would change exactly one thing — a TEST binary does not +/// go through logFn, and its stderr is the build runner's failure signal. +const log = std.log.scoped(.detached); + +/// nested.zig owns the socket conventions this file shares — the directory, +/// its vetting, the stat that will not follow a symlink, CLOEXEC — and its +/// module comment carries the reasoning for each. Imported and not copied: +/// see the module header. +const nested = @import("../nested.zig"); + +/// `pub` for client.zig, which needs the same platform answer for the same +/// reason: SIGPIPE is per-write on linux and per-socket on darwin. +pub const darwin = nested.darwin; + +/// Same two ingredients as nested.zig needs, minus the ancestor walk: unix +/// sockets and a per-user runtime directory. Anywhere else there is no detached +/// session and `listen` says so. +const supported = nested.supported; + +/// `sun_path` is 108 bytes on linux and 104 on darwin, taken from the struct so +/// that the buffers, the fit checks and the memcpy cannot disagree with the +/// kernel or with each other. +const sun_path_len = nested.sun_path_len; + +/// How many frontends may be attached at once. The number and the shape are +/// fuse.zig's park table: 32 slots, and overflow is a refusal rather than a +/// queue. A session with 32 frontends on it is not a session, it is a mistake, +/// and the 33rd gets told so instead of waiting in a backlog nobody drains. +pub const max_clients = 32; + +/// Bytes of un-drained CONTROL messages a client may owe before it is closed. +/// Frames are not in here (see the module header), so this is a backlog of +/// ACTIONS — spawns, clipboard mirrors, file writes — and a frontend that has +/// not taken 1 MiB of those has stopped reading its socket. Checked before an +/// append rather than after, so one oversized message is never the thing that +/// trips it. +const out_backlog = 1 << 20; + +/// One read per client per poll round (see `receive`). 16 KiB is two orders of +/// magnitude past a keystroke and small enough to sit on the loop's stack; a +/// 4 MiB paste arrives across several rounds, which is the point. +const read_chunk = 16 * 1024; + +/// Bytes of client traffic — every in-queue and out-queue together — this +/// session may hold before it starts closing the peers holding it. +/// `out_backlog` bounds ONE slot and this bounds the table, which is not the +/// same ceiling: 32 clients each a byte under their own cap is 32 MiB of a +/// daemon nobody is looking at. 4 MiB is one whole paste in flight plus every +/// frame queue a real session builds, and past it the fattest peer is the peer +/// that stopped reading. The mirrors are NOT in this number: a mirror is this +/// session's own bookkeeping for a client it chose to serve, not something a +/// peer can grow. +const session_backlog = 4 << 20; + +/// What a DRAINED client is allowed to keep. `in` grows to hold one whole +/// message, so a single 4 MiB paste otherwise leaves 4 MiB resident in that +/// slot for the life of the session — 128 MiB across a full table, for +/// something that happened once. Anything above one `read_chunk` is handed +/// back the moment the buffer empties, and the next message pays one +/// allocation for it; below that it is kept, so a session of keystrokes never +/// asks the allocator at all. +const idle_retain = read_chunk; + +/// How long the listener is left out of the poll set after an `accept` that +/// failed for a reason that persists (EMFILE above all). See `accept`: the +/// alternative was sleeping 100 ms inside the core. +const accept_pause_ms = 100; + +/// Why a client's connection ended. Only ever logged (`PARDES_LOG=1`), and +/// spelled out because "connection closed" is the one diagnostic that has never +/// helped anybody. +const Closed = enum { bye, peer, protocol, backlog, silent, write, read, oom, refused, quitting }; + +const Client = struct { + fd: c_int = -1, + /// The `hello` landed and was accepted. Before that the connection exists + /// but votes on nothing and is sent no frames: its geometry is unknown. + attached: bool = false, + /// A `welcome` is owed, and is sent once this round's geometry has settled + /// so the number in it is the one the next frame will use. + greet: bool = false, + /// This frontend's own window, as its last `hello`/`resize` said. One vote + /// in `reconcile`'s minimum, never the session's grid by itself. + cols: u16 = 0, + rows: u16 = 0, + /// Bytes read and not yet a whole message. + in: std.ArrayListUnmanaged(u8) = .empty, + /// Bytes owed to the kernel. + out: std.ArrayListUnmanaged(u8) = .empty, + /// What this client's grid holds, so the next frame can be a diff. Advanced + /// only when a frame is actually queued for it, which is what makes a + /// skipped frame correct rather than lost. + mirror: std.ArrayListUnmanaged(pardes.Cell) = .empty, + /// The next frame must be full: freshly attached, or the session geometry + /// moved under it. + need_full: bool = true, + /// Monotonic milliseconds at `accept`, and the only thing an un-greeted + /// connection is timed against. See `Session.greet_deadline_ms`. + accepted_ms: i64 = 0, +}; + +/// A pane's shell: which frontend was asked to fork it, and where. +/// +/// WHY THE SESSION REMEMBERS THIS. A spawn is the one primary-only call that +/// has to survive having no frontend to serve it. Every boot layout creates its +/// panes before the socket exists, so a `--detach` performs its startup spawns +/// with nobody attached — and dropping them meant a session that opened with +/// panes whose shells had never been forked, forever, in silence. So a spawn +/// with no primary is OWED, and asked of whoever attaches next. +/// +/// It is also what makes `pty_write` reach the right process. A pane's pty +/// lives in the frontend that forked it, which is not always the primary: A +/// attaches and forks the shells, B attaches, A leaves — the panes are re-owed +/// to B, and then C attaching into A's freed slot becomes primary while the +/// ptys are in B. Routing a pane's bytes by its OWNER rather than by the +/// primary is the difference between typing into a shell and typing into +/// nothing. +const Shell = struct { + /// The slot that was asked to fork this pane's shell, or null when nobody + /// has been. + owner: ?u8 = null, + /// A spawn owed to whoever attaches next: either it was never asked, or the + /// frontend holding it left and took the pty with it. + owed: bool = false, + /// Copied, because `push_spawn`'s `cwd` borrows the core's memory for the + /// length of that one call and this outlives it by definition. + cwd: std.ArrayListUnmanaged(u8) = .empty, +}; + +pub const Session = struct { + gpa: std.mem.Allocator, + core: *pardes.Pardes, + /// -1 when nothing is bound: an unsupported platform, or a bind that + /// failed. A session with no listener is a session nobody can attach to, + /// which still runs. + listener: c_int = -1, + /// The bound path, kept so teardown unlinks exactly what was created and + /// nothing else — guarded on the fd, like nested.zig's `unlisten`. + path_buf: [sun_path_len]u8 = undefined, + path_len: usize = 0, + clients: [max_clients]Client = @splat(.{}), + /// The session grid: the smallest common one across attached frontends. + /// Seeded from the core's own startup size so the first attach of an + /// identically sized frontend posts no resize at all. + cols: u16, + rows: u16, + /// Whose input was applied last, for the two calls that must go back to one + /// particular frontend. See the module header. + origin: ?u8 = null, + /// One encode buffer, reused. Grown to whatever the largest message so far + /// needed rather than sized from `wire.max_payload`, which would be 16 MiB + /// of resident memory for a session whose frames are six kilobytes. + scratch: std.ArrayListUnmanaged(u8) = .empty, + /// Where each pane's shell lives, and which spawns are still owed. See + /// `Shell`. + shells: [pardes.MAX_PANES]Shell = @splat(.{}), + /// How long a connection may stay silent before the session takes its slot + /// back. `Client.open` writes its `hello` in the same call that connects, + /// so a peer that has said nothing for five seconds is not a frontend that + /// was slow, and thirty-two of them used to fill the table and lock every + /// real frontend out with a `refuse .full`. + /// + /// A field rather than a constant for exactly one reason: the test for that + /// would otherwise have to sleep five seconds. Nothing else changes it. + greet_deadline_ms: u32 = 5_000, + /// Monotonic milliseconds until which the LISTENER is left out of the poll + /// set, because an `accept` failed for a reason that persists. See `accept`. + accept_paused_ms: i64 = 0, + + // ---- lifetime --------------------------------------------------------- + + pub fn deinit(s: *Session) void { + // Tell everyone the session is over before the socket disappears, so a + // frontend exits on a `quit` rather than on a read error whose meaning + // it has to guess. Best effort by construction: these descriptors are + // non-blocking, so a frontend that is not reading gets the EOF instead + // — which is a case it has to handle regardless. + for (&s.clients) |*c| if (c.attached) s.send(c, .quit); + for (&s.clients) |*c| if (c.fd >= 0) s.close(c, .quitting); + s.unlisten(); + s.scratch.deinit(s.gpa); + for (&s.shells) |*sh| sh.cwd.deinit(s.gpa); + } + + /// Bind and listen. False when there is no socket, and a session without + /// one is simply one nobody can attach to — the same posture nested.zig + /// takes, and for the same reason: a failed bind must not cost a launch. + pub fn listen(s: *Session, name: []const u8) bool { + if (comptime !supported) return false; + var dir_buf: [sun_path_len:0]u8 = undefined; + const dir = nested.socketDir(&dir_buf) orelse return false; + if (!nested.ensureSocketDir(dir)) return false; + sweep(dir); + const path = socketPath(&s.path_buf, dir, name) orelse return false; + var addr: libc.sockaddr.un = .{ .path = @splat(0) }; + @memcpy(addr.path[0 .. path.len + 1], path[0 .. path.len + 1]); + const fd = libc.socket(libc.AF.UNIX, libc.SOCK.STREAM, 0); + if (fd < 0) return false; + nested.setCloexec(fd); + // `bind` IS the exclusive create — it fails with EADDRINUSE the moment + // the path exists — so it, and nothing else, decides who owns a name. + // There is no unlink before it: unlinking unconditionally is how a + // second `pardes --detach=work` took the socket away from a live + // session, leaving every frontend attached to a file no new frontend + // could reach. + if (libc.bind(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) != 0) { + // The one case that is not a collision: a session killed rather + // than quit ran no teardown, so its file outlived it. `alive` is + // the only thing that may say so, and it says so only about a + // connect that was REFUSED. + if (alive(path)) { + log.debug("a detached session is already listening on {s}", .{path}); + _ = libc.close(fd); + return false; + } + _ = libc.unlink(path); + if (libc.bind(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) != 0) { + _ = libc.close(fd); + return false; + } + } + // Owner-only, and BEFORE listen(2), which is the first moment anyone + // could connect. The directory is already private; this is the second + // wall, and this socket carries keystrokes into a live editor. + _ = libc.chmod(path, 0o600); + // A backlog of max_clients: past that the kernel refuses the connect + // itself, which is the same answer `accept` would give. + if (libc.listen(fd, max_clients) != 0) { + _ = libc.close(fd); + return false; + } + setNonblock(fd); + s.listener = fd; + s.path_len = path.len; + return true; + } + + fn unlisten(s: *Session) void { + if (s.listener < 0) return; + _ = libc.close(s.listener); + s.listener = -1; + // Guarded on the fd, so a bind that FAILED cannot unlink a path this + // process never created. + var z: [sun_path_len:0]u8 = undefined; + @memcpy(z[0..s.path_len], s.path_buf[0..s.path_len]); + z[s.path_len] = 0; + _ = libc.unlink(z[0..s.path_len :0]); + } + + pub fn host(s: *Session) host_api.Host { + return .{ .ctx = s, .vtable = &vtable }; + } + + fn of(ctx: ?*anyopaque) *Session { + return @ptrCast(@alignCast(ctx.?)); + } + + /// Thirteen methods, and the seven that are missing are missing on purpose + /// — see wire.zig's header for each one's reason. `push_poll_frame` and + /// `push_post_present` carry no information a frame does not; the four + /// synchronous or dispatched pulls and `push_fs_reply` belong to whoever + /// owns the core, which is this process. + const vtable: host_api.Host.VTable = .{ + .pull_wait_input = waitInput, + .push_present = present, + .push_spawn = spawn, + .push_pty_write = ptyWrite, + .push_pty_resize = ptyResize, + .push_write_file = writeFile, + .push_write_dump = writeDump, + .push_watch_file = watchFile, + .push_watch_theme = watchTheme, + .push_dump_themes = dumpThemes, + .push_set_clipboard = setClipboard, + .pull_read_clipboard = readClipboard, + .push_open_link = openLink, + }; + + // ---- routing ---------------------------------------------------------- + + /// The oldest surviving attachment. No election and no state: slots are + /// filled lowest-first, so the lowest attached one is the oldest that is + /// still here. + fn primary(s: *Session) ?*Client { + for (&s.clients) |*c| if (c.attached) return c; + return null; + } + + /// ...and the frontend whose input we are answering, when there is one. + fn origins(s: *Session) ?*Client { + if (s.origin) |i| { + const c = &s.clients[i]; + if (c.attached) return c; + } + return s.primary(); + } + + /// The frontend holding pane `pane`'s pty, which is NOT the primary in + /// general — see `Shell`. Null when no frontend holds it, which is a pane + /// with no child: the core's own answer to that is silence, and so is this. + fn holder(s: *Session, pane: u8) ?*Client { + if (pane >= s.shells.len) return null; + const owner = s.shells[pane].owner orelse return null; + const c = &s.clients[owner]; + return if (c.attached) c else null; + } + + /// Which slot this client is. From the pointer because every caller here + /// holds a `*Client` and not its index. + fn slotOf(s: *Session, c: *const Client) u8 { + return @intCast(@divExact(@intFromPtr(c) - @intFromPtr(&s.clients[0]), @sizeOf(Client))); + } + + fn broadcast(s: *Session, msg: wire.ServerMsg) void { + for (&s.clients) |*c| if (c.attached) s.send(c, msg); + } + + // ---- the host methods ------------------------------------------------- + + fn spawn(ctx: ?*anyopaque, pane: u8, cwd: []const u8) void { + const s = of(ctx); + if (pane >= s.shells.len) return; // the core indexes its own panes + const sh = &s.shells[pane]; + // Kept whether or not there is somebody to ask, because the pane now + // exists either way and the cwd is the only thing that cannot be + // reconstructed later. + sh.cwd.clearRetainingCapacity(); + sh.cwd.appendSlice(s.gpa, cwd) catch {}; + if (s.primary()) |c| { + sh.owner = s.slotOf(c); + sh.owed = false; + return s.send(c, .{ .spawn = .{ .pane = pane, .cwd = cwd } }); + } + // Nobody can fork a shell right now — a startup layout, or every + // frontend gone. NOT dropped: `reconcile` asks the next arrival. + sh.owner = null; + sh.owed = true; + } + + fn ptyWrite(ctx: ?*anyopaque, pane: u8, bytes: []const u8) void { + const s = of(ctx); + // To the frontend that forked this pane's shell, not to the primary: + // the pty is in that process and nowhere else. A pane whose holder is + // gone is silent, which is what the core does with a null method. + if (s.holder(pane)) |c| s.send(c, .{ .pty_write = .{ .pane = pane, .bytes = bytes } }); + } + + fn ptyResize(ctx: ?*anyopaque, pane: u8, cols: u16, rows: u16) void { + const s = of(ctx); + if (s.holder(pane)) |c| s.send(c, .{ .pty_resize = .{ .pane = pane, .cols = cols, .rows = rows } }); + } + + fn writeFile(ctx: ?*anyopaque, pane: u8, path: []const u8, bytes: []const u8) void { + const s = of(ctx); + if (s.primary()) |c| return s.send(c, .{ .write_file = .{ .pane = pane, .path = path, .bytes = bytes } }); + // Nobody attached, and a Put must not evaporate. This method being + // non-null means the core did NOT reach for its own filesystem, so the + // obligation a null method would have discharged is discharged here by + // hand — the same shape `readClipboard` below has, and host.zig's rule + // that a zero-method host is a complete pardes. + s.core.fallback.writeFile(path, bytes); + } + + fn writeDump(ctx: ?*anyopaque, bytes: []const u8) void { + const s = of(ctx); + if (s.primary()) |c| return s.send(c, .{ .write_dump = bytes }); + // ...and the same for a Dump, including the part that makes the bytes + // reachable again: a real host reports where it landed, which is what + // puts `Restore <path>` in the topbar (pardes.zig `write_dump`). + s.core.fallback.writeFile(pardes.fallback_dump_path, bytes); + s.core.setLastDump(pardes.fallback_dump_path); + } + + fn watchFile(ctx: ?*anyopaque, pane: u8, path: []const u8, on: bool) void { + const s = of(ctx); + if (s.primary()) |c| return s.send(c, .{ .watch_file = .{ .pane = pane, .path = path, .on = on } }); + // No frontend to watch a path, so the core's own record of what was + // asked is the whole of what a watch means here — exactly what a null + // method leaves behind. + if (pane < s.core.fallback.watched.len) s.core.fallback.watched[pane] = on; + } + + fn watchTheme(ctx: ?*anyopaque, generation: u32, on: bool) void { + const s = of(ctx); + // Dropped with nobody attached, and that is the whole of it: the core's + // own `theme_file` effect does nothing for a null method either, so + // there is no obligation left over. Same for `dump_themes` below. + if (s.primary()) |c| s.send(c, .{ .watch_theme = .{ .generation = generation, .on = on } }); + } + + fn dumpThemes(ctx: ?*anyopaque, pane: u8) void { + const s = of(ctx); + if (s.primary()) |c| s.send(c, .{ .dump_themes = .{ .pane = pane } }); + } + + fn setClipboard(ctx: ?*anyopaque, text: []const u8) void { + const s = of(ctx); + // Mirrored into the core's own clipboard ALWAYS, not only when nobody + // is attached: `readClipboard` answers from it when there is no + // frontend, and a frontend can leave between the yank and the paste. A + // yank that a detached session then pasted as the previous yank is the + // bug this one line is. + s.core.fallback.setClipboard(text); + s.broadcast(.{ .set_clipboard = text }); + } + + /// The one `pull_` that crosses the wire, and it stays a pull for exactly + /// the reason host.zig gives: two frontends answering would paste the + /// clipboard twice for one Ctrl-V. + fn readClipboard(ctx: ?*anyopaque) void { + const s = of(ctx); + if (s.origins()) |c| return s.send(c, .read_clipboard); + // Nobody attached. This method being non-null means the core will NOT + // reach for its own fallback, so an unanswered request would leave + // `clip_pending` armed forever — host.zig's note that a null method + // answers immediately is the obligation being met here by hand. + s.core.update(.{ .paste = s.core.fallback.clipboard.items }); + } + + fn openLink(ctx: ?*anyopaque, url: []const u8) void { + const s = of(ctx); + if (s.origins()) |c| return s.send(c, .{ .open_link = url }); + // No desktop in reach, so the link goes where a host with no browser + // puts it: the core's record of the last one asked for, which is what + // `Fallback.setLink` is and what the acme filesystem reads back. + s.core.fallback.setLink(url); + } + + // ---- the frame -------------------------------------------------------- + + fn present(ctx: ?*anyopaque, surface: *const pardes.Surface) void { + const s = of(ctx); + for (&s.clients) |*c| { + if (!c.attached) continue; + // A client that has not drained what it already owes does not get + // this frame, and its mirror is deliberately left where it is: the + // next frame it does get is a diff against what it really has. A + // slow frontend gets fewer, larger frames rather than a queue. + if (c.out.items.len != 0) continue; + s.sendFrame(c, surface); + } + } + + fn sendFrame(s: *Session, c: *Client, surface: *const pardes.Surface) void { + const cells = surface.cells; + const want = wire.frameBound(surface.cols, surface.rows); + s.scratch.ensureTotalCapacity(s.gpa, want) catch return s.close(c, .oom); + // Nothing comparable on the far side is the LATE JOINER and the RESIZE + // in one test: either way the whole grid has to be described. + const prev: []const pardes.Cell = if (c.need_full or c.mirror.items.len != cells.len) + &.{} + else + c.mirror.items; + const cursor: ?wire.Cursor = if (surface.cursor) |cur| + .{ .x = cur.x, .y = cur.y, .bar = cur.bar } + else + null; + const bytes = wire.encodeFrame( + s.scratch.allocatedSlice()[0..want], + surface.cols, + surface.rows, + cursor, + cells, + prev, + ) catch |err| { + // A frame this protocol cannot carry is a grid past `max_cols` / + // `max_rows`, or a cursor the core placed outside its own surface. + // Dropping the frame keeps the session alive with a stale screen, + // which is strictly better than dropping the frontend — a frontend + // REFUSES such a frame and hangs up — and the log says which. + log.debug("frame {d}x{d} not encodable: {t}", .{ surface.cols, surface.rows, err }); + return; + }; + s.queue(c, bytes); + if (c.fd < 0) return; // the queue closed it; the mirror went with it + // The mirror advances only now, and only because the bytes are on the + // wire or in the kernel's buffer for it. + c.mirror.resize(s.gpa, cells.len) catch return s.close(c, .oom); + @memcpy(c.mirror.items, cells); + c.need_full = false; + } + + // ---- the loop --------------------------------------------------------- + + /// The only place this process sleeps, which is what `pull_wait_input`'s + /// comment in host.zig requires of whoever serves it. One `poll(2)` covers + /// the listener and every attached frontend; there is no thread per client + /// and nothing here blocks on a single peer. + fn waitInput(ctx: ?*anyopaque, timeout_ms: u32) void { + const s = of(ctx); + // Push what the kernel will take before sleeping: a client that becomes + // writable while we are inside poll(2) would otherwise be a frame late, + // and a frame late is a frame skipped (see `present`). + for (&s.clients) |*c| if (c.fd >= 0) s.flush(c); + + const now = monotonicMs(); + var fds: [max_clients + 1]libc.pollfd = undefined; + var slots: [max_clients + 1]u8 = undefined; + var n: usize = 0; + // The listener is left OUT of the set while accepting is paused, which + // is how an EMFILE is waited out without the core sleeping (see + // `accept`). Every frontend already attached goes on being served. + const watching_listener = s.listener >= 0 and now >= s.accept_paused_ms; + if (watching_listener) { + fds[n] = .{ .fd = s.listener, .events = poll_in, .revents = 0 }; + slots[n] = 0; + n += 1; + } + for (&s.clients, 0..) |*c, i| { + if (c.fd < 0) continue; + fds[n] = .{ + .fd = c.fd, + .events = if (c.out.items.len != 0) poll_in | poll_out else poll_in, + .revents = 0, + }; + slots[n] = @intCast(i); + n += 1; + } + // A detached session with no listener and no clients has no event + // source at all. Returning immediately would spin the outer + // `while (!core.quit)` at full speed, so sleep the interval the core + // offered and, when it offered none, a frame's worth. + if (n == 0) return nap(if (timeout_ms == 0) 16 else timeout_ms); + // Zero is the core's word for "sleep until something happens" (see + // pardes.zig `pump`: it passes a frame interval only while an animation + // is running). poll spells that -1. + var timeout: c_int = if (timeout_ms == 0) -1 else @intCast(@min(timeout_ms, std.math.maxInt(c_int))); + // Two things here are due on a CLOCK rather than on a descriptor: a + // handshake that has to expire, and a paused listener that has to come + // back. An indefinite poll would sit through both — and thirty-two + // peers that connect and then say nothing, with the session otherwise + // idle, IS the denial `greet_deadline_ms` exists to answer — so the + // wait is clamped to whichever is due first. + if (s.nextWake(now)) |due| timeout = if (timeout < 0) due else @min(timeout, due); + const ready = libc.poll(&fds, @intCast(n), timeout); + // Expired unconditionally: a slot held by silence comes back on a + // timeout exactly as it does on a wakeup, and a poll that returned + // nothing is the ordinary way this deadline is reached. + s.expire(monotonicMs()); + // A timeout is an ordinary frame boundary and EINTR is a signal we do + // not handle here; both simply come back next pump. + if (ready <= 0) return; + + var k: usize = 0; + if (watching_listener) { + if (fds[0].revents != 0) s.accept(); + k = 1; + } + while (k < n) : (k += 1) { + const c = &s.clients[slots[k]]; + // A slot closed earlier in this same pass (its peer hung up, a + // decode failed, its handshake expired) must not be touched + // through a stale revents. + if (c.fd < 0) continue; + if (fds[k].revents & poll_out != 0) s.flush(c); + if (c.fd < 0) continue; + if (fds[k].revents & poll_in != 0) { + s.receive(c, slots[k]); + } else if (fds[k].revents & (poll_hup | poll_err | poll_nval) != 0) { + // POLLIN wins when both are set: a peer that wrote and then + // closed has bytes still worth reading. + s.close(c, .peer); + } + } + s.reconcile(); + } + + /// Milliseconds until the next deadline that is kept by the CLOCK rather + /// than by a descriptor, or null when there is none. Floored at zero, so a + /// deadline already past polls once without blocking instead of blocking + /// forever on a negative timeout. + fn nextWake(s: *const Session, now: i64) ?c_int { + if (now == 0) return null; // no clock; see `monotonicMs` + var due: ?i64 = null; + for (&s.clients) |*c| { + if (c.fd < 0 or c.attached) continue; + const at = c.accepted_ms + @as(i64, s.greet_deadline_ms); + due = if (due) |d| @min(d, at) else at; + } + if (s.listener >= 0 and s.accept_paused_ms > now) + due = if (due) |d| @min(d, s.accept_paused_ms) else s.accept_paused_ms; + const at = due orelse return null; + return @intCast(@max(0, @min(at - now, std.math.maxInt(c_int)))); + } + + /// Take the slots of connections that never said `hello` back. A connection + /// that holds a slot in silence denies a real frontend exactly as a queue + /// would, and `Client.open` writes its hello in the same call that + /// connects, so there is nothing legitimate to wait for. + fn expire(s: *Session, now: i64) void { + if (now == 0) return; // no clock: enforce nothing rather than everything + for (&s.clients) |*c| { + if (c.fd < 0 or c.attached) continue; + if (now - c.accepted_ms >= s.greet_deadline_ms) s.close(c, .silent); + } + } + + /// Always accept, even with a full table: the tempting alternative — stop + /// accepting and let the kernel hold the surplus — is a spin, because + /// `poll` is level triggered and an unaccepted backlog reports ready + /// forever. fuse.zig's park table learned that as a deadlock; here it is + /// 100% of a core. + /// + /// BOUNDED all the same. `max_clients + 1` is enough to fill an empty table + /// and refuse one more, and past that the surplus waits in the backlog for + /// the next round — one pump later, with every frontend drawn in between. + /// The `while (true)` this replaces let a peer dialling in a loop hold the + /// core inside `accept` for as long as it kept dialling, and the core is + /// what draws every other frontend's screen. + fn accept(s: *Session) void { + for (0..max_clients + 1) |_| { + const fd = libc.accept(s.listener, null, null); + if (fd < 0) { + switch (libc.errno(fd)) { + // The backlog is empty, which is this loop's ordinary exit. + .AGAIN, .INTR, .CONNABORTED => return, + // Anything else — EMFILE above all — persists until some + // other descriptor is freed, and `poll` is LEVEL + // triggered: coming straight back means poll reports the + // listener ready again immediately and the core spins at + // 100% until the condition clears. The old answer was a + // 100 ms nanosleep, which parks the CORE — every attached + // frontend stops being drawn for a tenth of a second + // because a descriptor ran out. So the LISTENER is dropped + // from the poll set for that beat instead, and the session + // goes on serving the frontends it has. + else => { + s.accept_paused_ms = monotonicMs() + accept_pause_ms; + return; + }, + } + } + nested.setCloexec(fd); + setNonblock(fd); + if (comptime darwin) { + // linux says MSG_NOSIGNAL per write; darwin says it once per + // socket. Either way a frontend that dies mid-frame must not + // take the session down with SIGPIPE. + const on: c_int = 1; + _ = libc.setsockopt(fd, libc.SOL.SOCKET, libc.SO.NOSIGPIPE, &on, @sizeOf(c_int)); + } + const slot = for (&s.clients, 0..) |*c, i| { + if (c.fd < 0) break i; + } else { + // Refused, and told why, on a connection accepted purely so + // that the listener stays quiet. + s.refuseFd(fd, .full); + _ = libc.close(fd); + continue; + }; + s.clients[slot] = .{ .fd = fd, .accepted_ms = monotonicMs() }; + } + } + + /// One read per client per round. A frontend that never stops talking gets + /// one turn and then the loop moves on to the others and to the frame — + /// which is fuse.zig's `retry` rule (one attempt per parked request per + /// frame) applied to sockets. + fn receive(s: *Session, c: *Client, slot: u8) void { + var buf: [read_chunk]u8 = undefined; + const got = libc.read(c.fd, &buf, buf.len); + if (got == 0) return s.close(c, .peer); // clean EOF: the frontend left + if (got < 0) return switch (libc.errno(got)) { + .INTR, .AGAIN => {}, + else => s.close(c, .read), + }; + c.in.appendSlice(s.gpa, buf[0..@intCast(got)]) catch return s.close(c, .oom); + // The table's own ceiling, checked where the table grows: a peer that + // sends the first half of a 16 MiB message and stops is holding memory + // no per-message check can see. See `session_backlog`. + s.account(); + if (c.fd < 0) return; // it was this one + s.consume(c, slot); + } + + fn consume(s: *Session, c: *Client, slot: u8) void { + var off: usize = 0; + while (true) { + const found = wire.framed(c.in.items[off..]) catch return s.close(c, .protocol); + const msg = found orelse break; + // The decoded Event BORROWS these bytes, so the buffer is not + // compacted until every message already in it has been applied — + // the same borrow window the tty host gives a pty chunk. + s.apply(c, slot, msg.tag, msg.payload) catch return s.close(c, .protocol); + if (c.fd < 0) return; // apply closed it, buffers and all + off += msg.total; + } + if (off == 0) return; + if (off == c.in.items.len) { + c.in.clearRetainingCapacity(); + return retire(s.gpa, &c.in); + } + std.mem.copyForwards(u8, c.in.items, c.in.items[off..]); + c.in.items.len -= off; + } + + fn apply(s: *Session, c: *Client, slot: u8, tag: u8, payload: []const u8) wire.Error!void { + var scratch: wire.Scratch = .{}; + switch (try wire.decodeClient(tag, payload, &scratch)) { + .hello => |h| { + // A second hello on one connection is not a resize; it is a + // peer that is not speaking this protocol. + if (c.attached) return error.BadValue; + if (h.version != wire.version) { + log.debug("frontend speaks protocol {d}, this session speaks {d}", .{ h.version, wire.version }); + return s.refuse(c, .version); + } + if (s.core.quit) return s.refuse(c, .quitting); + c.cols = h.cols; + c.rows = h.rows; + c.attached = true; + c.need_full = true; + // Greeted after `reconcile`, so the geometry in the welcome is + // the one this client's first frame will actually use. + c.greet = true; + }, + .bye => s.close(c, .bye), + .event => |ev| { + // Input before a handshake has no geometry behind it and no + // version agreement either. + if (!c.attached) return error.BadValue; + switch (ev) { + // A frontend's resize is about ITS window. The core only + // ever sees the smallest common grid, which `reconcile` + // posts once per round when it moves — forwarding this raw + // would let whichever frontend resized last win. + .resize => |r| { + c.cols = r.cols; + c.rows = r.rows; + }, + else => { + s.origin = slot; + s.core.update(ev); + }, + } + }, + } + } + + /// Settle the session grid and greet whoever arrived, once per poll round + /// rather than once per message: three frontends attaching in the same + /// round are one resize, not three reflows of every pane. + fn reconcile(s: *Session) void { + var cols: u16 = 0; + var rows: u16 = 0; + for (&s.clients) |*c| { + if (!c.attached) continue; + cols = if (cols == 0) c.cols else @min(cols, c.cols); + rows = if (rows == 0) c.rows else @min(rows, c.rows); + } + // Nobody attached: keep the grid we had. A detached session is not a + // session of no size, it is one nobody is looking at, and reflowing + // every pane to nothing for zero readers is work with no reader. + if (cols != 0 and (cols != s.cols or rows != s.rows)) { + s.cols = cols; + s.rows = rows; + // Every mirror is now the wrong shape. `encodeFrame` reaches the + // same conclusion from the cell count alone, but saying it here is + // what makes a reshape with the SAME cell count (80x24 -> 48x40) + // safe too. + for (&s.clients) |*c| c.need_full = true; + s.core.update(.{ .resize = .{ .cols = cols, .rows = rows } }); + } + for (&s.clients, 0..) |*c, i| { + if (!c.greet) continue; + c.greet = false; + s.send(c, .{ .welcome = .{ .slot = @intCast(i), .cols = s.cols, .rows = s.rows } }); + } + s.flushOwed(); + } + + /// Hand every owed spawn to the frontend that can serve it. Runs at the end + /// of a poll round, so a frontend that has just been greeted is asked for + /// its panes' shells in the same round it arrived — and a session that was + /// started with panes and no frontend (which is every `--detach`) is a + /// session whose panes get their shells from the first attach rather than + /// never. See `Shell`. + fn flushOwed(s: *Session) void { + const c = s.primary() orelse return; + const slot = s.slotOf(c); + for (&s.shells, 0..) |*sh, pane| { + if (!sh.owed) continue; + sh.owed = false; + sh.owner = slot; + s.send(c, .{ .spawn = .{ .pane = @intCast(pane), .cwd = sh.cwd.items } }); + // The send closed it, and `close` put its panes back on the owed + // list; the ones this loop has not reached are still owed anyway. + if (c.fd < 0) return; + } + } + + // ---- bytes ------------------------------------------------------------ + + fn send(s: *Session, c: *Client, msg: wire.ServerMsg) void { + const want = wire.serverBound(msg); + s.scratch.ensureTotalCapacity(s.gpa, want) catch return s.close(c, .oom); + const bytes = wire.encodeServer(s.scratch.allocatedSlice()[0..want], msg) catch |err| { + // The only reachable case is a payload past `max_payload`: a save + // of a pane holding more text than this protocol carries. The + // session keeps it (the core's own filesystem already has it) and + // the frontend's copy does not happen — said out loud rather than + // silently. + log.debug("message {t} not encodable: {t}", .{ msg, err }); + return; + }; + s.queue(c, bytes); + } + + fn queue(s: *Session, c: *Client, bytes: []const u8) void { + // BEFORE the append, so one oversized message always goes out whole and + // what this refuses is a client that has stopped draining. + if (c.out.items.len > out_backlog) return s.close(c, .backlog); + // ...and the table as a whole, which `out_backlog` does not bound: 32 + // slots one byte under it each is 32 MiB. See `session_backlog`. + s.account(); + if (c.fd < 0) return; // the fattest peer was this one + c.out.appendSlice(s.gpa, bytes) catch return s.close(c, .oom); + // Try immediately: on a local socket this empties the queue in one + // write, and `present` skips a client whose queue is not empty. + s.flush(c); + } + + /// Close the peer holding the most of the table when the table as a whole + /// is over `session_backlog`. One peer per call, and the fattest one, + /// because this is only ever asked when the total is already over and the + /// peer holding the most of it is the peer that stopped reading. The next + /// append asks again, so a second offender is closed a message later rather + /// than in a loop that could empty the table on one bad frame. + fn account(s: *Session) void { + var total: usize = 0; + var worst: ?*Client = null; + var worst_bytes: usize = 0; + for (&s.clients) |*c| { + if (c.fd < 0) continue; + const held = c.in.capacity + c.out.capacity; + total += held; + if (held > worst_bytes) { + worst_bytes = held; + worst = c; + } + } + if (total <= session_backlog) return; + if (worst) |c| s.close(c, .backlog); + } + + fn flush(s: *Session, c: *Client) void { + var off: usize = 0; + while (off < c.out.items.len) { + const n = libc.send(c.fd, c.out.items.ptr + off, c.out.items.len - off, nosignal); + if (n < 0) switch (libc.errno(n)) { + .INTR => continue, + // The kernel's buffer is full: the rest waits for POLLOUT, and + // this client is skipped for frames until it drains. + .AGAIN => break, + else => return s.close(c, .write), + }; + if (n == 0) break; + off += @intCast(n); + } + if (off == 0) return; + if (off == c.out.items.len) { + c.out.clearRetainingCapacity(); + return retire(s.gpa, &c.out); + } + std.mem.copyForwards(u8, c.out.items, c.out.items[off..]); + c.out.items.len -= off; + } + + /// Say why, then hang up. The refusal is written with a plain blocking + /// write on a socket nobody has sent anything on yet: it is six bytes, and + /// queueing it would mean keeping a slot for a connection being rejected. + fn refuse(s: *Session, c: *Client, why: wire.Refusal) void { + s.refuseFd(c.fd, why); + s.close(c, .refused); + } + + /// Writes only. The descriptor belongs to the caller — `refuse` hands it to + /// `close`, and the full-table path in `accept` closes it itself — because + /// closing here as well is a double close, and the number is reusable the + /// instant the first one lands. + fn refuseFd(_: *Session, fd: c_int, why: wire.Refusal) void { + var buf: [wire.header_len + 1]u8 = undefined; + const bytes = wire.encodeServer(&buf, .{ .refuse = why }) catch unreachable; + var off: usize = 0; + while (off < bytes.len) { + const n = libc.send(fd, bytes.ptr + off, bytes.len - off, nosignal); + if (n < 0 and libc.errno(n) == .INTR) continue; + if (n <= 0) break; // it left before hearing why; nothing to do + off += @intCast(n); + } + } + + /// Free one slot. A frontend dying takes NOTHING with it: not the core, not + /// the listener, not another frontend's frames. Its buffers go back and the + /// slot is reusable on the next connect. + fn close(s: *Session, c: *Client, why: Closed) void { + if (c.fd < 0) return; + log.debug("frontend detached: {t}", .{why}); + _ = libc.close(c.fd); + c.in.deinit(s.gpa); + c.out.deinit(s.gpa); + c.mirror.deinit(s.gpa); + const gone = s.slotOf(c); + // Which slot this is, so a departing frontend cannot leave `origin` + // pointing at it and send the next `read_clipboard` to a stranger. + if (s.origin) |i| if (i == gone) { + s.origin = null; + }; + // ...and its panes' shells died with the process that forked them. They + // go back on the owed list, so the frontend that replaces this one is + // asked to fork them again in the directory they were forked in: the + // alternative — which is what this did — is a pane that looks alive, + // produces nothing, and swallows everything typed into it. Migrating a + // live pty between processes is the other answer and is a different + // feature; a fresh shell is the one this transport can keep. + for (&s.shells) |*sh| { + const owner = sh.owner orelse continue; + if (owner != gone) continue; + sh.owner = null; + sh.owed = true; + } + c.* = .{}; + } +}; + +// --------------------------------------------------------------------------- +// the process +// --------------------------------------------------------------------------- + +/// `pardes --detach[=<name>]`: one core, no terminal, a socket. The loop is the +/// core's own `pump`, exactly as the tty and gui shells run it — this frontend +/// simply has no window of its own. +/// +/// The pre-loop effect drain is here for the same reason tty.zig has one: the +/// startup spawns are already queued, and they have to be PERFORMED before the +/// loop rather than left in the queue. They reach no frontend — there is none +/// yet — and are remembered instead, then asked of the first attach; `Shell` +/// says why that is the only shape that works for a session whose panes exist +/// before its socket does. +pub fn run(init: std.process.Init, opts: pardes.Options, name: []const u8) !void { + const gpa = init.gpa; + const allocs = pardes.allocators.init(gpa); + defer pardes.allocators.deinit(); + var options = opts; + options.image_allocator = allocs.image; + options.pdf_allocator = allocs.pdf; + options.tree_sitter_allocator = allocs.tree_sitter; + options.frame_allocator = allocs.frame; + + // The core's own subsystems, not host work: a detached session syntax + // highlights and decodes images exactly like an attached one. + pardes.image.start(init.io, allocs.image); + if (comptime pardes.pdf_enabled) pardes.pdf.start(allocs.pdf); + pardes.syntax.start(allocs.tree_sitter); + defer { + pardes.image.stop(); + if (comptime pardes.pdf_enabled) pardes.pdf.stop(); + pardes.syntax.stop(); + } + + const core = if (options.load_path) |lp| blk: { + const bytes = try @import("../look.zig").readFile(gpa, lp); + defer gpa.free(bytes); + break :blk try pardes.Pardes.initFromDump(allocs.pardes, options, bytes); + } else try pardes.Pardes.init(allocs.pardes, options); + defer core.deinit(); + + var session: Session = .{ .gpa = gpa, .core = core, .cols = options.cols, .rows = options.rows }; + defer session.deinit(); + if (!session.listen(name)) { + // Loud, and on stderr rather than through the log: a `--detach` whose + // socket did not bind is a session nobody will ever find, and exiting + // is the only honest answer. + try std.Io.File.stderr().writeStreamingAll(init.io, "pardes: could not bind a detached session socket\n"); + return error.NoSocket; + } + + const h = session.host(); + core.host = h; + while (core.nextEffect()) |effect| core.perform(effect); + while (!core.quit) try core.pump(h); +} + +// --------------------------------------------------------------------------- +// the socket, nested.zig's way +// --------------------------------------------------------------------------- + +/// Re-exported so the frontend half of this transport (client.zig) has ONE +/// import for the socket conventions, and so that the file which owns the +/// convention is the file it asks. The definition and its reasoning are +/// nested.zig's. +pub const setCloexec = nested.setCloexec; + +/// Every descriptor in this transport is non-blocking, on both sides: the core +/// must never park on a peer (`waitInput`), and a frontend must never park on +/// the session (client.zig `wait`). `pub` for that second caller. +pub fn setNonblock(fd: c_int) void { + const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0)); + if (flags < 0) return; + var o: libc.O = @bitCast(@as(u32, @bitCast(flags))); + o.NONBLOCK = true; + _ = libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(o))))); +} + +/// A dead peer must never kill this process, and that is as true of a frontend +/// whose session ended as of a session whose frontend died — so client.zig +/// takes this one too. linux says it per write, darwin once per socket (see +/// `accept`); the `if (darwin)` is what keeps `MSG.NOSIGNAL`, which darwin's +/// headers do not have, out of that build. +pub const nosignal: u32 = if (darwin) 0 else libc.MSG.NOSIGNAL; + +pub const poll_in: i16 = @intCast(libc.POLL.IN); +pub const poll_out: i16 = @intCast(libc.POLL.OUT); +pub const poll_hup: i16 = @intCast(libc.POLL.HUP); +pub const poll_err: i16 = @intCast(libc.POLL.ERR); +pub const poll_nval: i16 = @intCast(libc.POLL.NVAL); + +/// Give a drained buffer's memory back, and only a big one's: see +/// `idle_retain`. Called where a queue empties rather than on a timer, because +/// that is the one moment the capacity is provably unused. +fn retire(gpa: std.mem.Allocator, list: *std.ArrayListUnmanaged(u8)) void { + if (list.items.len != 0 or list.capacity <= idle_retain) return; + list.clearAndFree(gpa); +} + +/// Monotonic milliseconds, the clock macos.zig's fling already times with and +/// for its reason: MONOTONIC and not REALTIME, because a handshake that expired +/// because NTP stepped the wall clock backwards is a bug nobody reproduces. +/// +/// Zero on failure, and every caller treats zero as "no clock" and enforces no +/// deadline at all — a session that cannot read a clock keeps every slot rather +/// than dropping every slot. +fn monotonicMs() i64 { + var ts: libc.timespec = undefined; + if (libc.clock_gettime(.MONOTONIC, &ts) != 0) return 0; + return @as(i64, ts.sec) * std.time.ms_per_s + @divTrunc(ts.nsec, std.time.ns_per_ms); +} + +/// Sleep, for the one case that has no descriptor to wait on (see `waitInput`). +fn nap(ms: u32) void { + var ts: libc.timespec = .{ + .sec = @intCast(ms / 1000), + .nsec = @intCast((ms % 1000) * std.time.ns_per_ms), + }; + _ = libc.nanosleep(&ts, null); +} + +/// `<dir>/pardes-detached-<name>.sock`. The prefix differs from nested.zig's +/// `pardes-<pid>.sock` on purpose: that file's sweeper unlinks the socket of any +/// name whose digits name a dead pid, and a session called `work` must never +/// look like one. The buffer is sun_path-sized, so a name that does not fit is +/// no address at all rather than a truncated one pointing somewhere else. +pub fn socketPath(buf: *[sun_path_len]u8, dir: []const u8, name: []const u8) ?[:0]const u8 { + // A name is one path component and nothing clever: a `/` would put the + // socket somewhere else entirely, and a NUL would truncate the address. + if (name.len == 0) return null; + if (std.mem.indexOfAny(u8, name, "/\x00") != null) return null; + return std.fmt.bufPrintSentinel(buf, "{s}/" ++ prefix ++ "{s}.sock", .{ dir, name }, 0) catch null; +} + +const prefix = "pardes-detached-"; + +/// The path a FRONTEND connects to for a session called `name`. Derived here +/// rather than in client.zig because this file owns the convention, and the +/// side that binds and the side that connects must not be able to disagree +/// about it. `path_max` is the buffer a caller has to supply. +pub const path_max = sun_path_len; + +pub fn sessionPath(buf: *[path_max]u8, name: []const u8) ?[:0]const u8 { + if (comptime !supported) return null; + var dir_buf: [sun_path_len:0]u8 = undefined; + const dir = nested.socketDir(&dir_buf) orelse return null; + return socketPath(buf, dir, name); +} + +/// The FRONTEND's half of the vetting this file does before it binds, and the +/// reason it is here rather than in client.zig: one convention, one predicate, +/// one file that owns both. +/// +/// Until this, the server refused a directory anyone else could write and a +/// socket anyone else could talk to, and the client connected to whatever it +/// found at the path it derived — which is the asymmetry this module's header +/// condemns in as many words. A socket planted at a path a frontend derives +/// from `$XDG_RUNTIME_DIR` receives every keystroke that frontend collects, and +/// answers with frames of its choosing. +/// +/// Checked and then connected, in that order, which is a TOCTOU only for +/// somebody who can already write the directory — and the directory is the +/// first thing this refuses. +pub fn vetted(path: [:0]const u8) bool { + if (comptime !supported) return false; + var dir_buf: [sun_path_len:0]u8 = undefined; + const dir = nested.socketDir(&dir_buf) orelse return false; + if (!ours(nested.statNoFollow(dir) orelse return false, s_ifdir)) return false; + return ours(nested.statNoFollow(path) orelse return false, s_ifsock); +} + +const s_ifmt: u32 = 0o170000; +const s_ifdir: u32 = 0o040000; +const s_ifsock: u32 = 0o140000; + +/// Is this a `kind` we own, with nothing granted to group or other? The three +/// questions `nested.ensureSocketDir` asks of the directory, asked of the +/// SOCKET too: the two walls are the directory's mode and the file's, and a +/// frontend that checks only one of them has checked neither. +fn ours(st: nested.DirFacts, kind: u32) bool { + if (st.mode & s_ifmt != kind) return false; + if (st.uid != libc.getuid()) return false; + return st.mode & 0o077 == 0; +} + +/// Is something LISTENING at `path`? The one place this file decides whether a +/// socket file is a corpse, asked by `listen` before it takes a name over and +/// by `sweep` before it unlinks anything. +/// +/// nested.zig can ask `kill(0)` because its filenames carry a pid; a detached +/// session is named by a PERSON, so the question is put to the socket: a +/// connect to a bound path with no listener is refused (ECONNREFUSED), and that +/// refusal is the ONLY evidence of death this accepts. Everything else is life, +/// including the case a blocking connect used to turn into a hang — a live +/// session busy inside the core has a full backlog and answers EAGAIN, which is +/// why this socket is NON-BLOCKING. EPERM, a socket() that failed and a path +/// that no longer fits are all "not proven dead" too, and leave the file alone. +/// +/// THE WINDOW THIS CANNOT SEE, stated because it is real: a session between its +/// own `bind` and its `listen(2)` also answers ECONNREFUSED and is alive. It is +/// two syscalls wide, it is only ever entered by another `pardes --detach` +/// starting in the same instant, and what the loser loses is a NAME (its +/// `listen` fails and it says so) rather than a session. Closing it needs a +/// lock file per session, which is a second thing to leak. +/// +/// The successful-connect case costs the live session one slot for one round: +/// closing this descriptor immediately turns the pending connection into an +/// EOF, which `receive` reads as a frontend that left. +fn alive(path: [:0]const u8) bool { + var addr: libc.sockaddr.un = .{ .path = @splat(0) }; + if (path.len + 1 > addr.path.len) return true; + @memcpy(addr.path[0 .. path.len + 1], path[0 .. path.len + 1]); + const fd = libc.socket(libc.AF.UNIX, libc.SOCK.STREAM, 0); + if (fd < 0) return true; + defer _ = libc.close(fd); + nested.setCloexec(fd); + setNonblock(fd); + const rc = libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))); + if (rc == 0) return true; + return libc.errno(rc) != .CONNREFUSED; +} + +/// Unlink the sockets of detached sessions that are gone — our own litter, +/// which `--attach`'s "the one session there is" would otherwise count as a +/// session (tty.zig `sessionName`). `alive` is the whole of the judgement. +/// +/// Bounded: one readdir of a directory only we write to, one connect each. +fn sweep(dir: [:0]const u8) void { + const d = libc.opendir(dir) orelse return; + defer _ = libc.closedir(d); + while (libc.readdir(d)) |ent| { + const name = std.mem.sliceTo(&ent.name, 0); + if (!std.mem.startsWith(u8, name, prefix) or !std.mem.endsWith(u8, name, ".sock")) continue; + var path_buf: [sun_path_len:0]u8 = undefined; + const path = std.fmt.bufPrintSentinel(&path_buf, "{s}/{s}", .{ dir, name }, 0) catch continue; + if (!alive(path)) _ = libc.unlink(path); + } +} diff --git a/src/detached/wire.zig b/src/detached/wire.zig new file mode 100644 index 00000000..f2030e26 --- /dev/null +++ b/src/detached/wire.zig @@ -0,0 +1,1755 @@ +//! THE DETACHED-SESSION WIRE FORMAT: one `Event` and one `Host.VTable` call per +//! message, byte for byte, with nothing native about the bytes. +//! +//! WHO OWNS THE CORE. The `Pardes` instance lives in the DETACHED process +//! (server.zig). A frontend (client.zig) owns a terminal and a socket and +//! nothing else: it sends the input it collects and draws the frames it is +//! sent. One core per session, N frontends attached to it, all looking at the +//! same screen — `screen -x`, not N sessions. +//! +//! WHY A CODEC AT ALL, when nested.zig's socket carries a builtin command line +//! and has nothing to version: a command line cannot carry a frame, and frames +//! and input are this transport's entire content. +//! +//! ARCHITECTURE-NEUTRAL, and not as decoration: the frontend on the other end +//! may be riscv32-freestanding (the ESP32-P4 board) while the core is x86_64 +//! linux. So: +//! * every integer is an explicit width, little-endian. No `usize` reaches +//! the wire — a pointer-sized field is 4 bytes on the board and 8 here, and +//! every field after it would then be read at the wrong offset. +//! * no native struct is ever blitted. `@bitCast`/`std.mem.asBytes` of a Zig +//! struct puts this compiler's field order and padding on a socket; every +//! field below is written and read by hand. +//! * every union and every enum gets a tag chosen HERE (`ClientTag`, +//! `ServerTag`, `ColorTag`, ...) and never `@intFromEnum` of a core type, +//! so reordering `Event` or `CellStyle.ul` cannot silently redefine the +//! protocol. The mapping switches are exhaustive: adding a variant to the +//! core is a compile error in this file, which is the point of them. +//! * every variable-length payload carries an explicit length prefix, and +//! `max_payload` bounds the lot. This is a parser on a socket: a malformed +//! frame must be REFUSED, never indexed past. +//! * a bool is one byte, 0 or 1. Any other value is a decode error rather +//! than "nonzero is true": a byte this protocol cannot mean is evidence +//! the stream is not the stream it claims to be. +//! * floats travel as their IEEE-754 binary32 bit pattern inside an explicit +//! u32. Both ends agree about binary32; neither agrees about struct layout. +//! +//! BUILD-NEUTRAL for the same reason. `Event.resize.cell_pixels` exists only +//! when native PDF placement is compiled in (pardes.zig `CellPixels`), and a +//! frontend must not have to have been built with the core's options — so it is +//! ALWAYS on the wire and dropped on arrival by a build with nowhere to put it. +//! +//! WHAT IS NOT HERE. The seam has twenty methods; this carries twelve of +//! them, and the eight it does not are named here with their reasons. +//! * `pull_wait_input` IS the server's poll loop, not a message. +//! * `push_poll_frame` and `push_post_present` carry no information. They are +//! per-frame bookkeeping ticks, and `frame` already arrives exactly once +//! per pump at the same place in the order — a frontend does its per-frame +//! work when a frame lands. Two more messages per frame per client would +//! say nothing the frame does not already say. +//! * `pull_tty_taken` and `pull_gpio_toggle` are answers the CALLER waits +//! for, and `pull_lsp`/`pull_pipe` are work dispatched off the loop. A +//! round trip inside `update` is the one thing this transport must never +//! do: the core would block on a socket, and `pull_wait_input`'s own +//! comment is that it is the only place this process may sleep. The +//! process that owns the core answers all four. +//! * `push_fs_reply` cannot be a broadcast. host.zig's rule is that the +//! transport which asked is the one holding the request; with N frontends, +//! N-1 would receive the answer to a request they never made. So the acme +//! mount stays in the detached process, where the `Event.fs_req` that +//! starts it is raised, and neither half of that pair is on the wire — +//! which is also why `Event.fs_req` has no `ClientTag`. +const std = @import("std"); +const pardes = @import("../pardes.zig"); + +/// Bumped whenever any layout below changes. Checked on connect and refused +/// loudly (see `Refusal.version`): two builds of pardes are routinely on one +/// machine — `zig build` replaces the binary under a running session — and a +/// frontend decoding another version's frame layout would paint garbage and +/// blame the terminal. +pub const version: u16 = 1; + +pub const Error = error{ + /// The message ended inside a field. + Truncated, + /// A length prefix, a run, or a grid dimension larger than this protocol + /// admits. Refused before anything is allocated or indexed. + Overlong, + /// A tag byte no version of this protocol has ever defined. + BadTag, + /// A tag this protocol does define, carrying a value it cannot mean: a + /// 3-in-a-bool, a zero-column resize, a pane past MAX_PANES. + BadValue, + /// The payload was decoded and bytes were left over. A message that says + /// more than its layout has room for is not this message. + Trailing, + /// The encoder ran out of caller-supplied buffer. + NoSpace, +}; + +// --------------------------------------------------------------------------- +// bounds +// --------------------------------------------------------------------------- + +/// The largest grid this protocol carries. `Surface.cols`/`rows` are u16, so +/// these are protocol bounds rather than type bounds, and they exist because +/// `max_payload` below is derived from them: a decoder that accepts 65535 +/// columns accepts a 25 GiB frame prefix. A 4K display at a 6-pixel font is +/// about 340 columns and 110 rows, so this is roughly 1.5x the largest grid +/// any real terminal has, and the board's own is 56x14. +pub const max_cols: u16 = 512; +pub const max_rows: u16 = 128; + +/// One cell at its largest: `default` false, a 7-byte grapheme, two rgb colors, +/// the attribute byte, the underline style and the font role. Written as the +/// sum of the fields rather than a number so that adding a field to `Cell` +/// moves it. +const cell_max = 1 + 1 + 7 + 4 + 4 + 1 + 1 + 1; + +/// `start:u32 + count:u16`. A run's cost, and therefore the break-even the +/// encoder coalesces against (see `encodeFrame`). +const run_header = 4 + 2; + +/// `kind:u8 + cols:u16 + rows:u16 + cursor(6) + nruns:u32`. +const frame_head = 1 + 2 + 2 + 6 + 4; + +/// The longest legal payload, and therefore the length prefix a decoder will +/// accept before it refuses the stream. Three messages set it: +/// * a full frame of the largest grid, worst case one run per cell: +/// 512*128 * (6 + 20) = 1.6 MiB. +/// * one paste, which the tty frontend already caps at 4 MiB (tty.zig +/// `max_paste_bytes`) on the grounds that anything larger is a mis-click. +/// * a `write_file`, whose bytes are a pane's whole text and are the only +/// genuinely open-ended payload here. +/// 16 MiB is past every source file anyone edits in this editor and is still a +/// buffer the receiving side can simply hold. A larger message is not sent and +/// a larger prefix is not read. +pub const max_payload: u32 = 16 << 20; + +/// Every message is `tag:u8, len:u32le, payload[len]`. A u32 because a full +/// frame and a paste both pass 64 KiB; a u16 would have needed the frame split +/// across messages, which is a second framing layer for no gain. +pub const header_len = 5; + +/// Bytes `encodeFrame` may need for this grid, worst case: every cell changed, +/// every cell in a run of its own, every cell at `cell_max`. The server sizes +/// one buffer from this per geometry rather than guessing. +pub fn frameBound(cols: u16, rows: u16) usize { + return header_len + frame_head + @as(usize, cols) * @as(usize, rows) * (run_header + cell_max); +} + +// --------------------------------------------------------------------------- +// tags +// --------------------------------------------------------------------------- + +/// Frontend -> core. Exhaustive on purpose, which is the opposite of +/// fuse.zig's `Opcode`: there, a newer KERNEL adds opcodes and a non-exhaustive +/// enum is the only way to receive one without undefined behaviour. Here both +/// ends are pardes and an unknown tag is not a newer peer — `version` already +/// refused that — so it is a corrupt or hostile stream and must be rejected. +/// `std.enums.fromInt` is how, at the one place a byte becomes a tag. +/// +/// The numbers are the PROTOCOL's, grouped session/input rather than derived +/// from `Event`'s declaration order, so reordering the union changes nothing. +pub const ClientTag = enum(u8) { + hello = 0x01, + bye = 0x02, + + key = 0x10, + mouse = 0x11, + resize = 0x12, + output = 0x13, + eof = 0x14, + lsp_resp = 0x15, + pipe_resp = 0x16, + file_changed = 0x17, + paste = 0x18, + command = 0x19, + pdf_scroll = 0x1a, + pinch = 0x1b, + touch_scroll = 0x1c, + pointer_leave = 0x1d, + tick = 0x1e, +}; + +/// Core -> frontend. 0x01..0x0f is the session, 0x10.. is one `push_` method +/// each, in `Host.VTable`'s own order so the two lists can be read side by +/// side. +pub const ServerTag = enum(u8) { + welcome = 0x01, + refuse = 0x02, + frame = 0x03, + quit = 0x04, + + spawn = 0x10, + pty_write = 0x11, + pty_resize = 0x12, + write_file = 0x13, + write_dump = 0x14, + watch_file = 0x15, + watch_theme = 0x16, + dump_themes = 0x17, + set_clipboard = 0x18, + read_clipboard = 0x19, + open_link = 0x1a, +}; + +/// Why the server hung up on a connect. Sent as a `refuse` and followed by a +/// close, so a frontend can say something specific instead of "connection +/// closed". +pub const Refusal = enum(u8) { + /// `Hello.version` is not `version`. The frontend and the core are two + /// builds of pardes. + version = 0x01, + /// Every client slot is taken (see server.zig `max_clients`). + full = 0x02, + /// The core has already quit; this session is ending. + quitting = 0x03, +}; + +const ColorTag = enum(u8) { default = 0x00, index = 0x01, rgb = 0x02 }; +const UlTag = enum(u8) { off = 0x00, single = 0x01, double = 0x02, curly = 0x03, dotted = 0x04, dashed = 0x05 }; +const FontTag = enum(u8) { body = 0x00, tagline = 0x01 }; +const ButtonTag = enum(u8) { + left = 0x00, + middle = 0x01, + right = 0x02, + wheel_up = 0x03, + wheel_down = 0x04, + wheel_left = 0x05, + wheel_right = 0x06, + none = 0x07, +}; +const KindTag = enum(u8) { press = 0x00, release = 0x01, motion = 0x02, drag = 0x03 }; + +/// A full frame resets the receiver's grid to unpainted cells and then applies +/// its runs; a diff applies its runs on top of what is already there. One bit +/// of semantics and one code path, and it is what makes a full frame of a +/// mostly-empty grid cheap. +pub const FrameKind = enum(u8) { full = 0x01, diff = 0x02 }; + +/// Bit per `CellStyle` bool, packed into one byte. Bit 7 is unassigned and a +/// set bit 7 is a decode error: it is a byte this protocol cannot mean. +const attr_bold: u8 = 1 << 0; +const attr_dim: u8 = 1 << 1; +const attr_italic: u8 = 1 << 2; +const attr_blink: u8 = 1 << 3; +const attr_reverse: u8 = 1 << 4; +const attr_invisible: u8 = 1 << 5; +const attr_strikethrough: u8 = 1 << 6; +const attr_reserved: u8 = 1 << 7; + +// --------------------------------------------------------------------------- +// messages +// --------------------------------------------------------------------------- + +/// First message on every connection, and `version` is its first field at a +/// fixed offset for exactly one reason: a mismatch has to be diagnosable even +/// when the rest of the layout is the part that changed. +pub const Hello = struct { + version: u16 = version, + /// This frontend's grid. Never zero — see `Cursor` for why zero is refused + /// rather than clamped. + cols: u16, + rows: u16, +}; + +pub const Welcome = struct { + version: u16 = version, + /// Which client slot this connection got. Carried because it is what the + /// server's own diagnostics name, so both sides say the same number. + slot: u8, + /// The session's grid as of this attach — the smallest common one, which + /// may be smaller than the `Hello` asked for. See server.zig `geometry`. + cols: u16, + rows: u16, +}; + +pub const Cursor = struct { x: u16, y: u16, bar: bool }; + +/// One frame, head decoded and runs left encoded. The runs are NOT expanded +/// into a slice of cells here: a frame of the largest grid is 1.6 MiB, this +/// union is passed by value, and the receiver already owns the grid the runs +/// belong in. `apply` is the bounds-checked walk. +pub const Frame = struct { + kind: FrameKind, + cols: u16, + rows: u16, + cursor: ?Cursor, + nruns: u32, + runs: []const u8, + + /// Paint this frame into `grid`, which must be exactly `cols * rows` cells + /// — the receiver resizes on a geometry change before applying, and a grid + /// of the wrong size is a receiver bug, not a wire condition. + /// + /// Every run is range-checked against the grid before a single cell is + /// written, so a run claiming to start past the end writes nothing. + pub fn apply(f: Frame, grid: []pardes.Cell) Error!void { + if (grid.len != @as(usize, f.cols) * @as(usize, f.rows)) return error.BadValue; + if (f.kind == .full) @memset(grid, .{}); + var r: Reader = .init(f.runs); + var i: u32 = 0; + while (i < f.nruns) : (i += 1) { + const start = try r.getU32(); + const count = try r.getU16(); + // A zero-length run is not something `encodeFrame` emits, and + // accepting one would let a peer spend the run budget saying + // nothing. + if (count == 0) return error.BadValue; + // THE bounds check. Written as `count > len - start` rather than + // `start + count > len` because the sum of two attacker-chosen + // 32-bit numbers is the classic way this check is bypassed. + if (start > grid.len or count > grid.len - start) return error.Overlong; + for (grid[start..][0..count]) |*c| c.* = try decodeCell(&r); + } + try r.end(); + } +}; + +/// Frontend -> core, decoded. The `Event`'s slices BORROW the payload buffer, +/// exactly like the pty chunks the tty host hands to `update`: valid for that +/// one call and no longer. +pub const ClientMsg = union(enum) { + hello: Hello, + bye, + event: pardes.Event, +}; + +/// Core -> frontend, decoded. Slices borrow the payload buffer the same way. +pub const ServerMsg = union(enum) { + welcome: Welcome, + refuse: Refusal, + frame: Frame, + /// The session is over. Sent before the listener closes so a frontend can + /// exit rather than report a broken pipe. + quit, + + spawn: struct { pane: u8, cwd: []const u8 }, + pty_write: struct { pane: u8, bytes: []const u8 }, + pty_resize: struct { pane: u8, cols: u16, rows: u16 }, + write_file: struct { pane: u8, path: []const u8, bytes: []const u8 }, + write_dump: []const u8, + watch_file: struct { pane: u8, path: []const u8, on: bool }, + watch_theme: struct { generation: u32, on: bool }, + dump_themes: struct { pane: u8 }, + set_clipboard: []const u8, + read_clipboard, + open_link: []const u8, +}; + +/// The one thing a decoder cannot put in a byte buffer: `pipe_resp.outputs` is +/// a `[]const []const u8`, so the outer array needs somewhere to live. Sized +/// from the core's own ceiling on selections (`MAX_SELS`), which is what bounds +/// the count a legitimate `pipe_resp` can carry. +pub const Scratch = struct { + outputs: [pardes.MAX_SELS][]const u8 = undefined, +}; + +// --------------------------------------------------------------------------- +// primitives +// --------------------------------------------------------------------------- + +pub const Writer = struct { + buf: []u8, + n: usize = 0, + + pub fn init(buf: []u8) Writer { + return .{ .buf = buf }; + } + + pub fn written(w: *const Writer) []const u8 { + return w.buf[0..w.n]; + } + + fn room(w: *Writer, k: usize) Error![]u8 { + if (w.buf.len - w.n < k) return error.NoSpace; + defer w.n += k; + return w.buf[w.n..][0..k]; + } + + fn putByte(w: *Writer, v: u8) Error!void { + (try w.room(1))[0] = v; + } + + fn putU16(w: *Writer, v: u16) Error!void { + std.mem.writeInt(u16, (try w.room(2))[0..2], v, .little); + } + + fn putU32(w: *Writer, v: u32) Error!void { + std.mem.writeInt(u32, (try w.room(4))[0..4], v, .little); + } + + /// One byte, 0 or 1, never "nonzero". `getBool` refuses anything else. + fn putBool(w: *Writer, v: bool) Error!void { + try w.putByte(@intFromBool(v)); + } + + /// IEEE-754 binary32, as its bit pattern in an explicit u32. A scalar + /// bitcast and not a struct blit: the format is the one thing a riscv32 + /// and an x86_64 do agree about. + fn putF32(w: *Writer, v: f32) Error!void { + try w.putU32(@bitCast(v)); + } + + fn putBytes(w: *Writer, v: []const u8) Error!void { + @memcpy(try w.room(v.len), v); + } + + /// Length-prefixed, u16: paths, command lines, a key's text. Nothing here + /// is legitimately longer than 64 KiB and a u16 says so on the wire. + fn putSlice16(w: *Writer, v: []const u8) Error!void { + if (v.len > std.math.maxInt(u16)) return error.Overlong; + try w.putU16(@intCast(v.len)); + try w.putBytes(v); + } + + /// ...and u32 for the ones that are: pty output, a paste, a file. + fn putSlice32(w: *Writer, v: []const u8) Error!void { + if (v.len > max_payload) return error.Overlong; + try w.putU32(@intCast(v.len)); + try w.putBytes(v); + } +}; + +pub const Reader = struct { + bytes: []const u8, + i: usize = 0, + + pub fn init(bytes: []const u8) Reader { + return .{ .bytes = bytes }; + } + + /// `i <= bytes.len` is the invariant every getter below preserves, which is + /// what makes the subtraction here safe. + fn take(r: *Reader, n: usize) Error![]const u8 { + if (r.bytes.len - r.i < n) return error.Truncated; + defer r.i += n; + return r.bytes[r.i..][0..n]; + } + + fn getByte(r: *Reader) Error!u8 { + return (try r.take(1))[0]; + } + + fn getU16(r: *Reader) Error!u16 { + return std.mem.readInt(u16, (try r.take(2))[0..2], .little); + } + + fn getU32(r: *Reader) Error!u32 { + return std.mem.readInt(u32, (try r.take(4))[0..4], .little); + } + + fn getBool(r: *Reader) Error!bool { + return switch (try r.getByte()) { + 0 => false, + 1 => true, + else => error.BadValue, + }; + } + + fn getF32(r: *Reader) Error!f32 { + const v: f32 = @bitCast(try r.getU32()); + // A NaN or an infinity here is not a scroll distance. `pinch` in + // particular multiplies into a zoom factor the pane keeps, so one bad + // value poisons that pane for the rest of the session. + if (!std.math.isFinite(v)) return error.BadValue; + return v; + } + + fn getSlice16(r: *Reader) Error![]const u8 { + return r.take(try r.getU16()); + } + + fn getSlice32(r: *Reader) Error![]const u8 { + const n = try r.getU32(); + if (n > max_payload) return error.Overlong; + return r.take(n); + } + + /// A pane id the core can actually index. `Event.output{ .pane = 200 }` + /// would reach `p.panes[200]`. + fn getPane(r: *Reader) Error!u8 { + const pane = try r.getByte(); + if (pane >= pardes.MAX_PANES) return error.BadValue; + return pane; + } + + /// A grid the core can render into. Zero is refused rather than clamped: a + /// frontend that has not been sized yet must not be allowed to collapse a + /// shared session to nothing (see server.zig `geometry`). + fn getCols(r: *Reader) Error!u16 { + const v = try r.getU16(); + if (v == 0 or v > max_cols) return error.BadValue; + return v; + } + + fn getRows(r: *Reader) Error!u16 { + const v = try r.getU16(); + if (v == 0 or v > max_rows) return error.BadValue; + return v; + } + + fn getTag(r: *Reader, comptime T: type) Error!T { + return std.enums.fromInt(T, try r.getByte()) orelse error.BadTag; + } + + fn end(r: *Reader) Error!void { + if (r.i != r.bytes.len) return error.Trailing; + } +}; + +/// Reserve a message header; `finishMessage` back-patches the length. Two +/// passes would mean walking a frame's runs twice to find out how long they +/// are, and the walk is the expensive half. +fn beginMessage(w: *Writer, tag: u8) Error!usize { + const at = w.n; + try w.putByte(tag); + try w.putU32(0); + return at; +} + +fn finishMessage(w: *Writer, at: usize) Error!void { + const len = w.n - at - header_len; + if (len > max_payload) return error.Overlong; + std.mem.writeInt(u32, w.buf[at + 1 ..][0..4], @intCast(len), .little); +} + +/// One complete message at the front of a stream buffer, or null when the rest +/// has not arrived. `total` is what the caller consumes. +pub const Framed = struct { tag: u8, payload: []const u8, total: usize }; + +pub fn framed(buf: []const u8) Error!?Framed { + if (buf.len < header_len) return null; + const len = std.mem.readInt(u32, buf[1..5], .little); + // Refused BEFORE the caller grows a buffer to hold it: an over-long prefix + // is the one field in this protocol that can ask for memory. + if (len > max_payload) return error.Overlong; + if (buf.len - header_len < len) return null; + return .{ .tag = buf[0], .payload = buf[header_len..][0..len], .total = header_len + len }; +} + +// --------------------------------------------------------------------------- +// cells +// --------------------------------------------------------------------------- + +fn putColor(w: *Writer, c: pardes.Color) Error!void { + switch (c) { + .default => try w.putByte(@intFromEnum(ColorTag.default)), + .index => |i| { + try w.putByte(@intFromEnum(ColorTag.index)); + try w.putByte(i); + }, + .rgb => |v| { + try w.putByte(@intFromEnum(ColorTag.rgb)); + try w.putBytes(&v); + }, + } +} + +fn getColor(r: *Reader) Error!pardes.Color { + return switch (try r.getTag(ColorTag)) { + .default => .default, + .index => .{ .index = try r.getByte() }, + .rgb => .{ .rgb = (try r.take(3))[0..3].* }, + }; +} + +const Ul = @FieldType(pardes.CellStyle, "ul"); + +/// The four enum mappings — underline, font role, mouse button, mouse kind — +/// and `clientTag`/`serverTag` further down all have one shape: the WIRE +/// member of the same NAME. So the byte on the socket is still `@intFromEnum` +/// of an explicitly numbered enum in THIS file and never of a core type — +/// reordering `CellStyle.ul` changes nothing — and adding a variant to the +/// core without adding one here does not compile: +/// +/// error: enum 'wire.UlTag' has no member named 'wavy' +/// +/// which is the whole property the six hand-written switches this replaced +/// existed for, in a form that cannot fall out of step. A variant that has to +/// travel under a DIFFERENT name than the core's gets an arm of its own before +/// the `inline else`, exactly as `clientTag` keeps `.fs_req`. +fn putUl(w: *Writer, u: Ul) Error!void { + try w.putByte(switch (u) { + inline else => |t| @intFromEnum(@field(UlTag, @tagName(t))), + }); +} + +fn getUl(r: *Reader) Error!Ul { + return switch (try r.getTag(UlTag)) { + inline else => |t| @field(Ul, @tagName(t)), + }; +} + +fn putFont(w: *Writer, f: pardes.FontRole) Error!void { + try w.putByte(switch (f) { + inline else => |t| @intFromEnum(@field(FontTag, @tagName(t))), + }); +} + +fn getFont(r: *Reader) Error!pardes.FontRole { + return switch (try r.getTag(FontTag)) { + inline else => |t| @field(pardes.FontRole, @tagName(t)), + }; +} + +/// How many bytes `putCell` will write. Exact, because `encodeFrame` weighs it +/// against `run_header` to decide whether to coalesce a run across it. +fn cellSize(c: *const pardes.Cell) usize { + if (c.default) return 1; + return 1 + 1 + c.len + colorSize(c.style.fg) + colorSize(c.style.bg) + 1 + 1 + 1; +} + +fn colorSize(c: pardes.Color) usize { + return switch (c) { + .default => 1, + .index => 2, + .rgb => 4, + }; +} + +/// An unpainted cell is ONE byte: `default` means "the shell renders the +/// terminal's default cell" (pardes.zig `Cell`), so its text and style are not +/// merely equal to the defaults, they are not part of the frame at all. +fn putCell(w: *Writer, c: *const pardes.Cell) Error!void { + try w.putBool(c.default); + if (c.default) return; + try w.putByte(c.len); + try w.putBytes(c.grapheme()); + const s = c.style; + try putColor(w, s.fg); + try putColor(w, s.bg); + var attrs: u8 = 0; + if (s.bold) attrs |= attr_bold; + if (s.dim) attrs |= attr_dim; + if (s.italic) attrs |= attr_italic; + if (s.blink) attrs |= attr_blink; + if (s.reverse) attrs |= attr_reverse; + if (s.invisible) attrs |= attr_invisible; + if (s.strikethrough) attrs |= attr_strikethrough; + try w.putByte(attrs); + try putUl(w, s.ul); + try putFont(w, s.font_role); +} + +fn decodeCell(r: *Reader) Error!pardes.Cell { + if (try r.getBool()) return .{}; + var c: pardes.Cell = .{ .default = false }; + const len = try r.getByte(); + // `Cell.text` is 7 bytes and `grapheme()` slices to `len`. A zero-length + // grapheme is a cell with nothing to draw and no way to advance a column. + if (len == 0 or len > c.text.len) return error.BadValue; + c.len = len; + @memcpy(c.text[0..len], try r.take(len)); + c.style.fg = try getColor(r); + c.style.bg = try getColor(r); + const attrs = try r.getByte(); + if (attrs & attr_reserved != 0) return error.BadValue; + c.style.bold = attrs & attr_bold != 0; + c.style.dim = attrs & attr_dim != 0; + c.style.italic = attrs & attr_italic != 0; + c.style.blink = attrs & attr_blink != 0; + c.style.reverse = attrs & attr_reverse != 0; + c.style.invisible = attrs & attr_invisible != 0; + c.style.strikethrough = attrs & attr_strikethrough != 0; + c.style.ul = try getUl(r); + c.style.font_role = try getFont(r); + return c; +} + +// --------------------------------------------------------------------------- +// frames +// --------------------------------------------------------------------------- + +fn putCursor(w: *Writer, cursor: ?Cursor, cols: u16, rows: u16) Error!void { + // Fixed six bytes whether or not there is a cursor. An optional field + // would save five bytes on a message that is already hundreds, and cost a + // branch on both sides of the wire. + try w.putBool(cursor != null); + const c = cursor orelse Cursor{ .x = 0, .y = 0, .bar = false }; + // Refused here as well as on decode, so this side cannot build a frame its + // own decoder rejects: a dropped frame (sendFrame logs it) leaves a stale + // screen, and a refused one takes the frontend's connection with it. + if (c.x >= cols or c.y >= rows) return error.BadValue; + try w.putU16(c.x); + try w.putU16(c.y); + try w.putBool(c.bar); +} + +fn getCursor(r: *Reader, cols: u16, rows: u16) Error!?Cursor { + const present = try r.getBool(); + const c: Cursor = .{ .x = try r.getU16(), .y = try r.getU16(), .bar = try r.getBool() }; + // THE cursor bounds check. Every other field of a frame is checked against + // the grid and these two were not, and they are the two a frontend indexes + // with directly — `Surface.at(cursor.x, cursor.y)` on a grid this frame + // says is 4x2 is an out-of-bounds write in EVERY frontend, not a wrong + // glyph. Checked whether or not the cursor is present, because the absent + // case is six bytes of padding this encoder writes as zero and a peer that + // fills them with anything else is not speaking this protocol. + if (c.x >= cols or c.y >= rows) return error.BadValue; + return if (present) c else null; +} + +/// Encode one frame of `cells`, as a DIFF against `prev` when `prev` is the +/// same grid this receiver last acknowledged, and as a FULL frame otherwise. +/// +/// A late joiner and a resize are the same case and are handled by the same +/// test: nothing on the far side is comparable to this grid, so `prev` is +/// empty or a different length and the frame becomes full. +/// +/// Why a diff at all. Measured against a real `pardes --detach` at 80x24 with +/// the default theme, which paints every cell and gives most of them an rgb +/// pair (14 bytes a cell rather than the 8 an uncoloured one costs): +/// * a full frame is 21965-26699 bytes, depending on what is on screen, +/// * a frame that changes nothing is 20 — the head, and no runs at all, +/// * a one-row change (`Msg hi`) is 54 bytes in one run, and a wordier one +/// 166, +/// * and an IDLE session sends nothing whatever, because the core is asleep +/// in `poll(2)` and produces no frame until something happens. +/// So the diff is worth roughly 500x on the traffic a session actually +/// generates. The byte-count test below asserts the arithmetic on the board's +/// own 56x14 grid with uncoloured cells, where it is exact: 6298 against 56. +/// This transport exists for a frontend on the far end of a slow link, and +/// shipping the Surface every frame would be 3 MiB/s at the animation tick. +pub fn encodeFrame( + out: []u8, + cols: u16, + rows: u16, + cursor: ?Cursor, + cells: []const pardes.Cell, + prev: []const pardes.Cell, +) Error![]const u8 { + // The protocol's ceiling, enforced by the ENCODER too, and BEFORE the + // assert below so a caller can be told rather than tripped. `max_payload` + // is derived from these two, so a larger grid is a frame this decoder + // refuses as Overlong — and sendFrame's log line already claims that this + // is the error it is catching, which was true of nothing until here. + if (cols == 0 or cols > max_cols or rows == 0 or rows > max_rows) return error.Overlong; + std.debug.assert(cells.len == @as(usize, cols) * @as(usize, rows)); + const full = prev.len != cells.len; + var w: Writer = .init(out); + const at = try beginMessage(&w, @intFromEnum(ServerTag.frame)); + try w.putByte(@intFromEnum(@as(FrameKind, if (full) .full else .diff))); + try w.putU16(cols); + try w.putU16(rows); + try putCursor(&w, cursor, cols, rows); + const nruns_at = w.n; + try w.putU32(0); + + var nruns: u32 = 0; + var i: usize = 0; + while (i < cells.len) { + if (!sendCell(cells, prev, full, i)) { + i += 1; + continue; + } + const start = i; + var run_end = i + 1; + i += 1; + while (i < cells.len) { + if (sendCell(cells, prev, full, i)) { + run_end = i + 1; + i += 1; + continue; + } + // A gap. Re-sending cells the far side already has is cheaper than + // a second run header whenever their encoded size adds up to less + // than one — true for short stretches of unpainted cells at a byte + // each, and false as soon as one painted cell (eight bytes at its + // smallest) is in the way. So the lookahead can never need to pass + // `run_header - 1` cells. + var gap: usize = 0; + var j = i; + while (j < cells.len and gap < run_header) : (j += 1) { + if (sendCell(cells, prev, full, j)) break; + gap += cellSize(&cells[j]); + } + if (j >= cells.len or gap >= run_header) break; + i = j; + } + try w.putU32(@intCast(start)); + try w.putU16(@intCast(run_end - start)); + for (cells[start..run_end]) |*c| try putCell(&w, c); + nruns += 1; + i = run_end; + } + std.mem.writeInt(u32, w.buf[nruns_at..][0..4], nruns, .little); + try finishMessage(&w, at); + return w.written(); +} + +fn sendCell(cells: []const pardes.Cell, prev: []const pardes.Cell, full: bool, i: usize) bool { + // Full: the receiver reset the grid, so unpainted cells are already right. + // Diff: anything the receiver cannot already be showing. + if (full) return !cells[i].default; + return !cells[i].visuallyEqual(&prev[i]); +} + +// --------------------------------------------------------------------------- +// frontend -> core +// --------------------------------------------------------------------------- + +/// Which tag a message travels under: the `ClientTag` of the same NAME, so a +/// new `Event` variant does not compile until it has a number here. See +/// `putUl` for the error it produces and why this is not a written-out table. +fn clientTag(msg: ClientMsg) ClientTag { + return switch (msg) { + .event => |ev| switch (ev) { + // Not on the wire, and not an omission: see the module header. + // The acme mount lives with the core, so this event is raised in + // the same process that answers it and never crosses a socket. + .fs_req => unreachable, + inline else => |_, t| @field(ClientTag, @tagName(t)), + }, + inline else => |_, t| @field(ClientTag, @tagName(t)), + }; +} + +/// Does the build this frontend was compiled with carry pixel dimensions in a +/// resize? The FIELD is comptime-conditional (pardes.zig `CellPixels`); the +/// WIRE is not. +const has_cell_pixels = @hasField(pardes.CellPixels, "w"); + +pub fn encodeClient(out: []u8, msg: ClientMsg) Error![]const u8 { + var w: Writer = .init(out); + const at = try beginMessage(&w, @intFromEnum(clientTag(msg))); + switch (msg) { + .hello => |h| { + try w.putU16(h.version); + try w.putU16(h.cols); + try w.putU16(h.rows); + }, + .bye => {}, + .event => |ev| switch (ev) { + .key => |k| { + try w.putU32(k.cp); + try w.putSlice16(k.text); + try w.putBool(k.ctrl); + try w.putBool(k.alt); + try w.putBool(k.shift); + }, + .mouse => |m| { + try w.putByte(switch (m.button) { + inline else => |t| @intFromEnum(@field(ButtonTag, @tagName(t))), + }); + try w.putByte(switch (m.kind) { + inline else => |t| @intFromEnum(@field(KindTag, @tagName(t))), + }); + try w.putU16(m.col); + try w.putU16(m.row); + try w.putBool(m.ctrl); + }, + .resize => |rs| { + try w.putU16(rs.cols); + try w.putU16(rs.rows); + // The conventional 1:2 cell aspect when this build has no + // pixels of its own, which is the same default the field + // carries where it exists. + try w.putU16(if (comptime has_cell_pixels) rs.cell_pixels.w else 8); + try w.putU16(if (comptime has_cell_pixels) rs.cell_pixels.h else 16); + }, + .output => |o| { + try w.putByte(o.pane); + try w.putSlice32(o.bytes); + }, + .eof => |e| try w.putByte(e.pane), + .lsp_resp => |l| { + try w.putU32(l.id); + try w.putSlice32(l.rows); + }, + .pipe_resp => |p| { + try w.putU32(p.id); + try w.putBool(p.success); + if (p.outputs.len > pardes.MAX_SELS) return error.Overlong; + try w.putU16(@intCast(p.outputs.len)); + for (p.outputs) |o| try w.putSlice32(o); + }, + .file_changed => |f| { + try w.putByte(f.pane); + try w.putSlice32(f.bytes); + }, + .paste => |b| try w.putSlice32(b), + .command => |line| try w.putSlice16(line), + .pdf_scroll => |s| { + try w.putByte(s.pane); + try w.putF32(s.delta_pixels); + }, + .pinch => |v| try w.putF32(v), + .touch_scroll => |v| try w.putF32(v), + .pointer_leave, .tick => {}, + .fs_req => unreachable, + }, + } + try finishMessage(&w, at); + return w.written(); +} + +pub fn decodeClient(tag: u8, payload: []const u8, scratch: *Scratch) Error!ClientMsg { + var r: Reader = .init(payload); + const msg: ClientMsg = switch (std.enums.fromInt(ClientTag, tag) orelse return error.BadTag) { + .hello => .{ .hello = .{ + .version = try r.getU16(), + .cols = try r.getCols(), + .rows = try r.getRows(), + } }, + .bye => .bye, + .key => blk: { + const cp = try r.getU32(); + // `Key.cp` is a u21, and the specials live in the private-use + // plane below this bound. A larger number is not a codepoint and + // @intCast of it would panic in a release build's own decoder. + if (cp > 0x10FFFF) return error.BadValue; + break :blk .{ .event = .{ .key = .{ + .cp = @intCast(cp), + .text = try r.getSlice16(), + .ctrl = try r.getBool(), + .alt = try r.getBool(), + .shift = try r.getBool(), + } } }; + }, + .mouse => .{ .event = .{ .mouse = .{ + .button = switch (try r.getTag(ButtonTag)) { + inline else => |t| @field(pardes.Mouse.Button, @tagName(t)), + }, + .kind = switch (try r.getTag(KindTag)) { + inline else => |t| @field(pardes.Mouse.Kind, @tagName(t)), + }, + .col = try r.getU16(), + .row = try r.getU16(), + .ctrl = try r.getBool(), + } } }, + .resize => blk: { + var ev: pardes.Event = .{ .resize = .{ .cols = try r.getCols(), .rows = try r.getRows() } }; + const px_w = try r.getU16(); + const px_h = try r.getU16(); + // Read either way — the bytes are on the wire — and kept only by a + // build that has somewhere to keep them. + if (comptime has_cell_pixels) ev.resize.cell_pixels = .{ .w = px_w, .h = px_h }; + break :blk .{ .event = ev }; + }, + .output => .{ .event = .{ .output = .{ .pane = try r.getPane(), .bytes = try r.getSlice32() } } }, + .eof => .{ .event = .{ .eof = .{ .pane = try r.getPane() } } }, + .lsp_resp => .{ .event = .{ .lsp_resp = .{ .id = try r.getU32(), .rows = try r.getSlice32() } } }, + .pipe_resp => blk: { + const id = try r.getU32(); + const success = try r.getBool(); + const n = try r.getU16(); + if (n > scratch.outputs.len) return error.Overlong; + for (scratch.outputs[0..n]) |*o| o.* = try r.getSlice32(); + break :blk .{ .event = .{ .pipe_resp = .{ + .id = id, + .success = success, + .outputs = scratch.outputs[0..n], + } } }; + }, + .file_changed => .{ .event = .{ .file_changed = .{ .pane = try r.getPane(), .bytes = try r.getSlice32() } } }, + .paste => .{ .event = .{ .paste = try r.getSlice32() } }, + .command => .{ .event = .{ .command = try r.getSlice16() } }, + .pdf_scroll => .{ .event = .{ .pdf_scroll = .{ .pane = try r.getPane(), .delta_pixels = try r.getF32() } } }, + .pinch => .{ .event = .{ .pinch = try r.getF32() } }, + .touch_scroll => .{ .event = .{ .touch_scroll = try r.getF32() } }, + .pointer_leave => .{ .event = .pointer_leave }, + .tick => .{ .event = .tick }, + }; + try r.end(); + return msg; +} + +// --------------------------------------------------------------------------- +// core -> frontend +// --------------------------------------------------------------------------- + +/// ...and the same rule in the same shape: the `ServerTag` of the same name. +fn serverTag(msg: ServerMsg) ServerTag { + return switch (msg) { + inline else => |_, t| @field(ServerTag, @tagName(t)), + }; +} + +/// Every server message EXCEPT a frame, which has its own encoder because its +/// payload is a walk of two grids rather than a value (see `encodeFrame`). +pub fn encodeServer(out: []u8, msg: ServerMsg) Error![]const u8 { + var w: Writer = .init(out); + const at = try beginMessage(&w, @intFromEnum(serverTag(msg))); + switch (msg) { + .welcome => |v| { + try w.putU16(v.version); + try w.putByte(v.slot); + try w.putU16(v.cols); + try w.putU16(v.rows); + }, + .refuse => |why| try w.putByte(@intFromEnum(why)), + // A frame's runs are not a value this union can hold; the server calls + // encodeFrame directly and this arm exists so the switch stays + // exhaustive over ServerMsg. + .frame => return error.BadValue, + .quit => {}, + .spawn => |s| { + try w.putByte(s.pane); + try w.putSlice16(s.cwd); + }, + .pty_write => |p| { + try w.putByte(p.pane); + try w.putSlice32(p.bytes); + }, + .pty_resize => |p| { + try w.putByte(p.pane); + try w.putU16(p.cols); + try w.putU16(p.rows); + }, + .write_file => |f| { + try w.putByte(f.pane); + try w.putSlice16(f.path); + try w.putSlice32(f.bytes); + }, + .write_dump => |b| try w.putSlice32(b), + .watch_file => |v| { + try w.putByte(v.pane); + try w.putSlice16(v.path); + try w.putBool(v.on); + }, + .watch_theme => |t| { + try w.putU32(t.generation); + try w.putBool(t.on); + }, + .dump_themes => |d| try w.putByte(d.pane), + .set_clipboard => |t| try w.putSlice32(t), + .read_clipboard => {}, + .open_link => |u| try w.putSlice16(u), + } + try finishMessage(&w, at); + return w.written(); +} + +/// Slack over a message's variable payload, covering every fixed field any +/// message here has plus its own header. One loose constant rather than a +/// field-by-field count: a bound that is 32 bytes generous costs one `memcpy` +/// worth of nothing, and a bound that is one byte tight is a bug that only +/// shows up on the message nobody tested. +const msg_slack = header_len + 32; + +/// An upper bound on `encodeServer`'s output, so a caller sizes its buffer +/// once instead of guessing and retrying. +pub fn serverBound(msg: ServerMsg) usize { + return msg_slack + switch (msg) { + .welcome, .refuse, .quit, .pty_resize, .watch_theme, .dump_themes, .read_clipboard => 0, + // A frame is bounded by its grid, not by this: see `frameBound`. + .frame => |f| frameBound(f.cols, f.rows), + .spawn => |s| s.cwd.len, + .pty_write => |p| p.bytes.len, + .write_file => |f| f.path.len + f.bytes.len, + .write_dump => |b| b.len, + .watch_file => |v| v.path.len, + .set_clipboard => |t| t.len, + .open_link => |u| u.len, + }; +} + +/// ...and the same for the frontend's side of the wire. +pub fn clientBound(msg: ClientMsg) usize { + return msg_slack + switch (msg) { + .hello, .bye => 0, + .event => |ev| switch (ev) { + .mouse, .resize, .eof, .pdf_scroll, .pinch, .touch_scroll, .pointer_leave, .tick => 0, + .key => |k| k.text.len, + .output => |o| o.bytes.len, + .lsp_resp => |l| l.rows.len, + .pipe_resp => |p| blk: { + // Each output carries its own u32 prefix, so the count is part + // of the bound and not just the bytes. + var total: usize = p.outputs.len * 4; + for (p.outputs) |o| total += o.len; + break :blk total; + }, + .file_changed => |f| f.bytes.len, + .paste => |b| b.len, + .command => |line| line.len, + .fs_req => unreachable, + }, + }; +} + +pub fn decodeServer(tag: u8, payload: []const u8) Error!ServerMsg { + var r: Reader = .init(payload); + const msg: ServerMsg = switch (std.enums.fromInt(ServerTag, tag) orelse return error.BadTag) { + .welcome => .{ .welcome = .{ + .version = try r.getU16(), + .slot = try r.getByte(), + .cols = try r.getCols(), + .rows = try r.getRows(), + } }, + .refuse => .{ .refuse = try r.getTag(Refusal) }, + .frame => blk: { + const kind = try r.getTag(FrameKind); + const cols = try r.getCols(); + const rows = try r.getRows(); + const cursor = try getCursor(&r, cols, rows); + const nruns = try r.getU32(); + // One run per cell is the worst an encoder emits, so a bigger + // count cannot be describing this grid. Checked HERE rather than + // in `apply`'s loop so the number is refused before it is used to + // bound anything. + if (nruns > @as(u32, cols) * @as(u32, rows)) return error.Overlong; + const runs = r.bytes[r.i..]; + r.i = r.bytes.len; + break :blk .{ .frame = .{ + .kind = kind, + .cols = cols, + .rows = rows, + .cursor = cursor, + .nruns = nruns, + .runs = runs, + } }; + }, + .quit => .quit, + .spawn => .{ .spawn = .{ .pane = try r.getPane(), .cwd = try r.getSlice16() } }, + .pty_write => .{ .pty_write = .{ .pane = try r.getPane(), .bytes = try r.getSlice32() } }, + .pty_resize => .{ .pty_resize = .{ + .pane = try r.getPane(), + .cols = try r.getU16(), + .rows = try r.getU16(), + } }, + .write_file => .{ .write_file = .{ + .pane = try r.getPane(), + .path = try r.getSlice16(), + .bytes = try r.getSlice32(), + } }, + .write_dump => .{ .write_dump = try r.getSlice32() }, + .watch_file => .{ .watch_file = .{ + .pane = try r.getPane(), + .path = try r.getSlice16(), + .on = try r.getBool(), + } }, + .watch_theme => .{ .watch_theme = .{ .generation = try r.getU32(), .on = try r.getBool() } }, + .dump_themes => .{ .dump_themes = .{ .pane = try r.getPane() } }, + .set_clipboard => .{ .set_clipboard = try r.getSlice32() }, + .read_clipboard => .read_clipboard, + .open_link => .{ .open_link = try r.getSlice16() }, + }; + try r.end(); + return msg; +} + +// --------------------------------------------------------------------------- +// tests +// --------------------------------------------------------------------------- +// +// Two obligations. Every message round-trips to a deep-equal value, because a +// codec written by hand is a codec whose two halves drift; and every malformed +// shape is REFUSED, because this parser is fed by a socket and a frame that is +// misread rather than refused is an out-of-bounds index. + +const testing = std.testing; + +/// Round-trip one frontend -> core message through the framing too, so a +/// length prefix that disagrees with the payload cannot pass. +fn roundClient(buf: []u8, msg: ClientMsg, scratch: *Scratch) !ClientMsg { + const bytes = try encodeClient(buf, msg); + const f = (try framed(bytes)).?; + try testing.expectEqual(bytes.len, f.total); + return decodeClient(f.tag, f.payload, scratch); +} + +fn roundServer(buf: []u8, msg: ServerMsg) !ServerMsg { + const bytes = try encodeServer(buf, msg); + const f = (try framed(bytes)).?; + try testing.expectEqual(bytes.len, f.total); + return decodeServer(f.tag, f.payload); +} + +test "detached wire: every Event variant round-trips" { + var buf: [4096]u8 = undefined; + var scratch: Scratch = .{}; + + // The tag space is the protocol's own, so assert the numbers themselves: + // a renumbering here breaks every deployed frontend and must be a diff + // somebody reads, not a silent change. + try testing.expectEqual(@as(u8, 0x01), @intFromEnum(ClientTag.hello)); + try testing.expectEqual(@as(u8, 0x10), @intFromEnum(ClientTag.key)); + try testing.expectEqual(@as(u8, 0x1e), @intFromEnum(ClientTag.tick)); + + { + const got = try roundClient(&buf, .{ .hello = .{ .cols = 80, .rows = 24 } }, &scratch); + try testing.expectEqual(version, got.hello.version); + try testing.expectEqual(@as(u16, 80), got.hello.cols); + try testing.expectEqual(@as(u16, 24), got.hello.rows); + } + try testing.expectEqual(ClientMsg.bye, try roundClient(&buf, .bye, &scratch)); + + { + const key: pardes.Key = .{ .cp = pardes.Key.page_down, .text = "ü", .ctrl = true, .alt = false, .shift = true }; + const got = (try roundClient(&buf, .{ .event = .{ .key = key } }, &scratch)).event.key; + try testing.expectEqual(key.cp, got.cp); + try testing.expectEqualStrings(key.text, got.text); + try testing.expectEqual(key.ctrl, got.ctrl); + try testing.expectEqual(key.alt, got.alt); + try testing.expectEqual(key.shift, got.shift); + } + { + // Every button and every kind, because the two mapping switches are + // the only place a value can be mistranslated and still decode. + for (std.enums.values(pardes.Mouse.Button)) |button| { + for (std.enums.values(pardes.Mouse.Kind)) |kind| { + const m: pardes.Mouse = .{ .button = button, .kind = kind, .col = 4200, .row = 7, .ctrl = true }; + const got = (try roundClient(&buf, .{ .event = .{ .mouse = m } }, &scratch)).event.mouse; + try testing.expectEqual(m.button, got.button); + try testing.expectEqual(m.kind, got.kind); + try testing.expectEqual(m.col, got.col); + try testing.expectEqual(m.row, got.row); + try testing.expectEqual(m.ctrl, got.ctrl); + } + } + } + { + const got = (try roundClient(&buf, .{ .event = .{ .resize = .{ .cols = 56, .rows = 14 } } }, &scratch)).event.resize; + try testing.expectEqual(@as(u16, 56), got.cols); + try testing.expectEqual(@as(u16, 14), got.rows); + // Pixels travel whether or not this build has them; where it does, the + // conventional 1:2 default survives the trip. + if (comptime has_cell_pixels) { + try testing.expectEqual(@as(u16, 8), got.cell_pixels.w); + try testing.expectEqual(@as(u16, 16), got.cell_pixels.h); + } + } + { + const got = (try roundClient(&buf, .{ .event = .{ .output = .{ .pane = 3, .bytes = "hi\x00there" } } }, &scratch)).event.output; + try testing.expectEqual(@as(u8, 3), got.pane); + try testing.expectEqualStrings("hi\x00there", got.bytes); + } + try testing.expectEqual(@as(u8, 15), (try roundClient(&buf, .{ .event = .{ .eof = .{ .pane = 15 } } }, &scratch)).event.eof.pane); + { + const got = (try roundClient(&buf, .{ .event = .{ .lsp_resp = .{ .id = 0xdeadbeef, .rows = "a:1:2-3 x" } } }, &scratch)).event.lsp_resp; + try testing.expectEqual(@as(u32, 0xdeadbeef), got.id); + try testing.expectEqualStrings("a:1:2-3 x", got.rows); + } + { + // Including an EMPTY output, which is what a filter that consumed a + // selection and printed nothing returns. + const outputs: []const []const u8 = &.{ "AA\n", "", "cc" }; + const got = (try roundClient(&buf, .{ .event = .{ .pipe_resp = .{ .id = 9, .success = true, .outputs = outputs } } }, &scratch)).event.pipe_resp; + try testing.expectEqual(@as(u32, 9), got.id); + try testing.expect(got.success); + try testing.expectEqual(@as(usize, 3), got.outputs.len); + for (outputs, got.outputs) |want, have| try testing.expectEqualStrings(want, have); + } + { + const got = (try roundClient(&buf, .{ .event = .{ .file_changed = .{ .pane = 0, .bytes = "" } } }, &scratch)).event.file_changed; + try testing.expectEqual(@as(u8, 0), got.pane); + try testing.expectEqualStrings("", got.bytes); + } + try testing.expectEqualStrings("clip", (try roundClient(&buf, .{ .event = .{ .paste = "clip" } }, &scratch)).event.paste); + try testing.expectEqualStrings("Look /x", (try roundClient(&buf, .{ .event = .{ .command = "Look /x" } }, &scratch)).event.command); + { + const got = (try roundClient(&buf, .{ .event = .{ .pdf_scroll = .{ .pane = 2, .delta_pixels = -12.5 } } }, &scratch)).event.pdf_scroll; + try testing.expectEqual(@as(u8, 2), got.pane); + try testing.expectEqual(@as(f32, -12.5), got.delta_pixels); + } + try testing.expectEqual(@as(f32, 1.25), (try roundClient(&buf, .{ .event = .{ .pinch = 1.25 } }, &scratch)).event.pinch); + try testing.expectEqual(@as(f32, -0.75), (try roundClient(&buf, .{ .event = .{ .touch_scroll = -0.75 } }, &scratch)).event.touch_scroll); + try testing.expectEqual( + std.meta.Tag(pardes.Event).pointer_leave, + (try roundClient(&buf, .{ .event = .pointer_leave }, &scratch)).event, + ); + try testing.expectEqual( + std.meta.Tag(pardes.Event).tick, + (try roundClient(&buf, .{ .event = .tick }, &scratch)).event, + ); +} + +test "detached wire: every server message round-trips" { + var buf: [4096]u8 = undefined; + + { + const got = (try roundServer(&buf, .{ .welcome = .{ .slot = 2, .cols = 56, .rows = 14 } })).welcome; + try testing.expectEqual(version, got.version); + try testing.expectEqual(@as(u8, 2), got.slot); + try testing.expectEqual(@as(u16, 56), got.cols); + try testing.expectEqual(@as(u16, 14), got.rows); + } + for (std.enums.values(Refusal)) |why| + try testing.expectEqual(why, (try roundServer(&buf, .{ .refuse = why })).refuse); + try testing.expectEqual(ServerMsg.quit, try roundServer(&buf, .quit)); + { + const got = (try roundServer(&buf, .{ .spawn = .{ .pane = 1, .cwd = "/home/x" } })).spawn; + try testing.expectEqual(@as(u8, 1), got.pane); + try testing.expectEqualStrings("/home/x", got.cwd); + } + { + const got = (try roundServer(&buf, .{ .pty_write = .{ .pane = 1, .bytes = "ls\r" } })).pty_write; + try testing.expectEqual(@as(u8, 1), got.pane); + try testing.expectEqualStrings("ls\r", got.bytes); + } + { + const got = (try roundServer(&buf, .{ .pty_resize = .{ .pane = 1, .cols = 80, .rows = 24 } })).pty_resize; + try testing.expectEqual(@as(u16, 80), got.cols); + try testing.expectEqual(@as(u16, 24), got.rows); + } + { + const got = (try roundServer(&buf, .{ .write_file = .{ .pane = 4, .path = "/tmp/a", .bytes = "body\n" } })).write_file; + try testing.expectEqual(@as(u8, 4), got.pane); + try testing.expectEqualStrings("/tmp/a", got.path); + try testing.expectEqualStrings("body\n", got.bytes); + } + try testing.expectEqualStrings(".{}", (try roundServer(&buf, .{ .write_dump = ".{}" })).write_dump); + { + const got = (try roundServer(&buf, .{ .watch_file = .{ .pane = 0, .path = "/tmp/b", .on = true } })).watch_file; + try testing.expectEqualStrings("/tmp/b", got.path); + try testing.expect(got.on); + } + { + const got = (try roundServer(&buf, .{ .watch_theme = .{ .generation = 7, .on = false } })).watch_theme; + try testing.expectEqual(@as(u32, 7), got.generation); + try testing.expect(!got.on); + } + try testing.expectEqual(@as(u8, 5), (try roundServer(&buf, .{ .dump_themes = .{ .pane = 5 } })).dump_themes.pane); + try testing.expectEqualStrings("yank", (try roundServer(&buf, .{ .set_clipboard = "yank" })).set_clipboard); + try testing.expectEqual(ServerMsg.read_clipboard, try roundServer(&buf, .read_clipboard)); + try testing.expectEqualStrings("https://x", (try roundServer(&buf, .{ .open_link = "https://x" })).open_link); +} + +/// A grid with something in every corner: a default cell, a plain ASCII cell, +/// an indexed pair, an rgb pair with every attribute on, and a multi-byte +/// grapheme — the five shapes `putCell` branches on. +fn sampleGrid(cells: []pardes.Cell) void { + @memset(cells, .{}); + cells[0] = .{ .text = "x".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; + cells[1] = .{ + .text = "y".* ++ @as([6]u8, @splat(0)), + .len = 1, + .default = false, + .style = .{ .fg = .{ .index = 3 }, .bg = .{ .index = 250 }, .bold = true, .ul = .curly }, + }; + cells[2] = .{ + .text = "→".* ++ @as([4]u8, @splat(0)), + .len = 3, + .default = false, + .style = .{ + .fg = .{ .rgb = .{ 1, 2, 3 } }, + .bg = .{ .rgb = .{ 250, 251, 252 } }, + .bold = true, + .dim = true, + .italic = true, + .blink = true, + .reverse = true, + .invisible = true, + .strikethrough = true, + .ul = .dashed, + .font_role = .tagline, + }, + }; + cells[cells.len - 1] = .{ .text = "z".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; +} + +fn expectGridEqual(want: []const pardes.Cell, have: []const pardes.Cell) !void { + try testing.expectEqual(want.len, have.len); + // `visuallyEqual` and not `std.meta.eql`: `Cell.text` past `len` is + // scratch left by an earlier grapheme, the decoder does not invent it, and + // pardes.zig says in as many words that it must never manufacture a diff. + for (want, have, 0..) |*a, *b, i| if (!a.visuallyEqual(b)) { + std.debug.print("cell {d} differs: {any} vs {any}\n", .{ i, a.*, b.* }); + return error.CellMismatch; + }; +} + +test "detached wire: a full frame carries the grid, a diff carries the change" { + const cols: u16 = 56; + const rows: u16 = 14; + const n = @as(usize, cols) * rows; + const gpa = testing.allocator; + + const cells = try gpa.alloc(pardes.Cell, n); + defer gpa.free(cells); + const mirror = try gpa.alloc(pardes.Cell, n); + defer gpa.free(mirror); + const buf = try gpa.alloc(u8, frameBound(cols, rows)); + defer gpa.free(buf); + + sampleGrid(cells); + @memset(mirror, .{}); + + // FULL: nothing on the far side is comparable, which is the late joiner + // and the resize both. + const full_bytes = try encodeFrame(buf, cols, rows, .{ .x = 3, .y = 4, .bar = true }, cells, &.{}); + { + const f = (try framed(full_bytes)).?; + const msg = (try decodeServer(f.tag, f.payload)).frame; + try testing.expectEqual(FrameKind.full, msg.kind); + try testing.expectEqual(cols, msg.cols); + try testing.expectEqual(rows, msg.rows); + try testing.expectEqual(@as(u16, 3), msg.cursor.?.x); + try testing.expectEqual(@as(u16, 4), msg.cursor.?.y); + try testing.expect(msg.cursor.?.bar); + try msg.apply(mirror); + try expectGridEqual(cells, mirror); + } + + // DIFF: two cells move. The mirror is already in sync, so this is what a + // steady-state frame looks like. + const before = try gpa.dupe(pardes.Cell, cells); + defer gpa.free(before); + cells[100] = .{ .text = "q".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; + cells[0] = .{}; // ...and one goes back to unpainted, which a diff must say + const diff_bytes = try encodeFrame(buf, cols, rows, null, cells, before); + { + const f = (try framed(diff_bytes)).?; + const msg = (try decodeServer(f.tag, f.payload)).frame; + try testing.expectEqual(FrameKind.diff, msg.kind); + try testing.expectEqual(@as(?Cursor, null), msg.cursor); + try testing.expectEqual(@as(u32, 2), msg.nruns); + try msg.apply(mirror); + try expectGridEqual(cells, mirror); + } + + // An unchanged frame is the head and nothing else: no runs, and the + // receiver keeps what it has. This is what makes an idle session silent. + { + const idle = try encodeFrame(buf, cols, rows, null, cells, cells); + try testing.expectEqual(@as(usize, header_len + frame_head), idle.len); + const f = (try framed(idle)).?; + const msg = (try decodeServer(f.tag, f.payload)).frame; + try testing.expectEqual(@as(u32, 0), msg.nruns); + try msg.apply(mirror); + try expectGridEqual(cells, mirror); + } +} + +test "detached wire: the diff is worth having, in bytes, on the board's grid" { + // The numbers quoted in `encodeFrame`'s comment, asserted so the claim + // cannot rot. 56x14 is the ESP32-P4 board's default grid (esp32p4.zig). + const cols: u16 = 56; + const rows: u16 = 14; + const n = @as(usize, cols) * rows; + const gpa = testing.allocator; + + const cells = try gpa.alloc(pardes.Cell, n); + defer gpa.free(cells); + const buf = try gpa.alloc(u8, frameBound(cols, rows)); + defer gpa.free(buf); + + // Every cell painted with a plain ASCII glyph and default colors, which is + // what a pardes frame overwhelmingly is: eight bytes a cell. + for (cells) |*c| c.* = .{ .text = "a".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; + const full = (try encodeFrame(buf, cols, rows, null, cells, &.{})).len; + try testing.expectEqual(@as(usize, header_len + frame_head + run_header + n * 8), full); + try testing.expectEqual(@as(usize, 6298), full); + + const before = try gpa.dupe(pardes.Cell, cells); + defer gpa.free(before); + // A keystroke: one glyph replaced, and the cursor's old and new cells + // repainted. Three cells, adjacent enough to coalesce into two runs. + cells[300] = .{ .text = "b".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; + cells[301] = .{ .text = "c".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; + cells[500] = .{ .text = "d".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; + const diff = (try encodeFrame(buf, cols, rows, null, cells, before)).len; + try testing.expectEqual(@as(usize, header_len + frame_head + 2 * run_header + 3 * 8), diff); + try testing.expectEqual(@as(usize, 56), diff); + // The whole reason this codec exists rather than shipping the Surface. + try testing.expect(full / diff > 100); +} + +test "detached wire: a run is coalesced across a gap only when that is cheaper" { + const cols: u16 = 8; + const rows: u16 = 1; + var cells: [8]pardes.Cell = @splat(.{}); + var prev: [8]pardes.Cell = @splat(.{}); + var buf: [512]u8 = undefined; + + // Two changes with three UNPAINTED cells between them. Re-sending those is + // three bytes; a second run header is six. One run. + cells[0] = .{ .text = "a".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; + cells[4] = .{ .text = "b".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; + { + const f = (try framed(try encodeFrame(&buf, cols, rows, null, &cells, &prev))).?; + try testing.expectEqual(@as(u32, 1), (try decodeServer(f.tag, f.payload)).frame.nruns); + } + // Now put a PAINTED cell in the gap that both sides already agree about. + // Eight bytes to re-send against six for a header: two runs. + const painted: pardes.Cell = .{ .text = "-".* ++ @as([6]u8, @splat(0)), .len = 1, .default = false }; + cells[2] = painted; + prev[2] = painted; + { + const f = (try framed(try encodeFrame(&buf, cols, rows, null, &cells, &prev))).?; + try testing.expectEqual(@as(u32, 2), (try decodeServer(f.tag, f.payload)).frame.nruns); + } + // Whichever it chose, the receiver ends up with the same grid. + var mirror: [8]pardes.Cell = prev; + const f = (try framed(try encodeFrame(&buf, cols, rows, null, &cells, &prev))).?; + try (try decodeServer(f.tag, f.payload)).frame.apply(&mirror); + try expectGridEqual(&cells, &mirror); +} + +test "detached wire: a truncated frame is refused at every length" { + var buf: [4096]u8 = undefined; + var scratch: Scratch = .{}; + + // Every prefix of a real message. The framing must say "not yet" for the + // ones that are short, and the decoder must say "truncated" for a payload + // whose header lies about its length — never read past the slice. + const whole = try encodeClient(&buf, .{ .event = .{ .key = .{ .cp = 'a', .text = "a" } } }); + var copy: [64]u8 = undefined; + @memcpy(copy[0..whole.len], whole); + for (header_len + 1..whole.len) |cut| { + try testing.expectEqual(@as(?Framed, null), try framed(copy[0..cut])); + // ...and the same bytes handed to the decoder with the header's length + // left claiming the whole message, which is how a decoder is walked + // off the end of its buffer. + try testing.expectError(error.Truncated, decodeClient(copy[0], copy[header_len..cut], &scratch)); + } + // Shorter than the header itself is not yet a message at all. + for (0..header_len + 1) |cut| + try testing.expectEqual(@as(?Framed, null), try framed(copy[0..cut])); + // A payload with bytes LEFT OVER is refused too: it is not this message. + try testing.expectError(error.Trailing, decodeClient(@intFromEnum(ClientTag.tick), "x", &scratch)); + try testing.expectError(error.Trailing, decodeServer(@intFromEnum(ServerTag.quit), "x")); +} + +test "detached wire: an unknown tag is refused, never guessed" { + var scratch: Scratch = .{}; + // 0x00 and 0xff have never been assigned, and 0x0f sits in the gap between + // the session tags and the input tags. All three are the same answer. + for ([_]u8{ 0x00, 0x0f, 0x1f, 0xff }) |tag| { + try testing.expectError(error.BadTag, decodeClient(tag, "", &scratch)); + try testing.expectError(error.BadTag, decodeServer(tag, "")); + } + // A tag NESTED in a payload gets the same treatment: a color, an + // underline style, a mouse button, a refusal reason. + try testing.expectError(error.BadTag, decodeServer(@intFromEnum(ServerTag.refuse), &.{0x7f})); + try testing.expectError(error.BadTag, decodeClient( + @intFromEnum(ClientTag.mouse), + &.{ 0x09, 0x00, 0, 0, 0, 0, 0 }, + &scratch, + )); +} + +test "detached wire: an over-long length prefix is refused before it is believed" { + // The prefix a hostile or corrupt peer sends to make the receiver allocate + // or index by it. `framed` must refuse rather than wait for 4 GiB. + var head: [header_len]u8 = .{ @intFromEnum(ClientTag.paste), 0, 0, 0, 0 }; + std.mem.writeInt(u32, head[1..5], max_payload + 1, .little); + try testing.expectError(error.Overlong, framed(&head)); + std.mem.writeInt(u32, head[1..5], std.math.maxInt(u32), .little); + try testing.expectError(error.Overlong, framed(&head)); + // At the boundary it is a legal prefix and simply has not all arrived. + std.mem.writeInt(u32, head[1..5], max_payload, .little); + try testing.expectEqual(@as(?Framed, null), try framed(&head)); + + // An INNER length prefix, past the payload it sits in but inside the + // protocol's cap — the one an outer-frame check cannot catch. + var scratch: Scratch = .{}; + var paste: [8]u8 = undefined; + std.mem.writeInt(u32, paste[0..4], 4096, .little); + @memcpy(paste[4..8], "abcd"); + try testing.expectError(error.Truncated, decodeClient(@intFromEnum(ClientTag.paste), &paste, &scratch)); + // ...and past the cap, which is refused rather than read. + std.mem.writeInt(u32, paste[0..4], max_payload + 1, .little); + try testing.expectError(error.Overlong, decodeClient(@intFromEnum(ClientTag.paste), &paste, &scratch)); +} + +test "detached wire: a frame that lies about its runs cannot walk out of the grid" { + const cols: u16 = 4; + const rows: u16 = 2; + var grid: [8]pardes.Cell = @splat(.{}); + + // start = 6, count = 4 on an 8-cell grid: the check that matters, and the + // one an `start + count > len` spelling would miss on overflow. + var runs: [10]u8 = undefined; + std.mem.writeInt(u32, runs[0..4], 6, .little); + std.mem.writeInt(u16, runs[4..6], 4, .little); + runs[6] = 1; + const f: Frame = .{ .kind = .diff, .cols = cols, .rows = rows, .cursor = null, .nruns = 1, .runs = runs[0..7] }; + try testing.expectError(error.Overlong, f.apply(&grid)); + + // A start past the end entirely, and the 32-bit wrap. + std.mem.writeInt(u32, runs[0..4], std.math.maxInt(u32) - 1, .little); + std.mem.writeInt(u16, runs[4..6], 4, .little); + try testing.expectError(error.Overlong, (Frame{ + .kind = .diff, + .cols = cols, + .rows = rows, + .cursor = null, + .nruns = 1, + .runs = runs[0..7], + }).apply(&grid)); + + // A zero-length run says nothing and is not something the encoder emits. + std.mem.writeInt(u32, runs[0..4], 0, .little); + std.mem.writeInt(u16, runs[4..6], 0, .little); + try testing.expectError(error.BadValue, (Frame{ + .kind = .diff, + .cols = cols, + .rows = rows, + .cursor = null, + .nruns = 1, + .runs = runs[0..6], + }).apply(&grid)); + + // A run count larger than the grid has cells is refused at DECODE, before + // it is used to bound the walk. + var head: [frame_head]u8 = undefined; + var w: Writer = .init(&head); + try w.putByte(@intFromEnum(FrameKind.diff)); + try w.putU16(cols); + try w.putU16(rows); + try putCursor(&w, null, cols, rows); + try w.putU32(9); + try testing.expectError(error.Overlong, decodeServer(@intFromEnum(ServerTag.frame), w.written())); + + // ...and a grid of the wrong size is the receiver's own bug, not a frame + // it should paint half of. + var small: [4]pardes.Cell = @splat(.{}); + try testing.expectError(error.BadValue, (Frame{ + .kind = .full, + .cols = cols, + .rows = rows, + .cursor = null, + .nruns = 0, + .runs = "", + }).apply(&small)); +} + +test "detached wire: a cursor outside the grid is refused, not painted" { + // The one field of a frame a frontend indexes with rather than copies: + // tty.zig moves the terminal's own cursor to `cursor.x`/`cursor.y`, and + // the ESP32-P4 panel writes the cell there. A frame that says 4x2 and puts + // the cursor at (9,0) is an out-of-bounds write in every frontend, so it + // has to be a decode error and not a clamp — a clamped cursor is a wrong + // screen that nobody reports. + const cols: u16 = 4; + const rows: u16 = 2; + var head: [frame_head]u8 = undefined; + for ([_][2]u16{ .{ cols, 0 }, .{ 0, rows }, .{ 0xffff, 0xffff } }) |xy| { + var w: Writer = .init(&head); + try w.putByte(@intFromEnum(FrameKind.diff)); + try w.putU16(cols); + try w.putU16(rows); + // Written by hand: `putCursor` now refuses this too, which is the + // other half of the same fix and is asserted below. + try w.putBool(true); + try w.putU16(xy[0]); + try w.putU16(xy[1]); + try w.putBool(false); + try w.putU32(0); + try testing.expectError( + error.BadValue, + decodeServer(@intFromEnum(ServerTag.frame), w.written()), + ); + } + // The last cell IS in the grid, and a frame carrying it decodes. + { + var w: Writer = .init(&head); + try w.putByte(@intFromEnum(FrameKind.diff)); + try w.putU16(cols); + try w.putU16(rows); + try putCursor(&w, .{ .x = cols - 1, .y = rows - 1, .bar = true }, cols, rows); + try w.putU32(0); + const got = (try decodeServer(@intFromEnum(ServerTag.frame), w.written())).frame; + try testing.expectEqual(@as(u16, cols - 1), got.cursor.?.x); + try testing.expectEqual(@as(u16, rows - 1), got.cursor.?.y); + } + // ...and the ENCODER refuses to build one, so a core bug is a dropped + // frame with a log line rather than a frame every frontend hangs up over. + var cells: [8]pardes.Cell = @splat(.{}); + var buf: [512]u8 = undefined; + try testing.expectError( + error.BadValue, + encodeFrame(&buf, cols, rows, .{ .x = cols, .y = 0, .bar = false }, &cells, &.{}), + ); + // A grid past the protocol's own ceiling is refused by the encoder for the + // same reason: `max_payload` is derived from it, so the decoder would. + try testing.expectError( + error.Overlong, + encodeFrame(&buf, max_cols + 1, 1, null, cells[0..0], &.{}), + ); +} + +test "detached wire: values a field cannot mean are refused" { + var scratch: Scratch = .{}; + + // A bool is 0 or 1. `2` used to be "true" in every hand-written codec that + // ever silently accepted a corrupt stream. + try testing.expectError(error.BadValue, decodeClient( + @intFromEnum(ClientTag.key), + &.{ 'a', 0, 0, 0, 0, 0, 2, 0, 0 }, + &scratch, + )); + // A codepoint past Unicode's last: @intCast into Key.cp's u21 would panic. + try testing.expectError(error.BadValue, decodeClient( + @intFromEnum(ClientTag.key), + &.{ 0x00, 0x00, 0x11, 0x00, 0, 0, 0, 0, 0 }, + &scratch, + )); + // A pane the core cannot index. + try testing.expectError(error.BadValue, decodeClient( + @intFromEnum(ClientTag.eof), + &.{pardes.MAX_PANES}, + &scratch, + )); + // A zero-column grid would collapse a shared session; an over-wide one is + // past what this protocol carries. + try testing.expectError(error.BadValue, decodeClient( + @intFromEnum(ClientTag.hello), + &.{ 1, 0, 0, 0, 24, 0 }, + &scratch, + )); + { + var hello: [6]u8 = undefined; + std.mem.writeInt(u16, hello[0..2], version, .little); + std.mem.writeInt(u16, hello[2..4], max_cols + 1, .little); + std.mem.writeInt(u16, hello[4..6], 24, .little); + try testing.expectError(error.BadValue, decodeClient(@intFromEnum(ClientTag.hello), &hello, &scratch)); + } + // A NaN scroll distance. `pinch` multiplies into a zoom the pane keeps. + { + var pinch: [4]u8 = undefined; + std.mem.writeInt(u32, &pinch, @bitCast(std.math.nan(f32)), .little); + try testing.expectError(error.BadValue, decodeClient(@intFromEnum(ClientTag.pinch), &pinch, &scratch)); + std.mem.writeInt(u32, &pinch, @bitCast(std.math.inf(f32)), .little); + try testing.expectError(error.BadValue, decodeClient(@intFromEnum(ClientTag.touch_scroll), &pinch, &scratch)); + } + // More pipe outputs than the core has selections to produce them. + { + var head: [7]u8 = undefined; + std.mem.writeInt(u32, head[0..4], 1, .little); + head[4] = 1; + std.mem.writeInt(u16, head[5..7], pardes.MAX_SELS + 1, .little); + try testing.expectError(error.Overlong, decodeClient(@intFromEnum(ClientTag.pipe_resp), &head, &scratch)); + } + // A cell with the reserved attribute bit set, and one with an impossible + // grapheme length. Both are bytes this protocol has no meaning for. + { + var grid: [1]pardes.Cell = @splat(.{}); + // default=0, len=1, 'x', fg default, bg default, attrs=0x80, ul, font + var runs = [_]u8{ 0, 0, 0, 0, 1, 0, 0, 1, 'x', 0, 0, 0x80, 0, 0 }; + const f: Frame = .{ .kind = .diff, .cols = 1, .rows = 1, .cursor = null, .nruns = 1, .runs = &runs }; + try testing.expectError(error.BadValue, f.apply(&grid)); + runs[11] = 0; + runs[7] = 8; // len past Cell.text + try testing.expectError(error.BadValue, f.apply(&grid)); + runs[7] = 0; // ...and a grapheme of no bytes at all + try testing.expectError(error.BadValue, f.apply(&grid)); + } +} + +test "detached wire: max_payload bounds every message this protocol can build" { + // The derivation in `max_payload`'s comment, asserted: the worst frame + // this protocol admits fits, so a receiver sized for max_payload can + // always hold one. + try testing.expect(frameBound(max_cols, max_rows) <= max_payload); + // ...and a 4 MiB paste, which is the tty frontend's own cap. + try testing.expect((4 << 20) + header_len + 4 <= max_payload); + // A slice past the width of its own length prefix is refused by the + // ENCODER, rather than written and rejected at the far end. A command line + // is u16-prefixed because nothing legitimately types 64 KiB of one. + const gpa = testing.allocator; + const huge = try gpa.alloc(u8, std.math.maxInt(u16) + 1); + defer gpa.free(huge); + @memset(huge, 'x'); + const room = try gpa.alloc(u8, huge.len + header_len + 2); + defer gpa.free(room); + try testing.expectError(error.Overlong, encodeClient(room, .{ .event = .{ .command = huge } })); + // ...and a buffer the caller sized too small is NoSpace, which is the + // server's signal to drop that one message rather than the client. + var small: [8]u8 = undefined; + try testing.expectError(error.NoSpace, encodeClient(&small, .{ .event = .{ .paste = "0123456789" } })); +} |
