From 31cb659ded4cf50af5903fc107f8c868ee3c7311 Mon Sep 17 00:00:00 2001 From: Gabriel Schneider Date: Tue, 22 Sep 2026 11:15:46 -0300 Subject: Answer 9P on the connection's task, so a session can open its own tree The editor's loop was the only thing that could answer a 9P request, which made the editor's own syscalls through a mount of its own tree -- a Look at /mnt/9p/pardes//anything under a `9ns --mntgen` view, a Save into it -- requests only the blocked loop could serve. The name-based refusal that followed (ownMountSuffix) and the in-process routing of a mount of oneself (Client.sameSession) were patches over that, and both are gone, with the mailbox that shipped every request to the editor's thread. One rule replaces them, `pardes.turn`: the core is single-threaded, the editor's thread has the turn by default and gives it up in two kinds of gap -- while it waits for input and while a step of it is out in a host syscall -- and a cloud9 connection task takes it in those gaps to answer. `out` counts the steps that are out, from any thread: while one is, the core reads consistently but that step still holds pointers into it, so a request that would change a pane (a write, a truncation, an rmdir) is parked in the engine and retried when the turn is next given up with nothing out, and the editor's own wake waits for the count to reach zero. It is never a write of its own that a step waits on out there -- writes come from a shell performing a save between steps -- so a parked request is never the syscall's own, and making a pane or rendering a screen need not park: every yield sits before its step's mutation, so the layout and the surface are whole under it. A changing request that queued effects is answered once the editor has performed them (`echo Save > exec` returns with the file written, as acme's `put` does), and it settles the way a step does, because without that a /log reader waited for the user's next keystroke. Every host syscall on a user path has to give the turn up, not fs.zig's alone: the first end-to-end run hung in `inotify_add_watch` performing the new pane's watch effect. PDFs and images are read whole at open, so no draw goes out into the host. The core's allocator takes its fixed buffer through the lock-free interface, since a connection task allocates while the editor's thread is out in a syscall that allocates too. A Restore puts the replacement in first and releases every task waiting on the old core. cloud9 (pinned at eb1a104) parks an open, a truncating wstat, a clunk and a remove on `again`, not only reads and writes, and answers a parked job whose fid was clunked without asking the backend. Verified: test/selfmount.py runs the editor under `9ns --mntgen` and Looks at, reads and Saves its own tree through the mount; a unit test pins that a change parks while the editor is out mid-step and lands when it rests, while a read is answered in the window. 9P over the Unix socket against a tty session, same machine, Debug builds: a read of /index 278us -> 61us, a truncating body write 1184us -> 609us, exec Save 718us -> 583us; the gesture benchmark is unchanged (geometric mean 0.997 over 53 cells). Also from the reviews: a notice chip over an image or PDF pane was painted out by the picture drawn after the cells, so pictures give up the rows; in the GUI a tree-sitter context band painted over the chip, so body layers are emitted first; a message is one row of printable text, its 256-byte cut never leaves half a glyph, and one wider than its pane keeps its tail (the file name, the reason) rather than its head. Co-Authored-By: Claude Fable 5.1 --- src/9p_io.zig | 511 +++++++++++++++++++++++++--------------------------------- 1 file changed, 220 insertions(+), 291 deletions(-) (limited to 'src/9p_io.zig') diff --git a/src/9p_io.zig b/src/9p_io.zig index 08deae7a..7e09fb27 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -62,7 +62,7 @@ pub fn ensureSocketDir(dir: [:0]const u8) bool { const prefix = "pardes-9p-"; -pub const msize: u32 = 8192; +pub const msize = ninep.msize; pub const max_conns = 4; @@ -143,76 +143,17 @@ 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 hands every backend request to the editor's thread -/// through `Slot` (see `tick`). QUIC keeps the poll loop below. +/// 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; -/// Answers one backend request on the editor's thread. -fn step(core: *pardes.Pardes, srv: *Srv, req: pardes.ctlfs.Req) void { - core.update(.{ .fs_req = req }); - var answered = false; - while (core.nextEffect()) |effect| { - if (effect == .fs_reply) { - const reply = effect.fs_reply; - if (reply.tag == req.tag) answered = true; - srv.reply(&reply, core.fsPayload(reply)); - } else core.perform(effect); - } - if (!answered) { - const reply = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); - srv.reply(&reply, ""); - } -} - -/// Pays the backend a request whose connection is gone; replies go nowhere. -fn pay(core: *pardes.Pardes, req: pardes.ctlfs.Req) void { - core.update(.{ .fs_req = req }); - while (core.nextEffect()) |effect| if (effect != .fs_reply) core.perform(effect); -} - -/// One runner slot's mailbox: the requests its connection task handed over, -/// answered by the editor's thread in `tick`. A hangup's releases (at most -/// one per fid) and one outstanding request are the most it ever holds, -/// because a closing connection waits for its entries to be paid before -/// the slot is reused. -const Slot = struct { - const capacity = ninep.editor.fid_capacity + 2; - const Entry = struct { gen: u32, req: pardes.ctlfs.Req }; - - mutex: std.Io.Mutex = .init, - queue: [capacity]Entry = undefined, - head: usize = 0, - len: usize = 0, - /// Bumped when a connection takes the slot; entries carry the value. - gen: u32 = 0, - closing: bool = false, - drained: std.Io.Event = .unset, - /// Requests refused for want of room; never expected. - refused: u32 = 0, - - fn push(s: *Slot, req: pardes.ctlfs.Req) bool { - if (s.len == s.queue.len) { - s.refused += 1; - return false; - } - s.queue[(s.head + s.len) % capacity] = .{ .gen = s.gen, .req = req }; - s.len += 1; - return true; - } - - fn pop(s: *Slot) ?Entry { - if (s.len == 0) return null; - const e = s.queue[s.head]; - s.head = (s.head + 1) % capacity; - s.len -= 1; - return e; - } -}; - const quic_slots = if (quic_enabled) max_conns else 0; /// A QUIC connection on the poll path. @@ -229,12 +170,10 @@ 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, - slots: [max_conns]Slot = @splat(.{}), stopping: std.atomic.Value(bool) = .init(false), - /// While `reset` runs: the generation each slot is paying off; anything - /// newer waits for the replacement core. - resetting: ?[max_conns]u32 = null, 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, @@ -264,47 +203,59 @@ pub const Listener = struct { return @ptrCast(@alignCast(ctx.?)); } - /// The runner's handler: on the connection's task. - fn onOpened(ctx: ?*anyopaque, conn: *Runner.Conn) void { - const l = of(ctx); - const s = &l.slots[conn.index]; - s.mutex.lockUncancelable(l.io); - s.gen +%= 1; - s.mutex.unlock(l.io); - } - + /// 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`). fn onServe(ctx: ?*anyopaque, conn: *Runner.Conn, req: pardes.ctlfs.Req) void { const l = of(ctx); - const s = &l.slots[conn.index]; - s.mutex.lockUncancelable(l.io); - const queued = s.push(req); - s.mutex.unlock(l.io); - if (!queued) { - const reply = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); - conn.reply(&reply, ""); - return; + // 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(); + if (!quiet and pardes.ctlfs.needsQuiet(req)) { + pardes.turn.parked = true; + const later: pardes.ctlfs.Reply = .{ .tag = req.tag, .status = .again }; + return conn.reply(&later, ""); } + const core = l.core; + const epoch = pardes.turn.epoch; + const restores = pardes.turn.restores; + const reply = core.serveFs(req); + if (req.op == .release and quiet) l.collectOs(); + if (!pardes.ctlfs.changesPane(req)) 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(); - } - - /// Holds the slot until the editor's thread has paid what the - /// connection still owes, so a successor never inherits its entries. - fn onClosed(ctx: ?*anyopaque, conn: *Runner.Conn) void { + const restoring = core.restore_req != null; + if (core.effects_len != 0) pardes.turn.awaitSettled(epoch); + // Once a wait returns `core` may be gone: a Restore meanwhile put a + // replacement in (`reset`) and is hanging this connection up, which + // is a Restore's answer. Only one that failed leaves this connection + // here to answer. + if (l.core != core) return; + if (restoring) pardes.turn.awaitRestored(restores); + if (l.core != core) return; + conn.reply(&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); - const s = &l.slots[conn.index]; if (l.stopping.load(.acquire)) return; - s.mutex.lockUncancelable(l.io); - if (s.len == 0) { - s.mutex.unlock(l.io); - return; - } - s.closing = true; - s.drained.reset(); - s.mutex.unlock(l.io); - l.kick(); - s.drained.wait(l.io) catch {}; + 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); @@ -411,68 +362,46 @@ pub const Listener = struct { c.draining = true; } - pub const Drained = struct { count: usize = 0, pending: bool = false }; - - /// Answers what every connection has asked since the last call, retries - /// their parked reads against the editor's current state and has the - /// runner send the replies. Call it after each editor update, on the - /// editor's thread; `pending` says QUIC has more to do right away. - pub fn drain(l: *Listener, core: *pardes.Pardes) Drained { - core.fs.socket_path = l.path(); - core.fs.tcp_address = l.tcp_address; - core.fs.quic_address = l.quic_address; - l.expire(); - var result: Drained = .{}; - if (comptime supported) { - for (&l.slots, 0..) |*s, i| { - const conn = &l.runner.conns[i]; - s.mutex.lockUncancelable(l.io); - const only: ?u32 = if (l.resetting) |gens| gens[i] else null; - var count: usize = 0; - if (s.len != 0 or conn.live()) { - conn.lock(); - while (s.len != 0) { - if (only) |gen| if (s.queue[s.head].gen != gen) break; - const e = s.pop().?; - if (e.gen == s.gen) step(core, &conn.engine, e.req) else if (e.req.op == .release) pay(core, e.req); - count += 1; - } - if (conn.live() and (only == null or only.? == s.gen)) { - while (conn.engine.retry()) |req| { - step(core, &conn.engine, req); - count += 1; - } - while (conn.engine.next()) |req| { - step(core, &conn.engine, req); - count += 1; - } - } - const owed = conn.engine.output().len != 0; - conn.unlock(); - if (owed) conn.flush(); - } - if (s.closing and s.len == 0) { - s.closing = false; - s.drained.set(l.io); - } - s.mutex.unlock(l.io); - result.count += count; - if (count != 0) l.collectOs(core); - } - } + /// 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 { + 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| { - step(core, &conn.srv, req); - l.collectOs(core); + l.step(&conn.srv, req); count += 1; } while (count < 64) { const req = conn.srv.next() orelse break; - step(core, &conn.srv, req); - l.collectOs(core); + l.step(&conn.srv, req); count += 1; } result.count += count; @@ -482,64 +411,44 @@ pub const Listener = struct { } else l.flush(@intCast(i)); if (conn.quic) |*connection| result.pending = result.pending or connection.pending(); } + l.arm(); } return result; } - /// `drain` plus the QUIC listener's events, accepts and reads. - pub fn tick(l: *Listener, core: *pardes.Pardes) Drained { - 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)); - } - const result = l.drain(core); - l.arm(); - return result; - } - - /// Hangs every connection up and pays `core` what they owed, so that a - /// replacement core starts with no client holding anything. A client - /// that connects meanwhile is not answered until the caller has put - /// the replacement in and calls `tick` again. - pub fn reset(l: *Listener, core: *pardes.Pardes) void { + /// 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: it + /// finds the core changed and answers nothing, its connection 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) { - var gens: [max_conns]u32 = undefined; - for (&l.slots, &gens) |*s, *gen| { - s.mutex.lockUncancelable(l.io); - gen.* = s.gen; - s.mutex.unlock(l.io); - } - l.resetting = gens; - defer l.resetting = null; l.runner.closeAll(); + pardes.turn.settle(); + pardes.turn.restoreSettled(); + pardes.turn.rest(); const deadline = nowMs() +| 2000; - while (true) { - const drained = l.drain(core); - var idle = !drained.pending; - for (&l.slots, &gens, 0..) |*s, gen, i| { - s.mutex.lockUncancelable(l.io); - if (s.len != 0 and s.queue[s.head].gen == gen) idle = false; - if (s.gen == gen and l.runner.conns[i].live()) idle = false; - s.mutex.unlock(l.io); - } - if (idle or nowMs() >= deadline) break; - Client.nap(1); - } - } else while (l.drain(core).pending) {} - l.collectOs(core); + while (l.runner.count() != 0 and nowMs() < deadline) Client.nap(1); + pardes.turn.wake(); + } + l.collectOs(); for (&l.conns) |*conn| conn.accepted_ms = 0; l.arm(); } - fn collectOs(l: *Listener, core: *pardes.Pardes) void { + /// 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]; @@ -552,17 +461,10 @@ pub const Listener = struct { } } - /// Whether any connection, or a release still to be paid, names `node`. + /// Whether any connection names `node`. fn references(l: *Listener, node: u64) bool { if (comptime supported) { - for (&l.slots, 0..) |*s, i| { - const conn = &l.runner.conns[i]; - s.mutex.lockUncancelable(l.io); - defer s.mutex.unlock(l.io); - var k: usize = 0; - while (k < s.len) : (k += 1) { - if (s.queue[(s.head + k) % Slot.capacity].req.node == node) return true; - } + for (&l.runner.conns) |*conn| { if (!conn.live()) continue; conn.lock(); defer conn.unlock(); @@ -576,16 +478,12 @@ pub const Listener = struct { return false; } - /// Where the runner's tasks report that a request waits for `tick`. - pub fn setWake(l: *Listener, ctx: ?*anyopaque, wake: *const fn (?*anyopaque) void) void { + /// 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; - } - - /// `setWake`, 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.setWake(ctx, wake); if (comptime !quic_enabled) return; if (l.quic == null or l.watcher != null) return; if (libc.pipe(&l.control) != 0) return error.PipeFailed; @@ -735,7 +633,15 @@ pub const Listener = struct { if (comptime quic_enabled) { if (l.quic) |*listener| listener.deinit(); } - if (comptime supported) l.runner.stop(); + 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(); + l.runner.stop(); + pardes.turn.wake(); + } + pardes.turn.wake_parked = 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]); @@ -752,7 +658,9 @@ pub fn socketPath(buf: *[sun_path_len]u8, dir: []const u8, name: []const u8) ?[: return std.fmt.bufPrintSentinel(buf, "{s}/" ++ prefix ++ "{s}.sock", .{ dir, name }, 0) catch null; } -pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: []const u8, tcp_dial: ?[]const u8, quic_dial: ?[]const u8) ?*Listener { +/// 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 { @@ -761,11 +669,14 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: [ }; if (!ensureSocketDir(dir)) return null; const l = gpa.create(Listener) catch return null; - l.* = .{ .io = io }; + l.* = .{ .io = io, .core = core }; + pardes.turn.start(io); + pardes.turn.wake_parked = Listener.wakeParked; + pardes.turn.wake_ctx = l; l.runner.init(.{ .io = io, .root = pardes.ctlfs.root, - .handler = .{ .ctx = l, .serve = Listener.onServe, .opened = Listener.onOpened, .closed = Listener.onClosed }, + .handler = .{ .ctx = l, .serve = Listener.onServe }, .greet_timeout_ms = Listener.greet_deadline_ms, }); const entry_name = if (named.len != 0) named else fallback; @@ -838,6 +749,9 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: [ } } 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; } @@ -1160,33 +1074,6 @@ test "TCP addresses are numeric and normalize mapped IPv4" { } else try testing.expectError(error.QuicUnavailable, Client.validateDial("quic!127.0.0.1!5640")); } -test "same-session TCP mounts compare canonical endpoints and local wildcard destinations" { - if (comptime !supported) return error.SkipZigTest; - const Case = struct { bound: []const u8, dial: []const u8, same: bool }; - for ([_]Case{ - .{ .bound = "tcp!127.0.0.1!5640", .dial = "tcp!::ffff:127.0.0.1!5640", .same = true }, - .{ .bound = "tcp!::ffff:127.0.0.1!5640", .dial = "tcp!127.0.0.1!5640", .same = true }, - .{ .bound = "tcp!::1!5640", .dial = "tcp!0:0:0:0:0:0:0:1!5640", .same = true }, - .{ .bound = "tcp!127.0.0.1!5640", .dial = "tcp!0.0.0.0!5640", .same = true }, - .{ .bound = "tcp!::1!5640", .dial = "tcp!::!5640", .same = true }, - .{ .bound = "tcp!0.0.0.0!5640", .dial = "tcp!127.0.0.2!5640", .same = true }, - .{ .bound = "tcp!::!5640", .dial = "tcp!::1!5640", .same = true }, - .{ .bound = "tcp!127.0.0.1!5640", .dial = "tcp!127.0.0.1!5641", .same = false }, - .{ .bound = "tcp!0.0.0.0!5640", .dial = "tcp!192.0.2.1!5640", .same = false }, - .{ .bound = "tcp!::!5640", .dial = "tcp!2001:db8::1!5640", .same = false }, - .{ .bound = "tcp!::!5640", .dial = "tcp!127.0.0.1!5640", .same = false }, - }) |c| try testing.expectEqual(c.same, Client.sameSession(c.dial, "", try networkAddress(c.bound, true), null)); - try testing.expect(Client.sameSession("/tmp/pardes-owned.sock", "/tmp/pardes-owned.sock", null, null)); - try testing.expect(Client.sameSession("unix!/tmp/pardes-owned.sock", "/tmp/pardes-owned.sock", null, null)); - try testing.expect(!Client.sameSession("tcp!127.0.0.1!5640", "/tmp/pardes-owned.sock", null, null)); - if (quic_enabled) { - const endpoint = try networkAddress("quic!127.0.0.1!5640", false); - try testing.expect(Client.sameSession("quic!::ffff:127.0.0.1!5640", "", null, endpoint)); - try testing.expect(!Client.sameSession("tcp!127.0.0.1!5640", "", null, endpoint)); - try testing.expect(!Client.sameSession("quic!127.0.0.1!5640", "", endpoint, null)); - } -} - extern "c" fn mkdtemp(template: [*:0]u8) ?[*:0]u8; extern "c" fn rmdir(path: [*:0]const u8) c_int; @@ -1225,7 +1112,9 @@ test "a listening editor posts itself into the 9P registry and unposts on stop" var link: [sun_path_len]u8 = undefined; { - const l = listen(testing.io, gpa, "unit", "", null, null) orelse return error.ListenFailed; + 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. @@ -1312,7 +1201,7 @@ test "Unix TCP and QUIC share one listener through reads writes reconnects and r 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, "roundtrip", "", if (tcp) bind else null, if (tcp) null else bind) orelse return error.ListenFailed; + 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); @@ -1331,39 +1220,108 @@ test "Unix TCP and QUIC share one listener through reads writes reconnects and r 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) { - _ = l.tick(p); + if (quic_enabled) { + pardes.turn.wake(); + _ = l.tick(); + pardes.turn.rest(); + } Client.nap(1); } + pardes.turn.wake(); try testing.expect(worker.done.load(.acquire)); if (worker.failure) |err| return err; l.reset(p); try testing.expectEqual(@as(usize, 0), l.runner.count()); - for (&l.slots) |*slot| try testing.expectEqual(@as(usize, 0), slot.len); for (&l.conns, 0..) |conn, i| try testing.expect(!l.live(@intCast(i)) and !conn.draining); for (p.fs.snapshots) |snapshot| try testing.expect(snapshot.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; + defer _ = rmdir(&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)); + + 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); +} + 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; const fallback = std.fmt.bufPrint(&name, "{d}", .{@as(u32, @intCast(libc.getpid()))}) catch unreachable; - const listener = listen(io, gpa, core.opts.ninep_name, fallback, core.opts.ninep_tcp, core.opts.ninep_quic) orelse { + 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; }; - core.fs.socket_path = listener.path(); - // So the editor can recognise its own tree by name and refuse to walk into - // it through the mount, which would deadlock the loop that serves it. - pardes.filesystem.noteOwnSocket(listener.path()); - core.fs.tcp_address = listener.tcp_address; - core.fs.quic_address = listener.quic_address; - return listener; } /// What a pane shell is told about the editor above it, set into this @@ -1423,7 +1381,7 @@ test "9P shell environment states being inside pardes apart from how to reach it const listener = try testing.allocator.create(Listener); defer testing.allocator.destroy(listener); - listener.* = .{ .io = testing.io }; + 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; @@ -1549,35 +1507,6 @@ pub const Client = struct { _ = try resolve(&buf, dial); } - pub fn sameSession(dial: []const u8, socket_path: []const u8, tcp_address: ?std.Io.net.IpAddress, quic_address: ?std.Io.net.IpAddress) bool { - if (comptime !supported) return false; - var buf: [sun_path_len]u8 = undefined; - const address = resolve(&buf, dial) catch return false; - switch (address) { - .unix => |path| return socket_path.len != 0 and std.mem.eql(u8, path, socket_path), - .tcp, .quic => |destination| { - var ip = destination; - switch (ip) { - .ip4 => |v4| if (std.mem.allEqual(u8, &v4.bytes, 0)) { - ip = .{ .ip4 = .loopback(v4.port) }; - }, - .ip6 => |v6| if (std.mem.allEqual(u8, &v6.bytes, 0)) { - ip = .{ .ip6 = .loopback(v6.port) }; - }, - } - const bound = canonicalIp((if (address == .tcp) tcp_address else quic_address) orelse return false); - if (ip.getPort() != bound.getPort() or @as(std.Io.net.IpAddress.Family, ip) != @as(std.Io.net.IpAddress.Family, bound)) return false; - if (ip.eql(&bound)) return true; - const wildcard = switch (bound) { - .ip4 => |v4| std.mem.allEqual(u8, &v4.bytes, 0), - .ip6 => |v6| std.mem.allEqual(u8, &v6.bytes, 0), - }; - if (wildcard) return localIp(ip); - return false; - }, - } - } - const Session = struct { fd: c_int, quic: if (quic_enabled) ?quic.Connection else void = if (quic_enabled) null else {}, -- cgit v1.3