const std = @import("std"); const libc = std.c; const builtin = @import("builtin"); const ninep = @import("9p.zig"); const cloud9 = @import("cloud9"); const transport = cloud9.transport; const pardes = @import("pardes.zig"); const limits = @import("memory.zig").limits; pub const quic_enabled = @import("9p_options").quic; const quic = if (quic_enabled) @import("9p_quic.zig") else struct {}; const log = std.log.scoped(.ninep); pub const darwin = switch (builtin.os.tag) { .macos, .ios, .tvos, .watchos, .visionos => true, else => false, }; pub const supported = builtin.os.tag == .linux or darwin; pub const sun_path_len = @typeInfo(@FieldType(libc.sockaddr.un, "path")).array.len; pub fn setCloexec(fd: c_int) void { _ = libc.fcntl(fd, libc.F.SETFD, @as(c_int, 1)); } pub fn socketDir(buf: *[sun_path_len:0]u8) ?[:0]const u8 { if (libc.getenv("XDG_RUNTIME_DIR")) |path| return std.fmt.bufPrintSentinel(buf, "{s}", .{std.mem.span(path)}, 0) catch null; const home = libc.getenv("HOME") orelse return null; return std.fmt.bufPrintSentinel(buf, "{s}/.local/state/pardes", .{std.mem.span(home)}, 0) catch null; } pub const FileFacts = struct { mode: u32, uid: libc.uid_t }; pub fn statNoFollow(path: [:0]const u8) ?FileFacts { if (comptime darwin) { var stat: libc.Stat = undefined; if (libc.fstatat(libc.AT.FDCWD, path, &stat, libc.AT.SYMLINK_NOFOLLOW) != 0) return null; return .{ .mode = stat.mode, .uid = stat.uid }; } else { const linux = std.os.linux; var stat: linux.Statx = undefined; const fields: linux.STATX = .{ .TYPE = true, .MODE = true, .UID = true }; if (libc.statx(linux.AT.FDCWD, path, linux.AT.SYMLINK_NOFOLLOW, fields, &stat) != 0) return null; return .{ .mode = stat.mode, .uid = stat.uid }; } } pub fn ensureSocketDir(dir: [:0]const u8) bool { if (dir.len == 0) return false; var partial: [sun_path_len:0]u8 = undefined; @memcpy(partial[0 .. dir.len + 1], dir[0 .. dir.len + 1]); for (1..dir.len) |i| { if (dir[i] != '/') continue; partial[i] = 0; _ = libc.mkdir(partial[0..i :0], 0o700); partial[i] = '/'; } _ = libc.mkdir(dir, 0o700); const stat = statNoFollow(dir) orelse return false; return stat.mode & 0o170000 == 0o040000 and stat.uid == libc.getuid() and stat.mode & 0o077 == 0; } const prefix = "pardes-9p-"; pub const msize = ninep.msize; /// Scripts, a mount and an agent or two at once; each slot holds its /// buffers (a few msize) whether used or not. pub const max_conns = 16; extern "c" fn inet_pton(family: c_int, src: [*:0]const u8, dst: *anyopaque) c_int; pub fn networkAddress(dial: []const u8, allow_zero_port: bool) error{BadDial}!std.Io.net.IpAddress { if (comptime !supported) return error.BadDial; const host_start: usize = if (std.mem.startsWith(u8, dial, "tcp!")) 4 else if (std.mem.startsWith(u8, dial, "quic!")) 5 else return error.BadDial; const split = std.mem.lastIndexOfScalar(u8, dial, '!') orelse return error.BadDial; if (split <= host_start or split + 1 == dial.len) return error.BadDial; const port_text = dial[split + 1 ..]; for (port_text) |c| if (c < '0' or c > '9') return error.BadDial; const port = std.fmt.parseInt(u16, port_text, 10) catch return error.BadDial; if (port == 0 and !allow_zero_port) return error.BadDial; const host = dial[host_start..split]; if (std.mem.indexOfScalar(u8, host, 0) != null) return error.BadDial; var host_buf: [46]u8 = undefined; const host_z = std.fmt.bufPrintSentinel(&host_buf, "{s}", .{host}, 0) catch return error.BadDial; var ip4: std.Io.net.Ip4Address = .{ .port = port, .bytes = undefined }; if (inet_pton(libc.AF.INET, host_z, &ip4.bytes) == 1) return .{ .ip4 = ip4 }; var ip6: std.Io.net.Ip6Address = .{ .port = port, .bytes = undefined }; if (inet_pton(libc.AF.INET6, host_z, &ip6.bytes) == 1) return canonicalIp(.{ .ip6 = ip6 }); return error.BadDial; } fn canonicalIp(address: std.Io.net.IpAddress) std.Io.net.IpAddress { if (address == .ip6) { if (std.Io.net.Ip4Address.fromIp6(address.ip6)) |ip4| return .{ .ip4 = ip4 }; } return address; } fn sockaddrIp(address: *const libc.sockaddr) ?std.Io.net.IpAddress { return switch (address.family) { libc.AF.INET => blk: { const addr: *const libc.sockaddr.in = @ptrCast(@alignCast(address)); break :blk .{ .ip4 = .{ .port = std.mem.bigToNative(u16, addr.port), .bytes = @bitCast(addr.addr) } }; }, libc.AF.INET6 => blk: { const addr: *const libc.sockaddr.in6 = @ptrCast(@alignCast(address)); break :blk canonicalIp(.{ .ip6 = .{ .port = std.mem.bigToNative(u16, addr.port), .bytes = addr.addr } }); }, else => null, }; } const IfAddr = extern struct { next: ?*IfAddr, name: ?[*:0]u8, flags: c_uint, address: ?*libc.sockaddr, netmask: ?*libc.sockaddr, destination: ?*libc.sockaddr, data: ?*anyopaque, }; extern "c" fn getifaddrs(out: *?*IfAddr) c_int; extern "c" fn freeifaddrs(first: *IfAddr) void; fn localIp(address: std.Io.net.IpAddress) bool { switch (address) { .ip4 => |ip| if (ip.bytes[0] == 127 or std.mem.allEqual(u8, &ip.bytes, 0)) return true, .ip6 => |ip| if (ip.isLoopBack() or std.mem.allEqual(u8, &ip.bytes, 0)) return true, } var first: ?*IfAddr = null; if (getifaddrs(&first) != 0) return false; defer if (first) |head| freeifaddrs(head); var next = first; while (next) |entry| : (next = entry.next) { var local = sockaddrIp(entry.address orelse continue) orelse continue; local.setPort(address.getPort()); if (address.eql(&local)) return true; } return false; } const Srv = ninep.Server(pardes.ctlfs, ninep.editor); /// The Unix and TCP listeners run on cloud9's `std.Io` runner: it accepts, /// reads frames, steps each connection's engine and writes replies on its /// own tasks, and every backend request is answered right there, on the /// connection's task, by taking the editor's turn (`pardes.turn`): while the /// editor waits for input, or while a step of it is out in a syscall -- /// which is what lets the editor read its own tree through a mount. QUIC /// keeps the poll loop below, on the editor's thread. const Runner = if (supported) cloud9.serve.Runner(pardes.ctlfs, ninep.editor, .{ .msize = msize, .connections = max_conns, .listeners = 2, }) else void; const quic_slots = if (quic_enabled) max_conns else 0; /// A QUIC connection on the poll path. const Conn = struct { quic: if (quic_enabled) ?quic.Connection else void = if (quic_enabled) null else {}, draining: bool = false, accepted_ms: i64 = 0, srv: Srv = undefined, in: [msize]u8 = undefined, out: [2 * msize]u8 = undefined, }; const accept_pause_ms: i64 = 100; pub const Listener = struct { io: std.Io, /// The core that answers; `reset` puts a replacement in. core: *pardes.Pardes, runner: Runner = undefined, stopping: std.atomic.Value(bool) = .init(false), /// A connection wrote the Restore and has yet to answer it. restore_writer: std.atomic.Value(bool) = .init(false), /// Clients turned away with every slot taken, since the log last said so. refused: std.atomic.Value(u32) = .init(0), /// Connections accepted so far, and the count each slot's connection /// was accepted at: a Restore cuts those accepted before it (`reset`), /// and a new client in a freed slot has a later stamp. Written on the /// connection's task as it opens, read with the turn. accepts: std.atomic.Value(u64) = .init(0), accepted: [max_conns]std.atomic.Value(u64) = @splat(.init(0)), /// Connections accepted before this count are being cut by a Restore. cut_before: u64 = 0, tcp_address: ?std.Io.net.IpAddress = null, quic: if (quic_enabled) ?quic.Listener else void = if (quic_enabled) null else {}, quic_address: ?std.Io.net.IpAddress = null, paused_ms: i64 = 0, path_buf: [sun_path_len]u8 = undefined, path_len: usize = 0, /// The registry entry this editor posted /// (`$XDG_RUNTIME_DIR/9p/pardes/`), zero when it did not. posted_buf: [sun_path_len]u8 = undefined, posted_len: usize = 0, conns: [quic_slots]Conn = @splat(.{}), control: [2]c_int = .{ -1, -1 }, watcher: ?std.Thread = null, watch_stop: std.atomic.Value(bool) = .init(false), watch_lock: std.atomic.Mutex = .unlocked, watch_fds: [2]libc.pollfd = undefined, watch_len: usize = 0, watch_timeout: c_int = -1, wake_ctx: ?*anyopaque = null, wake: ?*const fn (?*anyopaque) void = null, pub fn path(l: *const Listener) []const u8 { return l.path_buf[0..l.path_len]; } fn of(ctx: ?*anyopaque) *Listener { return @ptrCast(@alignCast(ctx.?)); } /// The runner's handler, on the connection's task: takes the turn and /// answers. A request that would change a pane while the editor is out /// in a syscall mid-step is parked in the engine instead, and retried /// when the turn is next given up between steps (`wakeParked`). /// On the accepting task, without the turn: counted, and logged with /// the turn (`answerHeld`). fn onRefused(ctx: ?*anyopaque) void { const l = of(ctx); _ = l.refused.fetchAdd(1, .acq_rel); l.kick(); } fn onOpened(ctx: ?*anyopaque, conn: *Runner.Conn) void { const l = of(ctx); l.accepted[conn.index].store(l.accepts.fetchAdd(1, .acq_rel) + 1, .release); } /// A live connection accepted before the Restore under way. fn cutting(l: *Listener, conn: *Runner.Conn) bool { return conn.live() and l.accepted[conn.index].load(.acquire) <= l.cut_before and l.cut_before != 0; } fn onServe(ctx: ?*anyopaque, conn: *Runner.Conn, req: pardes.ctlfs.Req) void { const l = of(ctx); // Once `stop` has begun the editor is tearing down and may never // rest again: answer without waiting on it. if (l.stopping.load(.acquire)) { const refused = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); return conn.reply(&refused, ""); } const quiet = pardes.turn.take(); defer pardes.turn.give(); // A connection a Restore is cutting (`reset`) is the old editor's // client: its fids name the old editor's panes and opens. if (l.cutting(conn)) { const refused = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); return conn.reply(&refused, ""); } if (!quiet and pardes.ctlfs.needsQuiet(l.core, req)) { pardes.turn.parked = true; const later: pardes.ctlfs.Reply = .{ .tag = req.tag, .status = .again }; return conn.reply(&later, ""); } const core = l.core; core.fs.write_room = if (conn.engine.protocol.msize > 24) conn.engine.protocol.msize - 24 else 0; const epoch = pardes.turn.epoch; const restores = pardes.turn.restores; // Asked first: a release that runs a held line has none after. const changes = pardes.ctlfs.changesPane(core, req); const reply = core.serveFs(req); if (req.op == .release and quiet) l.collectOs(); // A read with nothing yet stays parked in the engine, and the core // keeps its ticket to answer it when what it waits on has something. if (reply.status == .again and req.op == .read) { // Only an open's record can hold a read; one that waits with // none would sit parked until some unrelated write parks. const o = pardes.ctlfs.openOf(core, req) orelse { log.err("a read of node {x} waits with no open record to hold it", .{req.node}); std.debug.assert(false); return conn.reply(&reply, ""); }; // One read waits on an open at a time, as acme's window keeps // one `eventx`; a second is refused rather than left parked // where nothing would ever answer it. The first asked again by // a retry is the same request, and holds its place. if (o.held) |held| if (held.req.tag != req.tag) { const other: *Runner.Conn = @ptrCast(@alignCast(held.asker)); if (other.waiting(held.ticket)) return conn.reply(&pardes.ctlfs.failText(req.tag, pardes.ctlfs.E.BUSY, pardes.ctlfs.e_in_use), ""); }; o.held = null; const ticket = conn.hold(&reply) orelse return; o.held = .{ .asker = conn, .req = req, .ticket = ticket }; return; } if (!changes) return conn.reply(&reply, core.fsPayload(reply)); // The editor draws the change and performs what it asked for; when // it asked for something -- a save, a shell, a watch -- the answer // waits until that is done, so `echo Save > exec` returns with the // file written. A change carries no payload, so the reply is still // whole after the wait. l.kick(); // A refusal answers at once: its text may be in the core's one // buffer for it, which a request run while this one waited would // write over. if (reply.status == .err) return conn.reply(&reply, ""); // A write that quits the editor (Kill) is answered now: the editor // is on its way out and will settle nothing this could wait for, // and `deinit` lets the answer out before it cuts the connections. if (core.quit) return conn.reply(&reply, ""); const restoring = core.restore_req != null; // Set with the turn still held, so `reset` sees it and waits. if (restoring) l.restore_writer.store(true, .release); defer if (restoring) l.restore_writer.store(false, .release); core.fs.late_failure_len = 0; // What fails while this waits is this write's err, never a msg too. core.fs.write_waits = true; if (core.effects_len != 0) pardes.turn.awaitSettled(epoch); if (core.fs.lsp_answer_at) |n| { core.fs.lsp_answer_at = null; pardes.turn.awaitLsp(n); } core.fs.write_waits = false; // ponytail: one slot, so a failure of another client's effects that // settle in the same wait is told to this write too. if (l.core == core and core.fs.late_failure_len != 0) { const late = core.fs.late_failure[0..core.fs.late_failure_len]; // What is not there (`no such directory`) is ENOENT, as a // builtin's failure saying so is (ctl.failureErrno). const errno = if (std.mem.indexOf(u8, late, "no such") != null or std.mem.indexOf(u8, late, "not found") != null) pardes.ctlfs.E.NOENT else if (std.mem.indexOf(u8, late, "no space") != null) pardes.ctlfs.E.NOSPC else pardes.ctlfs.E.IO; const failed = pardes.ctlfs.failText(req.tag, errno, late); // Its err record says it (the path in it), the one record: it // was said with no msg while this waited (fs.write_waits). pardes.ctlfs.events.noteError(core, req, failed); return conn.reply(&failed, ""); } // The Restore's own write is answered once it is done, and before // `reset` hangs every connection up, this one too, so the writer // hears that it happened rather than a cut; a failed one leaves the // connection and says why on the message row and in the log. if (restoring) { if (l.core == core) pardes.turn.awaitRestored(restores); return conn.reply(&reply, ""); } // Once a wait returns `core` may be gone: a Restore meanwhile put a // replacement in and is hanging this connection up. if (l.core != core) return; conn.reply(&reply, ""); } /// With the turn, as it is given up: answers each read the core holds /// whose file now has something, on the connection that asked, with no /// retry of anything else parked there. Only while its ticket still /// waits, asked in the same hold of the engine as the answer: Tflush /// and clunk drop a park without a word to the backend, a `retry` takes /// one out to ask again (and that asking answers it), and a record spent /// on any of them would be lost to the reader that comes next. fn answerHeld(ctx: ?*anyopaque) void { if (comptime !supported) return; const l = of(ctx); if (l.stopping.load(.acquire)) return; const core = l.core; const refused = l.refused.swap(0, .acq_rel); if (refused != 0) { var why: [64]u8 = undefined; pardes.ctlfs.events.notePath(core, "err -", std.fmt.bufPrint(&why, "9p: too many connections ({d} turned away)", .{refused}) catch "9p: too many connections"); } if (!core.fs.news) return; core.fs.news = false; for (&core.fs.opens) |*o| { const held = o.held orelse continue; const conn: *Runner.Conn = @ptrCast(@alignCast(held.asker)); if (!conn.answerWith(held.ticket, Held{ .core = core, .req = held.req }, Held.make)) o.held = null; } } /// A held read asked again, to make its answer while its park waits. const Held = struct { core: *pardes.Pardes, req: pardes.ctlfs.Req, fn make(h: Held) ?Runner.Conn.Answer { const reply = h.core.serveFs(h.req); if (reply.status == .again) return null; return .{ .reply = reply, .bytes = h.core.fsPayload(reply) }; } }; /// The turn was given up quiet with a request parked: every connection /// retries what it parked. fn wakeParked(ctx: ?*anyopaque) void { const l = of(ctx); if (l.stopping.load(.acquire)) return; if (comptime supported) l.runner.wakeAll(); } /// A request changed the core: the editor wakes to draw it and perform /// what it asked for. fn kick(l: *Listener) void { if (l.stopping.load(.acquire)) return; if (l.wake) |f| f(l.wake_ctx); } /// Accepts and reads QUIC connections (the runner does its own). pub fn accept(l: *Listener) void { if (comptime !quic_enabled) return; if (l.quic) |*listener| for (0..max_conns + 1) |_| { var connection = (listener.accept() catch |err| { log.warn("QUIC accept failed: {s}", .{@errorName(err)}); return; }) orelse break; const c = for (&l.conns, 0..) |*cand, i| { if (!l.live(@intCast(i)) and !cand.draining) break cand; } else { connection.deinit(); continue; }; c.quic = connection; c.draining = false; c.accepted_ms = nowMs(); c.srv = .init(.{ .in = &c.in, .out = &c.out, .root = pardes.ctlfs.root }); }; } pub const greet_deadline_ms: i64 = 5000; pub fn expire(l: *Listener) void { if (comptime !quic_enabled) return; const now = nowMs(); if (now == 0) return; for (&l.conns, 0..) |*c, i| { if (!l.live(@intCast(i)) or c.srv.protocol.msize != 0) continue; if (now - c.accepted_ms < greet_deadline_ms) continue; log.debug("QUIC slot {d} never sent Tversion; taking it back", .{i}); l.drop(@intCast(i)); } } pub fn nextDue(l: *const Listener) ?i32 { if (comptime !quic_enabled) return null; const now = nowMs(); if (now == 0) return null; var due: ?i64 = null; if (l.quic) |*listener| if (listener.nextDue()) |ms| { const at = now + ms; due = if (due) |d| @min(d, at) else at; }; for (&l.conns, 0..) |*c, i| { if (!l.live(@intCast(i))) continue; if (c.quic) |*connection| if (connection.pending()) return 0; if (c.srv.protocol.msize != 0) continue; const at = c.accepted_ms + greet_deadline_ms; due = if (due) |d| @min(d, at) else at; } const at = due orelse return null; return @intCast(@max(0, at - now)); } const nowMs = transport.nowMs; pub fn fill(l: *Listener, i: u8) void { if (comptime !quic_enabled) return; const c = &l.conns[i]; if (c.srv.protocol.dead) return l.drop(i); const room = c.srv.protocol.in.len - c.srv.protocol.in_len; if (room == 0) return; var buf: [msize]u8 = undefined; const got = (c.quic.?.read(buf[0..@min(room, buf.len)]) catch return l.drop(i)) orelse return; if (got == 0) return l.drop(i); const n = c.srv.push(buf[0..@intCast(got)]); if (c.srv.protocol.dead) return l.drop(i); std.debug.assert(n == @as(usize, @intCast(got))); } pub fn flush(l: *Listener, i: u8) void { if (comptime !quic_enabled) return; const c = &l.conns[i]; if (!l.live(i)) return; while (true) { const bytes = c.srv.output(); if (bytes.len == 0) return; const n = c.quic.?.write(bytes) catch return l.drop(i); if (n == 0) return; c.srv.wrote(@intCast(n)); } } pub fn live(l: *const Listener, i: u8) bool { if (comptime !quic_enabled) return false; return l.conns[i].quic != null; } pub fn drop(l: *Listener, i: u8) void { if (comptime !quic_enabled) return; const c = &l.conns[i]; if (c.quic) |*connection| { connection.deinit(); c.quic = null; } if (c.draining or c.accepted_ms == 0) return; c.srv.hangup(); c.draining = true; } /// What one QUIC poll pass did: `pending` says there is more to do at /// once. pub const Polled = struct { count: usize = 0, pending: bool = false }; /// Answers one QUIC request, on the editor's thread. fn step(l: *Listener, srv: *Srv, req: pardes.ctlfs.Req) void { l.core.fs.write_room = if (srv.protocol.msize > 24) srv.protocol.msize - 24 else 0; const reply = l.core.serveFs(req); srv.reply(&reply, l.core.fsPayload(reply)); if (req.op == .release) l.collectOs(); } /// QUIC's events, accepts, reads, requests and replies, on the editor's /// thread after each of its steps; `pending` says there is more to do /// right away. Unix and TCP need none of this. pub fn tick(l: *Listener) Polled { var result: Polled = .{}; // Between the editor's steps, which is when the table may be swept: // the editor's own directory listings leave nodes in it too. l.collectOs(); if (comptime quic_enabled) { if (l.quic) |*listener| listener.events() catch |err| { log.warn("QUIC listener stopped: {s}", .{@errorName(err)}); for (&l.conns, 0..) |*conn, i| if (conn.quic != null) l.drop(@intCast(i)); listener.deinit(); l.quic = null; l.quic_address = null; }; l.accept(); for (0..quic_slots) |i| if (l.live(@intCast(i))) l.fill(@intCast(i)); l.expire(); for (&l.conns, 0..) |*conn, i| { if (!l.live(@intCast(i)) and !conn.draining) continue; var count: usize = 0; while (conn.srv.retry()) |req| { l.step(&conn.srv, req); count += 1; } while (count < 64) { const req = conn.srv.next() orelse break; l.step(&conn.srv, req); count += 1; } result.count += count; result.pending = result.pending or count >= 64; if (conn.draining) { if (count == 0) conn.draining = false else result.pending = true; } else l.flush(@intCast(i)); if (conn.quic) |*connection| result.pending = result.pending or connection.pending(); } l.arm(); } return result; } /// Hangs every connection up, so that `replacement` starts with no /// client holding anything. From here requests are answered from the /// replacement -- a client that connects meanwhile is a client of the /// replacement -- and every task waiting on the old core is let go: the /// Restore's own writer answers, the rest find the core changed and /// answer nothing, their connections being hung up. The connections pay their releases on their own tasks, so /// the editor rests while they do. pub fn reset(l: *Listener, replacement: *pardes.Pardes) void { for (0..quic_slots) |i| l.drop(@intCast(i)); if (comptime quic_enabled) while (l.tick().pending) {}; l.core = replacement; replacement.fs.socket_path = l.path(); replacement.fs.tcp_address = l.tcp_address; replacement.fs.quic_address = l.quic_address; if (comptime supported) { // The Restore's writer answers as it wakes (`onServe`), and its // answer gets 200 ms to leave before the cut; every other request // from the old editor's clients meanwhile is refused, since the // fids it names are the old editor's. A client that connects now // is the replacement's, and stays. l.cut_before = l.accepts.load(.acquire); pardes.turn.settle(); pardes.turn.restoreSettled(); pardes.turn.rest(); const flush_by = nowMs() +| 200; while (nowMs() < flush_by) { const pending = l.restore_writer.load(.acquire) or for (&l.runner.conns) |*conn| { if (!l.cutting(conn)) continue; conn.lock(); const n = conn.engine.output().len; conn.unlock(); if (n != 0) break true; } else false; if (!pending) break; Client.nap(1); } for (&l.runner.conns) |*conn| if (l.cutting(conn)) conn.close(); const deadline = nowMs() +| 2000; while (nowMs() < deadline) { const left = for (&l.runner.conns) |*conn| { if (l.cutting(conn)) break true; } else false; if (!left) break; Client.nap(1); } pardes.turn.wake(); l.cut_before = 0; } l.collectOs(); for (&l.conns) |*conn| conn.accepted_ms = 0; l.arm(); } /// Forgets the host paths no connection names any more. With the turn, /// and only between the editor's steps: a step out in a syscall may be /// reading one of these entries. fn collectOs(l: *Listener) void { const core = l.core; var i: usize = 0; while (i < core.fs.os_paths.items.len) { const entry = core.fs.os_paths.items[i]; if (l.references(entry.node)) { i += 1; } else { core.gpa.free(entry.path); _ = core.fs.os_paths.swapRemove(i); } } } /// Whether any connection names `node`. fn references(l: *Listener, node: u64) bool { if (comptime supported) { for (&l.runner.conns) |*conn| { if (!conn.live()) continue; conn.lock(); defer conn.unlock(); if (conn.engine.references(node)) return true; } } for (&l.conns, 0..) |*conn, j| { if (!l.live(@intCast(j)) and !conn.draining) continue; if (conn.srv.references(node)) return true; } return false; } /// Where the runner's tasks report that a request changed the core, /// and for QUIC a thread that watches its socket and calls `wake` when /// `tick` has something to read. pub fn watch(l: *Listener, ctx: ?*anyopaque, wake: *const fn (?*anyopaque) void) !void { l.wake_ctx = ctx; l.wake = wake; if (comptime !quic_enabled) return; if (l.quic == null or l.watcher != null) return; if (libc.pipe(&l.control) != 0) return error.PipeFailed; errdefer { _ = libc.close(l.control[0]); _ = libc.close(l.control[1]); l.control = .{ -1, -1 }; } for (l.control) |fd| { setCloexec(fd); _ = setNonblock(fd); } l.arm(); l.watcher = try std.Thread.spawn(.{}, watchQuic, .{l}); } fn arm(l: *Listener) void { if (comptime !quic_enabled) return; if (l.control[1] < 0) return; while (!l.watch_lock.tryLock()) std.atomic.spinLoopHint(); l.watch_fds[0] = .{ .fd = l.control[0], .events = @intCast(libc.POLL.IN), .revents = 0 }; l.watch_len = 1; if (l.quic) |*listener| { l.watch_fds[l.watch_len] = listener.poll(); l.watch_len += 1; } l.watch_timeout = l.nextDue() orelse -1; l.watch_lock.unlock(); _ = libc.write(l.control[1], "w", 1); } fn watchQuic(l: *Listener) void { var notified = false; while (!l.watch_stop.load(.acquire)) { var fds: [2]libc.pollfd = undefined; while (!l.watch_lock.tryLock()) std.atomic.spinLoopHint(); const len = if (notified) 1 else l.watch_len; @memcpy(fds[0..len], l.watch_fds[0..len]); const timeout = if (notified) -1 else l.watch_timeout; l.watch_lock.unlock(); const ready = libc.poll(&fds, @intCast(len), timeout); if (ready < 0) continue; if (fds[0].revents != 0) { var buf: [64]u8 = undefined; while (libc.read(l.control[0], &buf, buf.len) > 0) {} notified = false; continue; } if (!l.watch_stop.load(.acquire)) l.wake.?(l.wake_ctx); notified = true; } } /// Advertises this editor in the posted-9P registry so a /// `9ns --mntgen` mount lists it and dials it on a walk: /// `$XDG_RUNTIME_DIR/9p/pardes/` is a symlink to the socket /// pardes already binds. One directory for the program, one entry /// per running editor — the layout zmx uses for its sessions, so /// several editors group instead of crowding the registry root. /// /// Only the runtime-directory socket posts: an instance that fell /// back to `~/.local/state/pardes` stays out of the user's /// registry, the way a private `ZMX_DIR` does for zmx. /// /// The socket itself is still bound by `listen` above, not by /// `cloud9.post`: `post` claims flat names only, and its claim /// protocol derives its lock directory from the registry path, so /// a name inside a subdirectory cannot go through it yet. See /// docs/cloud9.md. fn postToRegistry(l: *Listener, name: []const u8) void { if (comptime !supported) return; if (name.len == 0 or std.mem.indexOfAny(u8, name, "/\x00") != null) return; const xdg_c = libc.getenv("XDG_RUNTIME_DIR") orelse return; const xdg = std.mem.span(xdg_c); if (xdg.len == 0) return; var reg_buf: [sun_path_len:0]u8 = undefined; const reg = std.fmt.bufPrintSentinel(®_buf, "{s}/9p", .{xdg}, 0) catch return; _ = libc.mkdir(reg, 0o750); var svc_buf: [sun_path_len:0]u8 = undefined; const svc = std.fmt.bufPrintSentinel(&svc_buf, "{s}/9p/pardes", .{xdg}, 0) catch return; if (libc.mkdir(svc, 0o750) != 0 and statNoFollow(svc) == null) { log.warn("registry post skipped: cannot create {s}", .{svc}); return; } sweepRegistry(l.io, svc, xdg); var entry_buf: [sun_path_len:0]u8 = undefined; const entry = std.fmt.bufPrintSentinel(&entry_buf, "{s}/{s}", .{ svc, name }, 0) catch { log.warn("registry post skipped: name too long: {s}", .{name}); return; }; const target = l.path_buf[0..l.path_len]; // Replace only what is provably not live: our own entry, or a // dead predecessor's symlink. `probe` treats uncertainty as // live, so it is only asked about an entry that *is* a symlink; // anything else is left strictly alone and the symlink below // simply fails. var link_buf: [sun_path_len]u8 = undefined; const n = libc.readlink(entry, &link_buf, link_buf.len); if (n >= 0) { const had = link_buf[0..@intCast(n)]; if (!std.mem.eql(u8, had, target) and probe(entry) == .live) { log.warn("registry entry pardes/{s} is live; not re-posted", .{name}); return; } _ = libc.unlink(entry); } var target_z: [sun_path_len:0]u8 = undefined; const target_zs = std.fmt.bufPrintSentinel(&target_z, "{s}", .{target}, 0) catch return; if (libc.symlink(target_zs, entry) != 0) { log.warn("registry post skipped: pardes/{s} is occupied", .{name}); return; } @memcpy(l.posted_buf[0..entry.len], entry); l.posted_len = entry.len; log.info("posted pardes/{s} -> {s}", .{ name, target }); } /// Removes the registry entry, but only while it is still ours: a /// name another editor has since claimed is never unlinked. fn unpostFromRegistry(l: *Listener) void { if (comptime !supported) return; if (l.posted_len == 0) return; var z: [sun_path_len:0]u8 = undefined; @memcpy(z[0..l.posted_len], l.posted_buf[0..l.posted_len]); z[l.posted_len] = 0; const entry = z[0..l.posted_len :0]; var link_buf: [sun_path_len]u8 = undefined; const n = libc.readlink(entry, &link_buf, link_buf.len); if (n >= 0 and std.mem.eql(u8, link_buf[0..@intCast(n)], l.path_buf[0..l.path_len])) { _ = libc.unlink(entry); } l.posted_len = 0; } pub fn deinit(l: *Listener, gpa: std.mem.Allocator) void { l.stopping.store(true, .release); l.unpostFromRegistry(); if (l.watcher) |thread| { l.watch_stop.store(true, .release); _ = libc.write(l.control[1], "q", 1); thread.join(); for (l.control) |fd| _ = libc.close(fd); } for (0..quic_slots) |i| l.drop(@intCast(i)); if (comptime quic_enabled) { if (l.quic) |*listener| listener.deinit(); } if (comptime supported) { // A connection task may be waiting for the turn; it answers // nothing more once stopping, but the wait itself must end. pardes.turn.rest(); // An answer given as the editor quit (a Kill written to /ctl) // may still be on its way out: let it go before the cut. const deadline = nowMs() +| 200; while (nowMs() < deadline) { const pending = for (&l.runner.conns) |*conn| { if (!conn.live()) continue; conn.lock(); const n = conn.engine.output().len; conn.unlock(); if (n != 0) break true; } else false; if (!pending) break; Client.nap(1); } Client.nap(1); // the last bytes leave the stream's own buffer l.runner.stop(); pardes.turn.wake(); } pardes.turn.wake_parked = null; pardes.turn.answer_held = null; pardes.turn.stop(); if (l.path_len != 0) { var z: [sun_path_len:0]u8 = undefined; @memcpy(z[0..l.path_len], l.path_buf[0..l.path_len]); z[l.path_len] = 0; _ = libc.unlink(z[0..l.path_len :0]); } gpa.destroy(l); } }; pub fn socketPath(buf: *[sun_path_len]u8, dir: []const u8, name: []const u8) ?[:0]const u8 { if (name.len == 0) return null; if (std.mem.indexOfAny(u8, name, "/\x00") != null) return null; return std.fmt.bufPrintSentinel(buf, "{s}/" ++ prefix ++ "{s}.sock", .{ dir, name }, 0) catch null; } /// Serves `core`. From here on the calling thread has the editor's turn /// (`pardes.turn`) and must give it up while it waits. pub fn listen(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes, named: []const u8, fallback: []const u8, tcp_dial: ?[]const u8, quic_dial: ?[]const u8) ?*Listener { if (comptime !supported) return null; var dir_buf: [sun_path_len:0]u8 = undefined; const dir = socketDir(&dir_buf) orelse { log.warn("no runtime directory for the socket", .{}); return null; }; if (!ensureSocketDir(dir)) return null; const l = gpa.create(Listener) catch return null; l.* = .{ .io = io, .core = core }; pardes.turn.start(io); pardes.turn.wake_parked = Listener.wakeParked; pardes.turn.answer_held = Listener.answerHeld; pardes.turn.wake_ctx = l; l.runner.init(.{ .io = io, .root = pardes.ctlfs.root, .handler = .{ .ctx = l, .serve = Listener.onServe, .opened = Listener.onOpened, .refused = Listener.onRefused }, .greet_timeout_ms = Listener.greet_deadline_ms, }); const entry_name = if (named.len != 0) named else fallback; const p = socketPath(&l.path_buf, dir, entry_name) orelse { gpa.destroy(l); return null; }; _ = l.runner.listen(.{ .unix = p }, max_conns) catch |err| retry: { const existing = @import("fs.zig").statPath(io, p, .{ .follow_symlinks = false }) catch null; if (err != error.AddressInUse or existing == null or existing.?.kind != .unix_domain_socket or alive(p)) { log.warn("something is already listening on {s}", .{p}); l.deinit(gpa); return null; } if (libc.unlink(p) != 0) { l.deinit(gpa); return null; } break :retry l.runner.listen(.{ .unix = p }, max_conns) catch { l.deinit(gpa); return null; }; }; l.path_len = p.len; if (libc.chmod(p, 0o600) != 0) { l.deinit(gpa); return null; } l.postToRegistry(entry_name); if (tcp_dial) |dial| { if (!std.mem.startsWith(u8, dial, "tcp!")) { l.deinit(gpa); return null; } const address = networkAddress(dial, true) catch { log.warn("invalid TCP address {s}", .{dial}); l.deinit(gpa); return null; }; const bound = l.runner.listen(.{ .tcp = address }, max_conns) catch |err| { log.warn("cannot listen on {s}: {s}", .{ dial, @errorName(err) }); l.deinit(gpa); return null; }; l.tcp_address = canonicalIp(bound); log.info("serving 9P2000 over TCP on {f}", .{l.tcp_address.?}); } if (quic_dial) |dial| { if (comptime quic_enabled) { if (!std.mem.startsWith(u8, dial, "quic!")) { l.deinit(gpa); return null; } const address = networkAddress(dial, true) catch { log.warn("invalid QUIC address {s}", .{dial}); l.deinit(gpa); return null; }; l.quic = quic.Listener.init(address) catch |err| { log.warn("cannot listen on {s}: {s}", .{ dial, @errorName(err) }); l.deinit(gpa); return null; }; l.quic_address = l.quic.?.address; log.info("serving 9P2000 over QUIC on {f}", .{l.quic_address.?}); } else { log.warn("QUIC is unavailable in this build", .{}); l.deinit(gpa); return null; } } log.info("serving 9P2000 on {s}", .{p}); core.fs.socket_path = l.path(); core.fs.tcp_address = l.tcp_address; core.fs.quic_address = l.quic_address; return l; } const alive = transport.isListening; /// What is at a socket path, in cloud9's three answers (`post.Probe`). /// Asked here with libc rather than through `cloud9.post.probe`, whose /// raw Linux syscalls darwin cannot compile, but the classification is /// that file's and must not become a second opinion: the one definite /// refusal is `.stale`, nothing at the path at all is `.none`, and /// everything else — connected, busy, refused permission, a surprise — /// is `.live`. Uncertainty belongs to the server that owns the socket, /// never to a sweeper deciding what to delete. const Probe = enum { none, stale, live }; /// `SOCK.STREAM`, and the kernel's own non-blocking bit where there is /// one. Darwin's `SOCK.NONBLOCK` is a Zig shim for `std.posix.socket` /// to unpack, not an ABI value, so handing it to the raw libc call /// would ask a kernel that has never heard of it; there it is an /// `fcntl` instead. const probe_socket_kind: c_uint = libc.SOCK.STREAM | (if (darwin) 0 else libc.SOCK.NONBLOCK); fn probe(path: [:0]const u8) Probe { if (path.len + 1 > sun_path_len) return .live; // cannot ask; assume occupied 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, probe_socket_kind, 0); if (fd < 0) return .live; defer _ = libc.close(fd); // Non-blocking is the whole safety of this function, so it is read // back rather than assumed: a BLOCKING connect to a live server // whose backlog is full parks in the kernel with no timeout to end // it — measured, it simply never returns — and a sweep that parks // takes the editor's startup with it. An fd that cannot be proven // non-blocking is never connected at all, which lands on `.live`, // the answer that deletes nothing. if (comptime darwin) _ = setNonblock(fd); if (!isNonblocking(fd)) return .live; if (libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) == 0) return .live; // cloud9's `post.probe` mapping, answer for answer. The connection // is never wanted: EAGAIN (a full Unix backlog) and EINPROGRESS say // somebody is listening, which is all that was asked, and the fd // closes without ever waiting for it to become writable. return switch (libc.errno(-1)) { .AGAIN, .INPROGRESS, .PERM, .ACCES => .live, .CONNREFUSED => .stale, .NOENT, .NOTDIR => .none, else => .live, }; } /// Reaps the registry entries whose editor is gone, once per post. /// /// A clean exit unposts itself (`unpostFromRegistry`), so everything /// left behind comes from an exit that could not: an aborted test, a /// kill, a crash. No code in the dead process can ever run, so the cure /// has to be somebody else's readdir, and the next editor to start is /// the somebody. It is only ever a symlink that is followed or /// unlinked, because that is the only thing `postToRegistry` makes and /// anything else under a name belongs to whoever put it there; and only /// a definite refusal counts as gone. The socket a stale entry points /// at goes too, but not before a stat agrees it is a socket of ours: a /// plain file answers a connect with the same refusal. /// Is this the kind of path this editor is allowed to delete? A registry /// entry is a symlink we wrote, but its target is just bytes on disk that /// anyone could have pointed anywhere, so the sweep only ever follows one /// into the directory our own sockets live in, and only to a name of the /// shape we give them. Everything else gets its entry removed and its target /// left strictly alone. fn ourSocket(target: []const u8, sockets: []const u8) bool { if (sockets.len == 0 or !std.mem.startsWith(u8, target, sockets)) return false; if (target.len <= sockets.len or target[sockets.len] != '/') return false; const base = target[sockets.len + 1 ..]; if (std.mem.indexOfScalar(u8, base, '/') != null) return false; return std.mem.startsWith(u8, base, "pardes-9p-") and std.mem.endsWith(u8, base, ".sock"); } fn sweepRegistry(io: std.Io, svc: [:0]const u8, sockets: []const u8) void { if (comptime !supported) return; // The names are staged before anything is unlinked, so the sweep // never asks a directory to keep reading while it is being edited. var names: [4096]u8 = undefined; var staged: usize = 0; { const dir = std.Io.Dir.openDirAbsolute(io, svc, .{ .iterate = true }) catch return; defer dir.close(io); var read_buf: [std.Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined; var reader: std.Io.Dir.Reader = .init(dir, &read_buf); while (true) { const listed = (reader.next(io) catch break) orelse break; if (listed.name.len == 0 or listed.name.len > 255) continue; if (staged + 1 + listed.name.len > names.len) break; names[staged] = @intCast(listed.name.len); @memcpy(names[staged + 1 ..][0..listed.name.len], listed.name); staged += 1 + listed.name.len; } } var reaped: usize = 0; var i: usize = 0; while (i < staged) { const name = names[i + 1 ..][0..names[i]]; i += 1 + name.len; var entry_buf: [sun_path_len:0]u8 = undefined; const entry = std.fmt.bufPrintSentinel(&entry_buf, "{s}/{s}", .{ svc, name }, 0) catch continue; const facts = statNoFollow(entry) orelse continue; if (facts.mode & 0o170000 != 0o120000) continue; // not a symlink: not ours to judge const state = probe(entry); if (state == .live) continue; var link_buf: [sun_path_len]u8 = undefined; const n = libc.readlink(entry, &link_buf, link_buf.len); if (libc.unlink(entry) != 0) continue; reaped += 1; // A relative target would resolve against this editor's working // directory, which says nothing about what the entry named, and a // readlink that exactly filled the buffer was truncated, so the path // it produced is some other file's. if (state != .stale or n <= 0 or link_buf[0] != '/') continue; if (@as(usize, @intCast(n)) >= link_buf.len) continue; var target_buf: [sun_path_len:0]u8 = undefined; const target = std.fmt.bufPrintSentinel(&target_buf, "{s}", .{link_buf[0..@intCast(n)]}, 0) catch continue; if (!ourSocket(target, sockets)) continue; const t = statNoFollow(target) orelse continue; if (t.mode & 0o170000 != 0o140000 or t.uid != libc.getuid()) continue; // Ask again, immediately before deleting. The first probe was of the // ENTRY and is by now several syscalls old; a socket that is bound but // has not reached listen(2) yet answers ECONNREFUSED exactly like a // dead one, and that window is every server's startup. if (probe(target) != .stale) continue; _ = libc.unlink(target); } if (reaped != 0) log.info("reaped {d} stale registry entries under {s}", .{ reaped, svc }); } /// Whether it worked, because `probe` is not allowed to find out the /// hard way: a caller that needs the guarantee has to be able to check. fn setNonblock(fd: c_int) bool { const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0)); if (flags < 0) return false; var o: libc.O = @bitCast(@as(u32, @bitCast(flags))); o.NONBLOCK = true; return libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(o))))) >= 0; } fn isNonblocking(fd: c_int) bool { const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0)); if (flags < 0) return false; const o: libc.O = @bitCast(@as(u32, @bitCast(flags))); return o.NONBLOCK; } const testing = std.testing; test "the socket name is a third prefix in the shared directory" { var buf: [sun_path_len]u8 = undefined; const p = socketPath(&buf, "/run/user/1000", "t9srv").?; try testing.expectEqualStrings("/run/user/1000/pardes-9p-t9srv.sock", p); try testing.expect(!std.mem.startsWith(u8, std.fs.path.basename(p), "pardes-detached-")); } test "a name that is not one path component is no address at all" { var buf: [sun_path_len]u8 = undefined; try testing.expect(socketPath(&buf, "/run", "") == null); try testing.expect(socketPath(&buf, "/run", "a/b") == null); try testing.expect(socketPath(&buf, "/run", "a\x00b") == null); } test "the registry sweep takes the dead entries and leaves everything else" { if (comptime !supported) return error.SkipZigTest; // Under the real runtime directory, because a sun_path is 108 bytes // and the test cache's temporary directories are longer than that. var dir_buf: [sun_path_len:0]u8 = undefined; const base = socketDir(&dir_buf) orelse return error.SkipZigTest; if (!ensureSocketDir(base)) return error.SkipZigTest; var svc_buf: [sun_path_len:0]u8 = undefined; const svc = std.fmt.bufPrintSentinel(&svc_buf, "{s}/sweep-{d}", .{ base, @as(u32, @intCast(libc.getpid())) }, 0) catch return error.SkipZigTest; if (libc.mkdir(svc, 0o700) != 0) return error.SkipZigTest; // The sockets live where the real ones do -- beside the registry, not in // it, and named the way a session names them -- because the sweep only // follows an entry to a target of exactly that shape and place. var paths: [6][sun_path_len:0]u8 = undefined; const pid: u32 = @intCast(libc.getpid()); const live_sock = try std.fmt.bufPrintSentinel(&paths[0], "{s}/pardes-9p-sweeplive-{d}.sock", .{ base, pid }, 0); const dead_sock = try std.fmt.bufPrintSentinel(&paths[1], "{s}/pardes-9p-sweepdead-{d}.sock", .{ base, pid }, 0); const live = try std.fmt.bufPrintSentinel(&paths[2], "{s}/live", .{svc}, 0); const dead = try std.fmt.bufPrintSentinel(&paths[3], "{s}/dead", .{svc}, 0); const stranger = try std.fmt.bufPrintSentinel(&paths[4], "{s}/stranger", .{svc}, 0); // A dead entry pointing at something that is NOT one of our sockets: the // entry goes, the file it named must not. const outsider_sock = try std.fmt.bufPrintSentinel(&paths[5], "{s}/sweep-outsider-{d}.sock", .{ base, pid }, 0); var outsider_buf: [sun_path_len:0]u8 = undefined; const outsider = try std.fmt.bufPrintSentinel(&outsider_buf, "{s}/outsider", .{svc}, 0); defer { for ([_][:0]const u8{ live_sock, dead_sock, live, dead, stranger, outsider_sock, outsider }) |p| _ = libc.unlink(p); _ = libc.rmdir(svc); } const listening = bindSocket(live_sock); try testing.expect(listening >= 0); defer _ = libc.close(listening); try testing.expectEqual(@as(c_int, 0), libc.listen(listening, 1)); // Bound and then dropped: the file stays, and nobody answers it — // exactly what an aborted test leaves behind. const abandoned = bindSocket(dead_sock); try testing.expect(abandoned >= 0); _ = libc.close(abandoned); try testing.expectEqual(@as(c_int, 0), libc.symlink(live_sock, live)); try testing.expectEqual(@as(c_int, 0), libc.symlink(dead_sock, dead)); // Dead too, but its target is not one of our sockets by name, so the // entry must go and the file it named must survive untouched. const outside = bindSocket(outsider_sock); try testing.expect(outside >= 0); defer _ = libc.close(outside); try testing.expectEqual(@as(c_int, 0), libc.symlink(outsider_sock, outsider)); // Not a symlink, so not this program's to reason about, even though // connecting to it is refused exactly like the dead socket. try std.Io.Dir.cwd().writeFile(testing.io, .{ .sub_path = stranger, .data = "" }); // The guarantee itself, asserted rather than inferred from the fact // that the test finished: if a probe socket ever stops being // non-blocking the connect below parks instead of failing, and a // parked test costs whoever is building far more than a red one. const checking = libc.socket(libc.AF.UNIX, probe_socket_kind, 0); try testing.expect(checking >= 0); if (comptime darwin) try testing.expect(setNonblock(checking)); try testing.expect(isNonblocking(checking)); _ = libc.close(checking); // The listener's backlog is filled before anything is asked, because // a full backlog is the one place a blocking connect parks forever // (`unix_wait_for_peer`, no timeout) and a sweep that only answers // while nobody is queued is the sweep that hangs a build. Nothing // below accepts any of these, so the queue stays full throughout. var queued: [8]c_int = @splat(-1); defer for (queued) |fd| { if (fd >= 0) _ = libc.close(fd); }; for (&queued) |*slot| { const fd = libc.socket(libc.AF.UNIX, probe_socket_kind, 0); if (fd < 0) break; if (comptime darwin) _ = setNonblock(fd); var addr: libc.sockaddr.un = .{ .path = @splat(0) }; @memcpy(addr.path[0 .. live_sock.len + 1], live_sock[0 .. live_sock.len + 1]); _ = libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))); slot.* = fd; } const started = transport.nowMs(); try testing.expectEqual(Probe.live, probe(live)); try testing.expectEqual(Probe.stale, probe(dead)); sweepRegistry(testing.io, svc, base); // Promptly, and not "eventually": a blocking probe never comes back // at all, so any wall-clock bound at all is the assertion that // matters. A second is several thousand times what three connects // and a readdir cost. try testing.expect(transport.nowMs() - started < 1000); try testing.expect(statNoFollow(dead) == null); try testing.expect(statNoFollow(dead_sock) == null); try testing.expect(statNoFollow(live) != null); try testing.expect(statNoFollow(live_sock) != null); try testing.expect(statNoFollow(stranger) != null); try testing.expect(statNoFollow(outsider) == null); try testing.expect(statNoFollow(outsider_sock) != null); try testing.expectEqual(Probe.live, probe(live)); } fn bindSocket(path: [:0]const u8) c_int { const fd = libc.socket(libc.AF.UNIX, libc.SOCK.STREAM, 0); if (fd < 0) return fd; var addr: libc.sockaddr.un = .{ .path = @splat(0) }; @memcpy(addr.path[0 .. path.len + 1], path[0 .. path.len + 1]); if (libc.bind(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) != 0) { _ = libc.close(fd); return -1; } return fd; } test "one connection's buffers are sized from the one msize constant" { try testing.expect(msize >= ninep.min_msize); try testing.expectEqual(@as(u32, msize), Runner.msize); try testing.expectEqual(@as(usize, msize), @typeInfo(@FieldType(Runner.Conn, "in")).array.len); try testing.expectEqual(@as(usize, 2 * msize), @typeInfo(@FieldType(Runner.Conn, "out")).array.len); try testing.expectEqual(@as(usize, max_conns), Runner.max_connections); if (quic_enabled) { const c: Conn = .{}; try testing.expectEqual(@as(usize, msize), c.in.len); } } test "TCP addresses are numeric and normalize mapped IPv4" { const loopback = try networkAddress("tcp!127.0.0.1!5640", false); try testing.expectEqualDeep(loopback, try networkAddress("tcp!::ffff:127.0.0.1!5640", false)); try testing.expectEqualDeep(loopback, try networkAddress("tcp!::ffff:7f00:1!5640", false)); try testing.expectEqualDeep(try networkAddress("tcp!::1!5640", false), try networkAddress("tcp!0:0:0:0:0:0:0:1!5640", false)); try testing.expectEqual(@as(u16, 0), (try networkAddress("tcp!127.0.0.1!0", true)).getPort()); for ([_][]const u8{ "tcp!localhost!5640", "tcp!127.0.0.1!0", "tcp!127.0.0.1!-1", "tcp!127.0.0.1!65536", "tcp!127.0.0.1!", "tcp!!5640" }) |dial| try testing.expectError(error.BadDial, networkAddress(dial, false)); try Client.validateDial("tcp!127.0.0.1!5640"); try Client.validateDial("/tmp/pardes-owned.sock"); try Client.validateDial("unix!/tmp/pardes-owned.sock"); try testing.expectError(error.BadDial, Client.validateDial("unix!work")); try testing.expectError(error.BadDial, Client.validateDial("unix!")); try testing.expectError(error.BadDial, Client.validateDial("unix!/tmp/a\x00b")); try testing.expectError(error.BadDial, Client.validateDial("tcp!localhost!5640")); try testing.expectError(error.BadDial, Client.validateDial("/tmp/a\x00b")); if (quic_enabled) { try Client.validateDial("quic!127.0.0.1!5640"); } else try testing.expectError(error.QuicUnavailable, Client.validateDial("quic!127.0.0.1!5640")); } extern "c" fn mkdtemp(template: [*:0]u8) ?[*:0]u8; extern "c" fn rmdir(path: [*:0]const u8) c_int; test "a listening editor posts itself into the 9P registry and unposts on stop" { if (comptime !supported) return error.SkipZigTest; const gpa = testing.allocator; var directory: [64:0]u8 = undefined; _ = try std.fmt.bufPrintSentinel(&directory, "/tmp/pardes-post-XXXXXX", .{}, 0); if (mkdtemp(&directory) == null) return error.TempDirectoryFailed; const dir = std.mem.span(@as([*:0]const u8, &directory)); const old_runtime = if (libc.getenv("XDG_RUNTIME_DIR")) |v| try gpa.dupeZ(u8, std.mem.span(v)) else null; defer { if (old_runtime) |v| { _ = setenv("XDG_RUNTIME_DIR", v, 1); gpa.free(v); } else _ = unsetenv("XDG_RUNTIME_DIR"); } try testing.expectEqual(@as(c_int, 0), setenv("XDG_RUNTIME_DIR", &directory, 1)); var entry_buf: [sun_path_len:0]u8 = undefined; const entry = try std.fmt.bufPrintSentinel(&entry_buf, "{s}/9p/pardes/unit", .{dir}, 0); var svc_buf: [sun_path_len:0]u8 = undefined; const svc = try std.fmt.bufPrintSentinel(&svc_buf, "{s}/9p/pardes", .{dir}, 0); var sock_buf: [sun_path_len:0]u8 = undefined; const sock = try std.fmt.bufPrintSentinel(&sock_buf, "{s}/" ++ prefix ++ "unit.sock", .{dir}, 0); defer { _ = libc.unlink(entry); _ = rmdir(svc); var reg_buf: [sun_path_len:0]u8 = undefined; if (std.fmt.bufPrintSentinel(®_buf, "{s}/9p", .{dir}, 0)) |reg| { _ = rmdir(reg); } else |_| {} _ = libc.unlink(sock); _ = rmdir(&directory); } var link: [sun_path_len]u8 = undefined; { const p = try pardes.Pardes.init(gpa, .{ .tty_only = true, .cols = 40, .rows = 12 }); defer p.deinit(); const l = listen(testing.io, gpa, p, "unit", "", null, null) orelse return error.ListenFailed; defer l.deinit(gpa); // The socket stays exactly where pardes has always bound it: // adopting the registry moves nothing, it only advertises. try testing.expect(statNoFollow(sock) != null); // And the registry holds a symlink to it one directory down, so // several editors group under /mnt/9p/pardes/ instead of // crowding the registry root — the layout zmx posts its // sessions in. const n = libc.readlink(entry, &link, link.len); try testing.expect(n > 0); try testing.expectEqualStrings(sock, link[0..@intCast(n)]); } // Stopping takes the name back out, so the next editor of that name // is not refused by its own corpse. try testing.expect(libc.readlink(entry, &link, link.len) < 0); } test "Unix TCP and QUIC share one listener through reads writes reconnects and reset" { if (comptime !supported) return error.SkipZigTest; const gpa = testing.allocator; var directory: [64:0]u8 = undefined; _ = try std.fmt.bufPrintSentinel(&directory, "/tmp/pardes-tcp-XXXXXX", .{}, 0); if (mkdtemp(&directory) == null) return error.TempDirectoryFailed; // What the listener posted and left inside goes with it. defer std.Io.Dir.cwd().deleteTree(testing.io, std.mem.sliceTo(&directory, 0)) catch |err| std.debug.print("deleteTree: {s}\n", .{@errorName(err)}); const old_runtime = if (libc.getenv("XDG_RUNTIME_DIR")) |v| try gpa.dupeZ(u8, std.mem.span(v)) else null; defer { if (old_runtime) |v| { _ = setenv("XDG_RUNTIME_DIR", v, 1); gpa.free(v); } else _ = unsetenv("XDG_RUNTIME_DIR"); } try testing.expectEqual(@as(c_int, 0), setenv("XDG_RUNTIME_DIR", &directory, 1)); const replacement = try gpa.alloc(u8, 3 * msize + 27); defer gpa.free(replacement); @memset(replacement, 'x'); @memcpy(replacement[0.."changed café λ\n".len], "changed café λ\n"); replacement[replacement.len - 1] = '\n'; const Worker = struct { dial: []const u8, body_path: []const u8, expected: []const u8, replacement: []const u8, done: std.atomic.Value(bool) = .init(false), failure: ?anyerror = null, fn run(w: *@This()) void { defer w.done.store(true, .release); w.check() catch |err| { w.failure = err; }; } fn check(w: *@This()) !void { const before = try Client.readLimit(testing.allocator, w.dial, w.body_path, w.body_path, w.expected.len); defer testing.allocator.free(before); try testing.expectEqualStrings(w.expected, before); try testing.expectError(error.FileTooLarge, Client.readLimit(testing.allocator, w.dial, w.body_path, w.body_path, w.expected.len - 1)); const listing = try Client.readLimit(testing.allocator, w.dial, "/pane", "/pane", 128); defer testing.allocator.free(listing); const exact_listing = try Client.readLimit(testing.allocator, w.dial, "/pane", "/pane", listing.len); defer testing.allocator.free(exact_listing); try testing.expectEqualStrings(listing, exact_listing); try testing.expectError(error.FileTooLarge, Client.readLimit(testing.allocator, w.dial, "/pane", "/pane", listing.len - 1)); try Client.write(testing.allocator, w.dial, w.body_path, w.replacement); const after = try Client.read(testing.allocator, w.dial, w.body_path, w.body_path); defer testing.allocator.free(after); try testing.expectEqualStrings(w.replacement, after); const screen = try Client.read(testing.allocator, w.dial, "/screen", "/screen"); defer testing.allocator.free(screen); const parsed = try std.json.parseFromSlice(struct { cols: u16, rows: u16 }, testing.allocator, screen, .{ .ignore_unknown_fields = true }); defer parsed.deinit(); try testing.expectEqual(@as(u16, 40), parsed.value.cols); try testing.expectEqual(@as(u16, 12), parsed.value.rows); } }; const protocols: []const []const u8 = if (quic_enabled) &.{ "tcp", "quic" } else &.{"tcp"}; for (protocols) |protocol| for ([_][]const u8{ "127.0.0.1", "::1" }) |host| { const p = try pardes.Pardes.init(gpa, .{ .tty_only = true, .cols = 40, .rows = 12 }); defer p.deinit(); const pane = try p.setTestFile("initial\n"); while (p.nextEffect()) |_| {} var bind_buf: [64]u8 = undefined; const bind = try std.fmt.bufPrint(&bind_buf, "{s}!{s}!0", .{ protocol, host }); const tcp = std.mem.eql(u8, protocol, "tcp"); const l = listen(testing.io, gpa, p, "roundtrip", "", if (tcp) bind else null, if (tcp) null else bind) orelse return error.ListenFailed; defer { l.reset(p); l.deinit(gpa); } try testing.expect(l.path().len != 0); try testing.expectEqual(tcp, l.tcp_address != null); try testing.expectEqual(@as(usize, max_conns), l.runner.conns.len); try testing.expect(l.watcher == null); const port = (if (tcp) l.tcp_address else l.quic_address).?.getPort(); try testing.expect(port != 0); var dial_buf: [64]u8 = undefined; const network_dial = try std.fmt.bufPrint(&dial_buf, "{s}!{s}!{d}", .{ protocol, host, port }); var body_buf: [64]u8 = undefined; const body = try std.fmt.bufPrint(&body_buf, "/pane/{d}/body", .{pane.serial}); for ([_][]const u8{ network_dial, l.path(), network_dial }, 0..) |dial, attempt| { var worker: Worker = .{ .dial = dial, .body_path = body, .expected = if (attempt == 0) "initial\n" else replacement, .replacement = replacement }; const thread = try std.Thread.spawn(.{}, Worker.run, .{&worker}); defer thread.join(); // The editor rests while the client works: its requests are // answered on the runner's tasks, not by this thread. const deadline = Client.nowMs() + 3 * Client.budget_ms; pardes.turn.rest(); while (!worker.done.load(.acquire) and Client.nowMs() < deadline) { if (quic_enabled) { pardes.turn.wake(); _ = l.tick(); pardes.turn.rest(); } Client.nap(1); } // The client has hung up, but its connection's task may not have // seen that yet. The releases it then owes are refused once // `reset` starts cutting (they name the old editor's opens), and // this test's "replacement" is the same editor, so they would // stay in `p.fs.opens`. Let the hangup pay them first, resting // so the task can take the turn: that is what is under test here. const settle_by = Client.nowMs() + 3 * Client.budget_ms; while (l.runner.count() != 0 and Client.nowMs() < settle_by) Client.nap(1); pardes.turn.wake(); try testing.expect(worker.done.load(.acquire)); if (worker.failure) |err| return err; try testing.expectEqual(@as(usize, 0), l.runner.count()); l.reset(p); try testing.expectEqual(@as(usize, 0), l.runner.count()); for (&l.conns, 0..) |conn, i| try testing.expect(!l.live(@intCast(i)) and !conn.draining); for (p.fs.opens) |o| try testing.expect(o.node == 0); } }; } test "a change waits while the editor is out mid-step, a read does not, and the change lands when the editor rests" { if (comptime !supported) return error.SkipZigTest; const gpa = testing.allocator; var directory: [64:0]u8 = undefined; _ = try std.fmt.bufPrintSentinel(&directory, "/tmp/pardes-park-XXXXXX", .{}, 0); if (mkdtemp(&directory) == null) return error.TempDirectoryFailed; // What the listener posted and left inside goes with it. defer std.Io.Dir.cwd().deleteTree(testing.io, std.mem.sliceTo(&directory, 0)) catch |err| std.debug.print("deleteTree: {s}\n", .{@errorName(err)}); const old_runtime = if (libc.getenv("XDG_RUNTIME_DIR")) |v| try gpa.dupeZ(u8, std.mem.span(v)) else null; defer { if (old_runtime) |v| { _ = setenv("XDG_RUNTIME_DIR", v, 1); gpa.free(v); } else _ = unsetenv("XDG_RUNTIME_DIR"); } try testing.expectEqual(@as(c_int, 0), setenv("XDG_RUNTIME_DIR", &directory, 1)); const p = try pardes.Pardes.init(gpa, .{ .tty_only = true, .cols = 40, .rows = 12 }); defer p.deinit(); const pane = try p.setTestFile("before\n"); while (p.nextEffect()) |_| {} // This thread is the editor's: it has the turn from here on. const l = listen(testing.io, gpa, p, "park", "", null, null) orelse return error.ListenFailed; defer { l.reset(p); l.deinit(gpa); } var body_buf: [64]u8 = undefined; const body = try std.fmt.bufPrint(&body_buf, "/pane/{d}/body", .{pane.serial}); const Writer = struct { dial: []const u8, path: []const u8, done: std.atomic.Value(bool) = .init(false), failure: ?anyerror = null, fn run(w: *@This()) void { defer w.done.store(true, .release); Client.write(testing.allocator, w.dial, w.path, "after\n") catch |err| { w.failure = err; }; } }; var writer: Writer = .{ .dial = l.path(), .path = body }; // The editor goes out mid-step -- a Look resolving a path, say -- and a // write arrives meanwhile. It must not land: the step still holds // pointers into the pane. A read is answered all the same. pardes.turn.yield(); const thread = try std.Thread.spawn(.{}, Writer.run, .{&writer}); defer thread.join(); Client.nap(150); try testing.expect(!writer.done.load(.acquire)); const seen = try Client.readLimit(gpa, l.path(), body, body, 64); defer gpa.free(seen); try testing.expectEqualStrings("before\n", seen); try testing.expect(!writer.done.load(.acquire)); pardes.turn.back(); try testing.expectEqualStrings("before\n", pane.file.?.content); // Resting is what lets it in: the parked write is retried, served, and // the client hears back. pardes.turn.rest(); const deadline = Client.nowMs() + Client.budget_ms; while (!writer.done.load(.acquire) and Client.nowMs() < deadline) Client.nap(1); pardes.turn.wake(); try testing.expect(writer.done.load(.acquire)); if (writer.failure) |err| return err; try testing.expectEqualStrings("after\n", pane.file.?.content); } test "a held read is answered when the log has news, and a flushed one spends nothing" { if (comptime !supported) return error.SkipZigTest; const gpa = testing.allocator; var directory: [64:0]u8 = undefined; _ = try std.fmt.bufPrintSentinel(&directory, "/tmp/pardes-held-XXXXXX", .{}, 0); if (mkdtemp(&directory) == null) return error.TempDirectoryFailed; // What the listener posted and left inside goes with it. defer std.Io.Dir.cwd().deleteTree(testing.io, std.mem.sliceTo(&directory, 0)) catch |err| std.debug.print("deleteTree: {s}\n", .{@errorName(err)}); const old_runtime = if (libc.getenv("XDG_RUNTIME_DIR")) |v| try gpa.dupeZ(u8, std.mem.span(v)) else null; defer { if (old_runtime) |v| { _ = setenv("XDG_RUNTIME_DIR", v, 1); gpa.free(v); } else _ = unsetenv("XDG_RUNTIME_DIR"); } try testing.expectEqual(@as(c_int, 0), setenv("XDG_RUNTIME_DIR", &directory, 1)); const p = try pardes.Pardes.init(gpa, .{ .tty_only = true, .cols = 40, .rows = 12 }); defer p.deinit(); _ = try p.setTestFile("held\n"); p.update(.tick); // the pane is logged, before the log is opened while (p.nextEffect()) |_| {} const l = listen(testing.io, gpa, p, "held", "", null, null) orelse return error.ListenFailed; defer { l.reset(p); l.deinit(gpa); } // This thread is the editor's and plays the client too, resting while // it does so that the runner's tasks can take the turn and answer. const s = try gpa.create(Client.Session); defer gpa.destroy(s); s.* = .{ .fd = -1, .deadline = Client.nowMs() + 3 * Client.budget_ms, .display_path = "/log" }; var sock_buf: [sun_path_len]u8 = undefined; pardes.turn.rest(); defer pardes.turn.wake(); s.fd = try Client.connect(try Client.resolve(&sock_buf, l.path()), s.deadline); defer _ = libc.close(s.fd); s.cl = .init(.{ .in = &s.in, .out = &s.out }); var remote: Client.RemoteError = .{}; _ = try s.ask(.{ .version = .{} }, &remote); _ = try s.ask(.{ .attach = .{ .fid = 0, .uname = "held" } }, &remote); _ = try s.ask(.{ .walk = .{ .fid = 0, .newfid = 1, .names = &.{"log"} } }, &remote); _ = try s.ask(.{ .open = .{ .fid = 1, .mode = ninep.ordwr } }, &remote); var frozen: u64 = 0; while (true) { const n = (try s.ask(.{ .read = .{ .fid = 1, .offset = frozen, .count = 4096 } }, &remote)).read.len; if (n == 0) break; frozen += n; } _ = try s.ask(.{ .write = .{ .fid = 1, .offset = 0, .data = "follow\n" } }, &remote); const heldNow = struct { /// Waits until the core holds `n` reads. fn count(core: *pardes.Pardes, n: usize) !void { const deadline = Client.nowMs() + Client.budget_ms; while (Client.nowMs() < deadline) { pardes.turn.wake(); var got: usize = 0; for (core.fs.opens) |o| got += @intFromBool(o.held != null); pardes.turn.rest(); if (got == n) return; Client.nap(1); } return error.Timeout; } }; const serial = p.panes[0].?.serial; var event_path: [32]u8 = undefined; var event_names: [3][]const u8 = .{ "pane", try std.fmt.bufPrint(&event_path, "{d}", .{serial}), "event" }; _ = try s.ask(.{ .walk = .{ .fid = 0, .newfid = 2, .names = &event_names } }, &remote); _ = try s.ask(.{ .open = .{ .fid = 2, .mode = ninep.oread } }, &remote); // Flushed while held, a read is gone from the engine without the core // hearing of it. Its tag, asked again at once for a read of the pane's // event, must not be answered with the log's next record, and that // record must reach the log's next read. const flushed = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); try s.flush(); try heldNow.count(p, 1); _ = try s.cl.submit(.{ .flush = .{ .oldtag = flushed } }); const interrupted = try s.settle(); try testing.expectEqual(flushed, interrupted.tag); try testing.expect(interrupted.result == .fail); try testing.expect((try s.settle()).result == .flush); const reused = try s.cl.submit(.{ .read = .{ .fid = 2, .offset = 0, .count = 4096 } }); try testing.expectEqual(flushed, reused); try s.flush(); try heldNow.count(p, 2); pardes.turn.wake(); p.setMessage(0, "after the flush"); pardes.turn.rest(); try heldNow.count(p, 1); pardes.turn.wake(); p.fs.origin = 'K'; _ = pardes.ctlfs.events.noteAction(p, 0, .body_exec, 0, 4, 0, "held"); pardes.turn.rest(); const clicked = try s.settle(); try testing.expectEqual(reused, clicked.tag); try testing.expectEqualStrings("KX0 4 0 4 held\n", clicked.result.read); try testing.expect(std.mem.indexOf(u8, (try s.ask(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }, &remote)).read, "after the flush") != null); // Held, a read is answered by the record that arrives, on its own // connection, with no retry of every parked request asked for; a // second read on the same open meanwhile is refused, not orphaned. const waiting = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); try s.flush(); try heldNow.count(p, 1); const second = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); const refused = try s.settle(); try testing.expectEqual(second, refused.tag); try testing.expectEqualStrings(pardes.ctlfs.e_in_use, refused.result.fail); pardes.turn.wake(); p.setMessage(0, "while held"); const retry_all = pardes.turn.parked; pardes.turn.rest(); try testing.expect(!retry_all); const answered = try s.settle(); try testing.expectEqual(waiting, answered.tag); try testing.expect(std.mem.indexOf(u8, answered.result.read, "while held") != null); // Clunked, the log's held read is interrupted by the engine and the // record it would have had is spent on nobody. _ = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); try s.flush(); try heldNow.count(p, 1); s.drop(1); s.drop(2); try heldNow.count(p, 0); } 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; pub fn start(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes) ?*Listener { var name: [16]u8 = undefined; // unreachable: a u32 is at most 10 digits const fallback = std.fmt.bufPrint(&name, "{d}", .{@as(u32, @intCast(libc.getpid()))}) catch unreachable; return listen(io, gpa, core, core.opts.ninep_name, fallback, core.opts.ninep_tcp, core.opts.ninep_quic) orelse { core.reportError(0, "9p listener", error.ListenFailed); return null; }; } /// What a pane shell is told about the editor above it, set into this /// process's environment just before `forkpty` so the child inherits it. /// /// Two independent facts, and they are separate variables because they answer /// separate questions. `PARDES_PID` says "you are inside this editor, and it /// will take a Look from you": a session whose 9P listener never came up still /// owns its children, so the child says so rather than looking, from its own /// side, exactly like no pardes at all. `PARDES_9P` and `PARDES_PANE` say how /// to reach it, and a `--nested` session exports them too — its socket stays /// open to scripts and to mounted shells — while withholding `PARDES_PID`, /// which is the whole of what `--nested` means. pub fn exportPaneEnv(listener: ?*const Listener, serial: u32, adopts: bool) void { var announced = false; if (adopts) announcing: { var buf: [16]u8 = undefined; const text = std.fmt.bufPrintSentinel(&buf, "{d}", .{@as(u32, @intCast(libc.getpid()))}, 0) catch break :announcing; announced = setenv("PARDES_PID", text, 1) == 0; } if (!announced) _ = unsetenv("PARDES_PID"); if (listener) |l| exporting: { var sock: [sun_path_len]u8 = undefined; const path = std.fmt.bufPrintSentinel(&sock, "{s}", .{l.path()}, 0) catch break :exporting; var buf: [16]u8 = undefined; const id = std.fmt.bufPrintSentinel(&buf, "{d}", .{serial}, 0) catch break :exporting; if (setenv("PARDES_9P", path, 1) != 0) break :exporting; if (setenv("PARDES_PANE", id, 1) == 0) return; } _ = unsetenv("PARDES_9P"); _ = unsetenv("PARDES_PANE"); } test "9P shell environment states being inside pardes apart from how to reach it" { const names = [_][*:0]const u8{ "XDG_RUNTIME_DIR", "HOME", "PARDES_PID", "PARDES_9P", "PARDES_PANE" }; var saved: [names.len]?[:0]u8 = @splat(null); for (names, &saved) |name, *value| { if (libc.getenv(name)) |old| value.* = try testing.allocator.dupeZ(u8, std.mem.span(old)); } defer for (names, saved) |name, value| { if (value) |old| { _ = setenv(name, old, 1); testing.allocator.free(old); } else _ = unsetenv(name); }; var dir: [sun_path_len:0]u8 = undefined; _ = setenv("XDG_RUNTIME_DIR", "/run/user/1000", 1); try testing.expectEqualStrings("/run/user/1000", socketDir(&dir).?); _ = unsetenv("XDG_RUNTIME_DIR"); _ = setenv("HOME", "/home/example", 1); try testing.expectEqualStrings("/home/example/.local/state/pardes", socketDir(&dir).?); _ = unsetenv("HOME"); try testing.expect(socketDir(&dir) == null); const listener = try testing.allocator.create(Listener); defer testing.allocator.destroy(listener); listener.* = .{ .io = testing.io, .core = undefined }; // only its path is read const path = "/tmp/pardes-example.sock"; @memcpy(listener.path_buf[0..path.len], path); listener.path_len = path.len; var own: [16]u8 = undefined; const own_pid = try std.fmt.bufPrint(&own, "{d}", .{@as(u32, @intCast(libc.getpid()))}); // A `--nested` session: reachable for scripts and mounts, but nobody's // parent, so a pardes started in one of its shells runs a session of // its own instead of handing its argument over. exportPaneEnv(listener, 7, false); try testing.expectEqualStrings(path, std.mem.span(libc.getenv("PARDES_9P").?)); try testing.expectEqualStrings("7", std.mem.span(libc.getenv("PARDES_PANE").?)); try testing.expect(libc.getenv("PARDES_PID") == null); exportPaneEnv(listener, 8, true); try testing.expectEqualStrings("8", std.mem.span(libc.getenv("PARDES_PANE").?)); try testing.expectEqualStrings(own_pid, std.mem.span(libc.getenv("PARDES_PID").?)); // A session whose listener never came up is still the session this shell // is inside: the pid stands on its own, and only the address is missing. exportPaneEnv(null, 0, true); try testing.expect(libc.getenv("PARDES_9P") == null); try testing.expect(libc.getenv("PARDES_PANE") == null); try testing.expectEqualStrings(own_pid, std.mem.span(libc.getenv("PARDES_PID").?)); exportPaneEnv(null, 0, false); try testing.expect(libc.getenv("PARDES_9P") == null); try testing.expect(libc.getenv("PARDES_PANE") == null); try testing.expect(libc.getenv("PARDES_PID") == null); } pub const Client = struct { pub const budget_ms: i64 = 2000; pub const max_depth: usize = 2 * ninep.max_welem; const uname = "pardes"; const Dial = union(enum) { unix: [:0]const u8, tcp: std.Io.net.IpAddress, quic: std.Io.net.IpAddress }; pub const Error = error{ PathTooDeep, BadDial, Dial, Hangup, Timeout, Botch, Remote, IsDirectory, NotFound, FileTooLarge, QuicUnavailable, }; pub fn read(gpa: std.mem.Allocator, dial: []const u8, path: []const u8, display_path: []const u8) ![]u8 { return readLimit(gpa, dial, path, display_path, limits.max_file_bytes); } pub fn readLimit(gpa: std.mem.Allocator, dial: []const u8, path: []const u8, display_path: []const u8, max_bytes: usize) ![]u8 { if (comptime !supported) return error.Unsupported; var names: [max_depth][]const u8 = undefined; const n = try elements(path, &names); var sock_buf: [sun_path_len]u8 = undefined; const sock = try resolve(&sock_buf, dial); var remote: RemoteError = .{}; return fetchBytes(gpa, sock, names[0..n], &remote, null, display_path, @min(max_bytes, limits.max_file_bytes)); } /// Whether a peer answers at `dial`: a version and attach, and its root /// read. Any answer from the peer, an error of its own included, is one. pub fn probe(gpa: std.mem.Allocator, dial: []const u8) !void { const got = readLimit(gpa, dial, "", "", 0) catch |err| switch (err) { error.IsDirectory, error.NotFound, error.Remote, error.FileTooLarge => return, else => return err, }; gpa.free(got); } pub fn write(gpa: std.mem.Allocator, dial: []const u8, path: []const u8, bytes: []const u8) !void { if (comptime !supported) return error.Unsupported; var names: [max_depth][]const u8 = undefined; const n = try elements(path, &names); var sock_buf: [sun_path_len]u8 = undefined; const sock = try resolve(&sock_buf, dial); var remote: RemoteError = .{}; const result = try fetchBytes(gpa, sock, names[0..n], &remote, bytes, path, limits.max_file_bytes); gpa.free(result); } const RemoteError = struct { buf: [ninep.errmax]u8 = undefined, len: usize = 0, fn set(r: *RemoteError, msg: []const u8) error{Remote} { r.len = @min(msg.len, r.buf.len); @memcpy(r.buf[0..r.len], msg[0..r.len]); return error.Remote; } }; fn elements(path: []const u8, out: *[max_depth][]const u8) Error!usize { var n: usize = 0; var it = std.mem.tokenizeScalar(u8, path, '/'); while (it.next()) |name| { if (n == out.len) return Error.PathTooDeep; out[n] = name; n += 1; } return n; } fn resolve(buf: *[sun_path_len]u8, dial: []const u8) error{ BadDial, QuicUnavailable }!Dial { if (dial.len == 0) return error.BadDial; if (std.mem.startsWith(u8, dial, "tcp!")) return .{ .tcp = try networkAddress(dial, false) }; if (std.mem.startsWith(u8, dial, "quic!")) { if (comptime !quic_enabled) return error.QuicUnavailable; return .{ .quic = try networkAddress(dial, false) }; } const explicit_unix = std.mem.startsWith(u8, dial, "unix!"); const path = if (explicit_unix) dial[5..] else dial; if (explicit_unix and !std.mem.startsWith(u8, path, "/")) return error.BadDial; if (std.mem.indexOfScalar(u8, path, '/') != null) { if (std.mem.indexOfScalar(u8, path, 0) != null) return error.BadDial; return .{ .unix = std.fmt.bufPrintSentinel(buf, "{s}", .{path}, 0) catch return error.BadDial }; } var dir_buf: [sun_path_len:0]u8 = undefined; const dir = socketDir(&dir_buf) orelse return error.BadDial; return .{ .unix = socketPath(buf, dir, dial) orelse return error.BadDial }; } pub fn validateDial(dial: []const u8) error{ BadDial, QuicUnavailable }!void { var buf: [sun_path_len]u8 = undefined; _ = try resolve(&buf, dial); } const Session = struct { fd: c_int, quic: if (quic_enabled) ?quic.Connection else void = if (quic_enabled) null else {}, deadline: i64, display_path: []const u8, cl: ninep.Client = undefined, in: [msize]u8 = undefined, out: [msize]u8 = undefined, stage: [msize]u8 = undefined, fn wait(s: *Session, events: i16) Error!void { while (true) { const left = s.deadline - nowMs(); if (left <= 0) return Error.Timeout; if (comptime quic_enabled) { if (s.quic) |*connection| { var fds = [1]libc.pollfd{connection.poll().?}; const timeout = @min(left, connection.nextDue() orelse budget_ms); const ready = libc.poll(&fds, 1, @intCast(timeout)); if (ready < 0) { if (libc.errno(ready) == .INTR) continue; return Error.Hangup; } if (nowMs() >= s.deadline) return Error.Timeout; if (fds[0].revents & @as(i16, @intCast(libc.POLL.NVAL)) != 0) return Error.Hangup; connection.events() catch return Error.Hangup; return; } } return transport.wait(s.fd, events, s.deadline) catch |err| switch (err) { error.Timeout => Error.Timeout, else => Error.Hangup, }; } } fn flush(s: *Session) Error!void { while (s.cl.output().len != 0) { if (nowMs() >= s.deadline) return Error.Timeout; const bytes = s.cl.output(); if (comptime quic_enabled) { if (s.quic) |*connection| { const sent = connection.write(bytes) catch return Error.Hangup; if (sent == 0) { try s.wait(poll_out); continue; } s.cl.wrote(sent); continue; } } try s.wait(poll_out); const sent = (transport.write(s.fd, bytes) catch return Error.Hangup) orelse continue; s.cl.wrote(@intCast(sent)); } } fn settle(s: *Session) Error!ninep.Client.Done { while (true) { if (nowMs() >= s.deadline) return Error.Timeout; try s.flush(); if (s.cl.take()) |done| return done; if (s.cl.dead) return Error.Botch; const room = s.cl.in.len - s.cl.in_len; if (room == 0) return Error.Botch; if (comptime quic_enabled) { if (s.quic) |*connection| { const got = (connection.read(s.stage[0..@min(room, s.stage.len)]) catch return Error.Hangup) orelse { try s.wait(poll_in); continue; }; if (got == 0) return Error.Hangup; const n = s.cl.push(s.stage[0..got]); std.debug.assert(n == got); continue; } } try s.wait(poll_in); const got = (transport.read(s.fd, s.stage[0..@min(room, s.stage.len)]) catch return Error.Hangup) orelse continue; if (got == 0) return Error.Hangup; const n = s.cl.push(s.stage[0..@intCast(got)]); std.debug.assert(n == @as(usize, @intCast(got))); } } fn ask(s: *Session, req: ninep.Client.Request, remote: *RemoteError) Error!ninep.Client.Result { _ = s.cl.submit(req) catch return Error.Botch; const done = try s.settle(); if (done.result == .fail) return remote.set(done.result.fail); if (std.mem.eql(u8, @tagName(done.result), @tagName(std.meta.activeTag(req)))) return done.result; return Error.Botch; } fn drop(s: *Session, fid: u32) void { _ = s.cl.submit(.{ .clunk = .{ .fid = fid } }) catch return; _ = s.settle() catch {}; } fn dropNoWait(s: *Session, fid: u32) void { _ = s.cl.submit(.{ .clunk = .{ .fid = fid } }) catch return; s.flush() catch {}; } }; /// Version, attach and a walk to `names`: the fid there and its qid. fn walkTo(s: *Session, names: []const []const u8, remote: *RemoteError) !struct { fid: u32, qid: ninep.Qid } { _ = try s.ask(.{ .version = .{} }, remote); if (s.cl.msize == 0) return Error.Botch; const root: u32 = 0; var here = (try s.ask(.{ .attach = .{ .fid = root, .uname = uname } }, remote)).attach; var cur: u32 = root; var next: u32 = 1; var i: usize = 0; while (i < names.len) { const n = @min(ninep.max_welem, names.len - i); const w = (try s.ask(.{ .walk = .{ .fid = cur, .newfid = next, .names = names[i..][0..n], } }, remote)).walk; if (w.nwqid != n) return Error.NotFound; here = w.wqid[n - 1]; if (cur != root) s.drop(cur); cur = next; next = if (next == 1) 2 else 1; i += n; } return .{ .fid = cur, .qid = here }; } fn transact( s: *Session, names: []const []const u8, out: *std.Io.Writer.Allocating, remote: *RemoteError, write_bytes: ?[]const u8, read_limit: usize, ) !void { const at = try walkTo(s, names, remote); const cur = at.fid; const here = at.qid; defer s.dropNoWait(cur); const directory = here.type & ninep.qtdir != 0; const mode: u8 = if (write_bytes != null) ninep.owrite else ninep.oread; const truncate = write_bytes != null and names.len > 0 and (std.mem.eql(u8, names[0], "os") or std.mem.eql(u8, names[names.len - 1], "body")); _ = try s.ask(.{ .open = .{ .fid = cur, .mode = mode | if (truncate) ninep.otrunc else 0 } }, remote); if (write_bytes) |bytes| { if (directory) return Error.IsDirectory; var written: usize = 0; while (written < bytes.len) { const chunk = bytes[written..][0..@min(bytes.len - written, s.cl.maxWrite())]; const count = (try s.ask(.{ .write = .{ .fid = cur, .offset = written, .data = chunk } }, remote)).write; if (count == 0 or count > chunk.len) return Error.Botch; written += count; } return; } const max_bytes: u64 = @min(read_limit, @as(usize, if (directory) limits.max_stream_bytes else limits.max_file_bytes)); const wire_limit: u64 = if (directory) limits.max_stream_bytes else max_bytes; var off: u64 = 0; while (true) { const want: u32 = @intCast(@min(@as(u64, s.cl.maxRead()), wire_limit + 1 - off)); const data = (try s.ask(.{ .read = .{ .fid = cur, .offset = off, .count = want } }, remote)).read; if (data.len == 0) return; if (off + data.len > wire_limit) return Error.FileTooLarge; if (directory) { var pos: usize = 0; while (pos < data.len) { if (data.len - pos < 2) return Error.Botch; const len: usize = 2 + @as(usize, std.mem.readInt(u16, data[pos..][0..2], .little)); if (len > data.len - pos) return Error.Botch; const entry = ninep.Stat.decode(data[pos..][0..len]) catch return Error.Botch; const display_dir = std.mem.trimEnd(u8, s.display_path, "/"); const row_len = display_dir.len + entry.name.len + 2 + @as(usize, @intFromBool(entry.qid.type & ninep.qtdir != 0)); if (row_len > max_bytes - out.written().len) return Error.FileTooLarge; try out.writer.print("{s}/{s}", .{ display_dir, entry.name }); if (entry.qid.type & ninep.qtdir != 0) try out.writer.writeByte('/'); try out.writer.writeByte('\n'); pos += len; } } else try out.writer.writeAll(data); off += data.len; } } fn fetchBytes( gpa: std.mem.Allocator, sock: Dial, names: []const []const u8, remote: *RemoteError, write_bytes: ?[]const u8, display_path: []const u8, read_limit: usize, ) ![]u8 { if (comptime !supported) return Error.Dial; const s = try startSession(gpa, sock, display_path); defer endSession(gpa, s); var out: std.Io.Writer.Allocating = .init(gpa); errdefer out.deinit(); try transact(s, names, &out, remote, write_bytes, read_limit); return out.toOwnedSlice(); } /// A connected session whose requests have the usual deadline. fn startSession(gpa: std.mem.Allocator, sock: Dial, display_path: []const u8) !*Session { const deadline = nowMs() +| budget_ms; const s = try gpa.create(Session); s.* = .{ .fd = -1, .deadline = deadline, .display_path = display_path }; errdefer endSession(gpa, s); if (sock == .quic) { if (comptime quic_enabled) { s.quic = quic.Connection.dial(sock.quic) catch return Error.Dial; s.fd = s.quic.?.fd; } else return Error.QuicUnavailable; } else s.fd = try connect(sock, deadline); s.cl = .init(.{ .in = &s.in, .out = &s.out }); return s; } fn endSession(gpa: std.mem.Allocator, s: *Session) void { if (quic_enabled and s.quic != null) { s.quic.?.deinit(); } else if (s.fd >= 0) _ = libc.close(s.fd); gpa.destroy(s); } /// Opens `path` on one connection and writes `first` to it (/log's /// `follow new`), all with the usual deadline; then, once `ready(ctx)` /// has said it is not done already, reads it with no deadline, a read at /// a time, until `record(ctx, bytes)` says done. An end of file or a /// dropped connection is Hangup. What `pardes --wait` blocks on. pub fn follow( gpa: std.mem.Allocator, dial: []const u8, path: []const u8, first: []const u8, ctx: anytype, comptime ready: fn (@TypeOf(ctx)) bool, comptime record: fn (@TypeOf(ctx), []const u8) bool, ) !void { if (comptime !supported) return error.Unsupported; var names: [max_depth][]const u8 = undefined; const n = try elements(path, &names); var sock_buf: [sun_path_len]u8 = undefined; const sock = try resolve(&sock_buf, dial); var remote: RemoteError = .{}; const s = try startSession(gpa, sock, path); defer endSession(gpa, s); const fid = (try walkTo(s, names[0..n], &remote)).fid; _ = try s.ask(.{ .open = .{ .fid = fid, .mode = ninep.ordwr } }, &remote); _ = try s.ask(.{ .write = .{ .fid = fid, .offset = 0, .data = first } }, &remote); if (ready(ctx)) return; s.deadline = std.math.maxInt(i64); while (true) { const data = (try s.ask(.{ .read = .{ .fid = fid, .offset = 0, .count = s.cl.maxRead() } }, &remote)).read; if (data.len == 0) return Error.Hangup; if (record(ctx, data)) return; } } fn connect(sock: Dial, deadline: i64) Error!c_int { const address: transport.Address = switch (sock) { .unix => |path| .{ .unix = path }, .tcp => |ip| .{ .tcp = ip }, .quic => return Error.QuicUnavailable, }; return transport.connectFd(address, deadline) catch return Error.Dial; } const nowMs = transport.nowMs; const poll_in: i16 = @intCast(libc.POLL.IN); const poll_out: i16 = @intCast(libc.POLL.OUT); fn nap(ms: c_int) void { _ = libc.poll(&[0]libc.pollfd{}, 0, ms); } test "a path becomes walk elements, normalised the way a shell would" { var out: [max_depth][]const u8 = undefined; try testing.expectEqual(@as(usize, 3), try elements("/pane/1/body", &out)); try testing.expectEqualStrings("pane", out[0]); try testing.expectEqualStrings("1", out[1]); try testing.expectEqualStrings("body", out[2]); try testing.expectEqual(@as(usize, 3), try elements("pane/1/body", &out)); try testing.expectEqual(@as(usize, 3), try elements("//pane//1//body//", &out)); try testing.expectEqual(@as(usize, 1), try elements("/index", &out)); try testing.expectEqual(@as(usize, 0), try elements("/", &out)); var deep: [8 * max_depth]u8 = @splat('/'); for (0..max_depth + 1) |i| deep[i * 2 + 1] = 'a'; try testing.expectError(Error.PathTooDeep, elements(deep[0 .. (max_depth + 1) * 2], &out)); } test "a bare dial resolves to the socket --9p binds, and a path is taken as given" { if (comptime !supported) return error.SkipZigTest; var buf: [sun_path_len]u8 = undefined; const named = (try resolve(&buf, "work")).unix; try testing.expect(std.mem.endsWith(u8, named, "/pardes-9p-work.sock")); var expect: [sun_path_len]u8 = undefined; var dir_buf: [sun_path_len:0]u8 = undefined; const dir = socketDir(&dir_buf).?; try testing.expectEqualStrings(socketPath(&expect, dir, "work").?, named); const path = (try resolve(&buf, "/tmp/somewhere.sock")).unix; try testing.expectEqualStrings("/tmp/somewhere.sock", path); try testing.expectEqualStrings("/tmp/somewhere.sock", (try resolve(&buf, "unix!/tmp/somewhere.sock")).unix); try testing.expectError(error.BadDial, resolve(&buf, "")); try testing.expectError(error.BadDial, resolve(&buf, "/tmp/a\x00b")); } test "a dial with nothing listening is one error and not a wait" { if (comptime !supported) return error.SkipZigTest; var names: [max_depth][]const u8 = undefined; const n = try elements("/pane/1/body", &names); var remote: RemoteError = .{}; const before = nowMs(); try testing.expectError( Error.Dial, fetchBytes(testing.allocator, .{ .unix = "/tmp/pardes-9p-no-such-socket.sock" }, names[0..n], &remote, null, "/pane/1/body", limits.max_file_bytes), ); try testing.expect(nowMs() - before < budget_ms); } test "one fetch has three msize buffers and bounded transport metadata" { const transport_bytes = if (quic_enabled) @sizeOf(?quic.Connection) else 0; try testing.expectEqual(@as(usize, msize), @as(usize, (Session{ .fd = -1, .deadline = 0, .display_path = "" }).in.len)); try testing.expect(transport_bytes <= 64); try testing.expect(@sizeOf(Session) <= 3 * msize + @sizeOf(ninep.Client) + 128 + transport_bytes); try testing.expect(@sizeOf(ninep.Client) <= 512); } test "expired sessions do not send or consume buffered protocol work" { var session: Session = .{ .fd = -1, .deadline = 0, .display_path = "" }; session.cl = .init(.{ .in = &session.in, .out = &session.out }); _ = try session.cl.submit(.{ .version = .{} }); const queued = session.cl.output().len; try testing.expect(queued > 0); try testing.expectError(Error.Timeout, session.flush()); try testing.expectEqual(queued, session.cl.output().len); try testing.expectError(Error.Timeout, session.settle()); try testing.expectError(Error.Timeout, session.wait(poll_in)); } };