diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/9p.zig | 19 | ||||
| -rw-r--r-- | src/9p_io.zig | 505 | ||||
| -rw-r--r-- | src/CHANGELOG.md | 13 | ||||
| -rw-r--r-- | src/detached/server.zig | 27 | ||||
| -rw-r--r-- | src/detached/wire.zig | 8 | ||||
| -rw-r--r-- | src/file_watch.zig | 6 | ||||
| -rw-r--r-- | src/fs.zig | 444 | ||||
| -rw-r--r-- | src/gui/gui.zig | 30 | ||||
| -rw-r--r-- | src/host_io.zig | 3 | ||||
| -rw-r--r-- | src/macos.zig | 125 | ||||
| -rw-r--r-- | src/memory.zig | 50 | ||||
| -rw-r--r-- | src/ninep/testing.zig | 7 | ||||
| -rw-r--r-- | src/ninep/tree.zig | 16 | ||||
| -rw-r--r-- | src/panes.zig | 45 | ||||
| -rw-r--r-- | src/pardes.zig | 282 | ||||
| -rw-r--r-- | src/tty/tty.zig | 11 |
16 files changed, 900 insertions, 691 deletions
@@ -52,9 +52,16 @@ pub const msize_min = cloud9.fs.msize_min; /// The smallest msize a native listener offers (src/9p_io.zig). pub const min_msize: u32 = 4096; +/// What a native listener negotiates. +pub const msize: u32 = 8192; -/// The editor's engine: native fid count, native file names. -pub const editor: cloud9.fs.Options = .{ .fid_capacity = max_fids, .name_capacity = 255 }; +/// The editor's engine: native fid count, native file names, and room to +/// park a write of any size a frame can carry -- a request that would change +/// a pane while the editor is out in a syscall waits in the engine +/// (src/9p_io.zig, `serve`), and a write it cannot keep would fail instead. +/// The largest write is msize less the Twrite header, 23 bytes (the test +/// below asks a client), one byte more than `iohdrsz` allows for. +pub const editor: cloud9.fs.Options = .{ .fid_capacity = max_fids, .name_capacity = 255, .park_data_max = msize - 23 }; /// The board's engine: the GPIO firmware's budget (src/esp32p4_9p.zig). pub const board: cloud9.fs.Options = .{ .fid_capacity = board_fids, .name_capacity = board_name_capacity }; @@ -70,6 +77,14 @@ const Stub = struct { pub const Reply = cloud9.fs.Reply; }; +test "9p server: the editor's engine can park the largest write a client sends" { + var in: [4096]u8 = undefined; + var out: [4096]u8 = undefined; + var client: Client = .init(.{ .in = &in, .out = &out }); + client.msize = msize; + try testing.expect(editor.park_data_max >= client.maxWrite()); +} + test "9p server: board and native capacities size the actual fid storage" { const Board = Server(Stub, board); const Native = Server(Stub, editor); 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(); + 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, ""); } - /// 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 { + /// 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 }; + /// 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 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); - } - } + /// 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 {}, diff --git a/src/CHANGELOG.md b/src/CHANGELOG.md index e518f97a..0083f97f 100644 --- a/src/CHANGELOG.md +++ b/src/CHANGELOG.md @@ -2,6 +2,19 @@ ## 0.0.3 +- Answer 9P on the connection's task instead of the editor's loop. The core + is single-threaded and `pardes.turn` says whose turn it is with it: the + editor's by default, given up while it waits for input and while a step of + it is out in a syscall, and taken by a connection task to answer a request. + A request that would change a pane while a step is out parks in the engine + and is retried when the editor rests. So a session can open, read and save + its own tree through a mount -- which used to hang it outright, and was + then refused by name -- and the name-based refusal, the in-process routing + of a mount of oneself and the mailbox that shipped every request to the + editor's thread are gone. +- A notice wider than its pane keeps its tail rather than its head: the end + of a message is the file name or the reason. + - Give pty children their own terminal identity. `TERM`, `COLORTERM` and `TERM_PROGRAM` are now written by pardes rather than inherited from whatever launched it, and the child is exec'd with that environment. A macOS `.app` diff --git a/src/detached/server.zig b/src/detached/server.zig index ac9af2ee..3e98fc9d 100644 --- a/src/detached/server.zig +++ b/src/detached/server.zig @@ -156,7 +156,8 @@ pub const Session = struct { inotify_fd: c_int = -1, ninep: ?*ninep_io.Listener = null, ninep_pending: bool = false, - /// Set by the 9P runner's tasks: a request waits for `tick`. + /// A 9P request changed the core since the wake byte was last drained: + /// the next poll returns at once, to draw it. ninep_wake: std.atomic.Value(bool) = .init(false), watches: file_watch.Table = @splat(null), check_files: bool = false, @@ -240,7 +241,7 @@ pub const Session = struct { file_watch.watchPane(s.inotify_fd, &s.watches, @intCast(pane), null, 0, .{ .text = 0 }); _ = file_watch.applyThemeEffect(s.core, s.gpa, s.inotify_fd, &s.watches, 0, false, false); s.check_files = false; - if (s.ninep) |listener| listener.reset(s.core); + if (s.ninep) |listener| listener.reset(replacement); s.ninep_pending = false; replacement.host = s.host(); s.core.deinit(); @@ -587,7 +588,7 @@ pub const Session = struct { if (s.origins()) |c| s.send(c, .detach); } - /// The 9P runner has a request for `tick`: wake the poll loop. + /// A 9P request changed the core: wake the poll loop to draw it. fn wakeNinep(ctx: ?*anyopaque) void { const s = of(ctx); s.ninep_wake.store(true, .release); @@ -596,7 +597,7 @@ pub const Session = struct { fn pollFrame(ctx: ?*anyopaque) void { const s = of(ctx); - if (s.ninep) |l| s.ninep_pending = l.tick(s.core).pending; + if (s.ninep) |l| s.ninep_pending = l.tick().pending; for (&s.ptys, 0..) |*pt, pane| { if (pt.fd < 0) continue; var lbuf: [pardes.memory.limits.host_path_cap + 1]u8 = undefined; @@ -805,11 +806,18 @@ pub const Session = struct { } } } - if (n == 0) return nap(if (timeout_ms == 0) 16 else timeout_ms); + // The wait is the 9P connections' turn with the core. + if (n == 0) { + pardes.turn.rest(); + defer pardes.turn.wake(); + return nap(if (timeout_ms == 0) 16 else timeout_ms); + } var timeout: c_int = if (timeout_ms == 0) -1 else @intCast(@min(timeout_ms, std.math.maxInt(c_int))); if (s.nextWake(now)) |due| timeout = if (timeout < 0) due else @min(timeout, due); if (completed or s.check_files or regridded or s.ninep_pending) timeout = 0; + pardes.turn.rest(); const ready = libc.poll(&fds, @intCast(n), timeout); + pardes.turn.wake(); if (ready > 0) s.dispatch(fds[0..n], src[0..n]); _ = s.drainCompletions(true); s.expire(monotonicMs()); @@ -1134,18 +1142,17 @@ pub fn run(init: std.process.Init, opts: pardes.Options, name: []const u8) !void return error.NoSocket; } - session.ninep = ninep_io.listen(init.io, gpa, opts.ninep_name, name, opts.ninep_tcp, opts.ninep_quic); + session.ninep = ninep_io.listen(init.io, gpa, core, opts.ninep_name, name, opts.ninep_tcp, opts.ninep_quic); if (session.ninep == null) return error.ListenFailed; - session.ninep.?.setWake(&session, Session.wakeNinep); - core.fs.socket_path = session.ninep.?.path(); - core.fs.tcp_address = session.ninep.?.tcp_address; - core.fs.quic_address = session.ninep.?.quic_address; + session.ninep.?.wake_ctx = &session; + session.ninep.?.wake = Session.wakeNinep; const h = session.host(); core.host = h; while (core.nextEffect()) |effect| core.perform(effect); session.in_loop = true; while (!session.core.quit) { + pardes.turn.restoreSettled(); try session.core.pump(h); if (session.core.quit) break; if (session.core.takeRestore()) |path| restore: { diff --git a/src/detached/wire.zig b/src/detached/wire.zig index 6805d980..20b5c093 100644 --- a/src/detached/wire.zig +++ b/src/detached/wire.zig @@ -509,7 +509,7 @@ const Ul = @FieldType(pardes.CellStyle, "ul"); /// which is the whole property the six hand-written switches this replaced /// existed for, in a form that cannot fall out of step. A variant that has to /// travel under a DIFFERENT name than the core's gets an arm of its own before -/// the `inline else`, exactly as `clientTag` keeps `.fs_req`. +/// the `inline else`, exactly as `clientTag` keeps the machine-local reports. fn putUl(w: *Writer, u: Ul) Error!void { try w.putByte(switch (u) { inline else => |t| @intFromEnum(@field(UlTag, @tagName(t))), @@ -819,7 +819,7 @@ fn clientTag(msg: ClientMsg) ClientTag { return switch (msg) { .event => |ev| switch (ev) { // Machine-local reports and 9P requests belong to the session owner. - .output, .eof, .lsp_resp, .pipe_resp, .file_changed, .tick, .fs_req => unreachable, + .output, .eof, .lsp_resp, .pipe_resp, .file_changed, .tick => unreachable, inline else => |_, t| @field(ClientTag, @tagName(t)), }, inline else => |_, t| @field(ClientTag, @tagName(t)), @@ -946,7 +946,7 @@ pub fn encodeClient(out: []u8, msg: ClientMsg) Error![]const u8 { .touch_scroll => |v| try w.putF32(v), .pointer_leave => {}, // See `clientTag`: no tag, so nothing to encode. - .output, .eof, .lsp_resp, .pipe_resp, .file_changed, .tick, .fs_req => unreachable, + .output, .eof, .lsp_resp, .pipe_resp, .file_changed, .tick => unreachable, }, } try finishMessage(&w, at); @@ -1076,7 +1076,7 @@ pub fn clientBound(msg: ClientMsg) usize { .paste => |b| b.len, .command => |line| line.len, // See `clientTag`: not on this wire in this direction. - .output, .eof, .lsp_resp, .pipe_resp, .file_changed, .tick, .fs_req => unreachable, + .output, .eof, .lsp_resp, .pipe_resp, .file_changed, .tick => unreachable, }, }; } diff --git a/src/file_watch.zig b/src/file_watch.zig index b2908eb7..dfcce65f 100644 --- a/src/file_watch.zig +++ b/src/file_watch.zig @@ -382,6 +382,10 @@ fn watchPath( } const declared_path = path orelse return; const watched_path = filesystem.localPath(declared_path) orelse return; + // The path may lie in a mount this editor serves, and marking it is a + // stat and a lookup out there: the turn goes out with them. + pardes.turn.yield(); + defer pardes.turn.back(); const stat = std.Io.Dir.cwd().statFile(std.Io.Threaded.global_single_threaded.io(), watched_path, .{}) catch null; const dir = if (stat != null and stat.?.kind == .directory) watched_path else std.fs.path.dirname(watched_path) orelse "."; var dir_buf: [4096:0]u8 = undefined; @@ -601,6 +605,8 @@ pub fn reloadTheme( pub fn identify(io: std.Io, path: []const u8) !Identity { const native = filesystem.localPath(path) orelse return error.NonLocalPath; + pardes.turn.yield(); + defer pardes.turn.back(); const stat = try std.Io.Dir.cwd().statFile(io, native, .{}); if (stat.kind != .file) return error.NotFile; return .{ @@ -64,10 +64,9 @@ pub const WriteError = error{ WriteFailed, }; +/// The host takes the bytes. The caller has given the turn up (`write`), +/// because the path may be a mount this editor serves. pub fn writeFile(path: []const u8, bytes: []const u8) WriteError!void { - // No core here to answer from, so the only safe answer is no answer: an - // open(2) into our own mount is the call that never returns. - if (isOwnMount(path)) return error.OpenFailed; var pathbuf: [4096:0]u8 = undefined; if (path.len >= pathbuf.len) return error.PathTooLong; if (std.mem.indexOfScalar(u8, path, 0) != null) return error.OpenFailed; @@ -124,12 +123,21 @@ fn osNode(p: *pardes.Pardes, path: []const u8) !u64 { return node; } +/// The host filesystem under /os. Every syscall here may go out through a +/// mount this editor serves, so the turn is given up around each and taken +/// back before the core is touched; `path` is the table entry's own copy and +/// outlives the yield, and nothing is staged into the shared reply buffer +/// until the turn is back. pub fn osHandle(p: *pardes.Pardes, req: Req) Reply { if (comptime !platform_has_fs) return Reply.fail(req.tag, E.NOENT); const path = osPath(p, req.node) orelse return Reply.fail(req.tag, E.NOENT); if (req.op == .release) return .{ .tag = req.tag }; const io = std.Io.Threaded.global_single_threaded.io(); - const stat = std.Io.Dir.cwd().statFile(io, path, .{}) catch return Reply.fail(req.tag, E.NOENT); + const stat = stat: { + pardes.turn.yield(); + defer pardes.turn.back(); + break :stat std.Io.Dir.cwd().statFile(io, path, .{}) catch return Reply.fail(req.tag, E.NOENT); + }; const attr: Reply.Attr = .{ .name = if (req.node == os_root) "os" else std.fs.path.basename(path), .node = req.node, .dir = stat.kind == .directory, .size = stat.size, .mode = if (stat.kind == .directory) 0o755 else 0o644, .mtime = std.math.cast(u32, stat.mtime.toSeconds()) orelse 0 }; switch (req.op) { .getattr => return .{ .tag = req.tag, .attr = attr }, @@ -150,52 +158,67 @@ pub fn osHandle(p: *pardes.Pardes, req: Req) Reply { return osHandle(p, .{ .tag = req.tag, .op = .getattr, .node = node }); }, .readdir => { - var dir = std.Io.Dir.cwd().openDir(io, path, .{ .iterate = true }) catch return Reply.fail(req.tag, E.NOTDIR); - defer dir.close(io); - var it = dir.iterate(); - const out = p.fs.stage(p.gpa); - var skip = req.off; - while (it.next(io) catch return Reply.fail(req.tag, E.IO)) |entry| { - var buf: [4096]u8 = undefined; - const joined = std.fmt.bufPrint(&buf, "{s}/{s}", .{ std.mem.trimEnd(u8, path, "/"), entry.name }) catch continue; - var normalized_buf: [4096]u8 = undefined; - const normalized = resolveOs(joined, &normalized_buf) orelse continue; - const child = dir.statFile(io, entry.name, .{}) catch continue; - if (skip > 0) { - skip -= 1; - continue; + // Listed into a buffer of this request's own first; the shared + // reply buffer is another request's to use while the turn is out. + var listing: [16 * 1024]u8 = undefined; + var fixed: std.heap.FixedBufferAllocator = .init(&listing); + var listed: std.ArrayList(u8) = .empty; + { + pardes.turn.yield(); + defer pardes.turn.back(); + var dir = std.Io.Dir.cwd().openDir(io, path, .{ .iterate = true }) catch return Reply.fail(req.tag, E.NOTDIR); + defer dir.close(io); + var it = dir.iterate(); + var skip = req.off; + while (it.next(io) catch return Reply.fail(req.tag, E.IO)) |entry| { + var buf: [4096]u8 = undefined; + const joined = std.fmt.bufPrint(&buf, "{s}/{s}", .{ std.mem.trimEnd(u8, path, "/"), entry.name }) catch continue; + var normalized_buf: [4096]u8 = undefined; + const normalized = resolveOs(joined, &normalized_buf) orelse continue; + const child = dir.statFile(io, entry.name, .{}) catch continue; + if (skip > 0) { + skip -= 1; + continue; + } + const node = if (std.mem.eql(u8, normalized.path, "/")) os_root else os_node | (std.hash.Wyhash.hash(0, normalized.path) & (os_node - 1)); + tree.stageDirent(&listed, fixed.allocator(), node, child.kind == .directory, entry.name); + if (listed.items.len >= @max(req.size, 512)) break; } - const node = if (std.mem.eql(u8, normalized.path, "/")) os_root else os_node | (std.hash.Wyhash.hash(0, normalized.path) & (os_node - 1)); - tree.stageDirent(out, p.gpa, node, child.kind == .directory, entry.name); - if (out.items.len >= @max(req.size, 512)) break; } + const out = p.fs.stage(p.gpa); + out.appendSlice(p.gpa, listed.items) catch return Reply.fail(req.tag, E.NOMEM); return .{ .tag = req.tag, .payload = .{ .staged = @intCast(out.items.len) } }; }, .read, .write, .setattr => { if (stat.kind != .file) return Reply.fail(req.tag, E.PERM); var z: [4096]u8 = undefined; const path_z = std.fmt.bufPrintSentinel(&z, "{s}", .{path}, 0) catch return Reply.fail(req.tag, E.NOENT); - const fd = libc.open(path_z, .{ .ACCMODE = if (req.op == .read) .RDONLY else .WRONLY, .NONBLOCK = true, .CLOEXEC = true }); - if (fd < 0) return Reply.fail(req.tag, E.PERM); - defer _ = libc.close(fd); - if (req.op == .setattr) { - if (!req.truncate or req.off != 0 or libc.ftruncate(fd, 0) != 0) return Reply.fail(req.tag, E.INVAL); - var truncated = attr; - truncated.size = 0; - return .{ .tag = req.tag, .attr = truncated }; - } - if (req.off > std.math.maxInt(i64)) return Reply.fail(req.tag, E.INVAL); - if (libc.lseek(fd, @intCast(req.off), libc.SEEK.SET) < 0) return Reply.fail(req.tag, E.IO); - if (req.op == .write) { - const written = libc.write(fd, req.data.ptr, req.data.len); - if (written < 0) return Reply.fail(req.tag, E.IO); - return .{ .tag = req.tag, .written = @intCast(written) }; - } + var buf: [ninep_io.msize]u8 = undefined; // a read never asks for more than a frame + const got = io: { + pardes.turn.yield(); + defer pardes.turn.back(); + const fd = libc.open(path_z, .{ .ACCMODE = if (req.op == .read) .RDONLY else .WRONLY, .NONBLOCK = true, .CLOEXEC = true }); + if (fd < 0) return Reply.fail(req.tag, E.PERM); + defer _ = libc.close(fd); + if (req.op == .setattr) { + if (!req.truncate or req.off != 0 or libc.ftruncate(fd, 0) != 0) return Reply.fail(req.tag, E.INVAL); + var truncated = attr; + truncated.size = 0; + return .{ .tag = req.tag, .attr = truncated }; + } + if (req.off > std.math.maxInt(i64)) return Reply.fail(req.tag, E.INVAL); + if (libc.lseek(fd, @intCast(req.off), libc.SEEK.SET) < 0) return Reply.fail(req.tag, E.IO); + if (req.op == .write) { + const written = libc.write(fd, req.data.ptr, req.data.len); + if (written < 0) return Reply.fail(req.tag, E.IO); + return .{ .tag = req.tag, .written = @intCast(written) }; + } + const got = libc.read(fd, &buf, @min(req.size, buf.len)); + if (got < 0) return Reply.fail(req.tag, E.IO); + break :io @as(usize, @intCast(got)); + }; const out = p.fs.stage(p.gpa); - out.resize(p.gpa, @min(req.size, 65536)) catch return Reply.fail(req.tag, E.NOMEM); - const got = libc.read(fd, out.items.ptr, out.items.len); - if (got < 0) return Reply.fail(req.tag, E.IO); - out.shrinkRetainingCapacity(@intCast(got)); + out.appendSlice(p.gpa, buf[0..got]) catch return Reply.fail(req.tag, E.NOMEM); return .{ .tag = req.tag, .payload = .{ .staged = @intCast(got) } }; }, else => return Reply.fail(req.tag, E.PERM), @@ -378,21 +401,22 @@ test "bounded reads accept exact OS file lengths and empty files" { const empty = try readLimit(p, explicit, 0); defer gpa.free(empty); try std.testing.expectEqual(@as(usize, 0), empty.len); - try std.testing.expectEqual(@as(usize, 0), p.fs.os_paths.items.len); + // The directory's node stays in the table for the listener to collect + // once no connection names it (src/9p_io.zig, collectOs): a client may + // have walked to the same directory meanwhile and hold it by that node. + try std.testing.expectEqual(@as(usize, 1), p.fs.os_paths.items.len); + try std.testing.expectEqualStrings(native, p.fs.os_paths.items[0].path); } test "bounded reads apply the same limits to self bodies and embedded sources" { const gpa = std.testing.allocator; const p = try Pardes.init(gpa, .{ .tty_only = true }); defer p.deinit(); - p.fs.socket_path = "/tmp/pardes-limited-in-process.sock"; - try p.fs.mount(gpa, "own", p.fs.socket_path); var contents: [8209]u8 = @splat('x'); contents[contents.len - 1] = '\n'; for ([_][]const u8{ &contents, "" }) |expected| { const pane = try p.setTestFile(expected); - for ([_][]const u8{ "/virtual", "/n/self", "/n/own" }) |prefix| { - if (!platform_has_fs and std.mem.eql(u8, prefix, "/n/own")) continue; + for ([_][]const u8{ "/virtual", "/n/self" }) |prefix| { var path_buf: [128]u8 = undefined; const path = try std.fmt.bufPrint(&path_buf, "{s}/pane/{d}/body", .{ prefix, pane.serial }); const bytes = try readLimit(p, path, expected.len); @@ -407,8 +431,7 @@ test "bounded reads apply the same limits to self bodies and embedded sources" { } if (limits.embedded_sources) { const source = findEmbeddedSource("src/look.zig", false).?; - for ([_][]const u8{ "/virtual/src/look.zig", "/n/self/src/look.zig", "/n/own/src/look.zig" }) |path| { - if (!platform_has_fs and std.mem.startsWith(u8, path, "/n/own/")) continue; + for ([_][]const u8{ "/virtual/src/look.zig", "/n/self/src/look.zig" }) |path| { const bytes = try readLimit(p, path, source.contents.len); defer gpa.free(bytes); try std.testing.expectEqualStrings(source.contents, bytes); @@ -425,10 +448,8 @@ test "bounded reads count rendered directory paths and unknown self lengths" { const gpa = std.testing.allocator; const p = try Pardes.init(gpa, .{ .tty_only = true }); defer p.deinit(); - p.fs.socket_path = "/tmp/pardes-limited-in-process.sock"; - try p.fs.mount(gpa, "own", p.fs.socket_path); - for ([_][]const u8{ "/n", "/virtual", "/n/self", "/n/own", "/n/own/pane", "/virtual/listeners", "/virtual/README" }) |path| { - if (!platform_has_fs and std.mem.startsWith(u8, path, "/n/own")) continue; + p.fs.socket_path = "/tmp/pardes-limited-in-process.sock"; // so listeners has a line + for ([_][]const u8{ "/n", "/virtual", "/n/self", "/virtual/pane", "/virtual/listeners", "/virtual/README" }) |path| { const expected = try read(p, path); defer gpa.free(expected); try std.testing.expect(expected.len > 0); @@ -573,49 +594,6 @@ test "virtual body writes can read their input from the same pane" { try std.testing.expectEqualStrings("pane contents\n", pane.file.?.content); } -test "same-core mount directories preserve their mount prefix across reads" { - if (!platform_has_fs) return error.SkipZigTest; - const gpa = std.testing.allocator; - const p = try pardes.Pardes.init(gpa, .{ .tty_only = true }); - defer p.deinit(); - const socket_path = "/tmp/pardes-in-process-mount.sock"; - p.fs.socket_path = socket_path; - try p.fs.mount(gpa, "own", socket_path); - const roots = try read(p, "/n/own"); - defer gpa.free(roots); - try std.testing.expect(std.mem.startsWith(u8, roots, "/n/own/README\n/n/own/index\n")); - try std.testing.expect(std.mem.indexOf(u8, roots, "/n/own/pane/\n/n/own/os/\n") != null); - try std.testing.expectError(error.NotADirectory, read(p, "/n/own/index/..")); - try std.testing.expectError(error.FileNotFound, read(p, "/n/own/self/index")); - - var tmp = std.testing.tmpDir(.{}); - defer tmp.cleanup(); - for (0..520) |i| { - var name: [32]u8 = undefined; - try tmp.dir.writeFile(std.testing.io, .{ .sub_path = try std.fmt.bufPrint(&name, "entry-{d:0>4}.txt", .{i}), .data = "child\n" }); - } - var directory_buffer: [4096]u8 = undefined; - const directory = directory_buffer[0..try tmp.dir.realPath(std.testing.io, &directory_buffer)]; - const held_node = try osNode(p, directory); - const path = try std.fmt.allocPrint(gpa, "/n/own/os{s}", .{directory}); - defer gpa.free(path); - const listing = try read(p, path); - defer gpa.free(listing); - var entries = std.mem.tokenizeScalar(u8, listing, '\n'); - var count: usize = 0; - while (entries.next()) |entry| { - try std.testing.expect(std.mem.startsWith(u8, entry, path)); - const child = try read(p, entry); - defer gpa.free(child); - try std.testing.expectEqualStrings("child\n", child); - try std.testing.expectEqual(@as(usize, 1), p.fs.os_paths.items.len); - count += 1; - } - try std.testing.expectEqual(@as(usize, 520), count); - try std.testing.expectEqual(held_node, p.fs.os_paths.items[0].node); - try std.testing.expectEqualStrings(directory, p.fs.os_paths.items[0].path); -} - pub fn isVirtual(path: []const u8) bool { return std.mem.eql(u8, path, "/virtual") or std.mem.startsWith(u8, path, "/virtual/") or std.mem.eql(u8, path, "/n") or std.mem.startsWith(u8, path, "/n/"); @@ -674,9 +652,6 @@ pub fn resolve(p: ?*pardes.Pardes, word: []const u8, cwd: []const u8, out: *[409 } if (std.mem.eql(u8, joined, "/virtual")) return resolveVirtual(p, "/", out); if (std.mem.startsWith(u8, joined, "/virtual/")) return resolveVirtual(p, joined[8..], out); - // This editor's own tree, seen through a mount: answer from the tree - // instead of walking out into the view and back in. - if (ownMountSuffix(joined)) |inner| return resolveVirtual(p, inner, out); if (resolveOs(joined, out)) |found| return found; if (resolveVirtual(p, joined, out)) |found| return found; if (resolveVirtual(p, word, out)) |found| return found; @@ -687,116 +662,15 @@ pub fn resolve(p: ?*pardes.Pardes, word: []const u8, cwd: []const u8, out: *[409 return null; } -/// The name this session is posted under, taken from its socket path. -/// As wide as a name a listener will accept, so there is no session whose -/// name is too long to recognise and therefore too long to protect. -var own_name_buf: [108]u8 = undefined; -var own_name_len: usize = 0; - -pub fn noteOwnSocket(socket_path: []const u8) void { - own_name_len = 0; - const base = std.fs.path.basename(socket_path); - const head = "pardes-9p-"; - const tail = ".sock"; - if (!std.mem.startsWith(u8, base, head) or !std.mem.endsWith(u8, base, tail)) return; - const name = base[head.len .. base.len - tail.len]; - if (name.len == 0 or name.len > own_name_buf.len) return; - @memcpy(own_name_buf[0..name.len], name); - own_name_len = name.len; -} - -/// What this path names inside this editor's OWN 9P tree, if it does: -/// `/mnt/9p/pardes/<me>/pane/3/body` -> `pane/3/body`, and the mount root -/// itself -> `/`. -/// -/// A mount is a VIEW of a tree this editor already holds. Going out through -/// the view to reach it deadlocks the session outright -- the realpath, the -/// stat and the read all leave through the mount and come back as 9P requests -/// only this editor's loop can answer, while that loop is blocked making them, -/// and the filesystem then stops answering anybody. So the view is recognised -/// by NAME, before any syscall (the syscall is the thing that never returns), -/// and the request is served from the tree directly. Same answer, no round -/// trip, and a session can drive itself through its own 9P namespace. -pub fn ownMountSuffix(path: []const u8) ?[]const u8 { - if (own_name_len == 0 or path.len == 0 or path[0] != '/') return null; - // Whole components, and the registry's own layout: a session is posted at - // `<runtime>/9p/pardes/<name>` and mounts group under `/mnt/9p/pardes/`. - // Matching a bare `/pardes/<name>` anywhere in the string would claim - // `~/src/pardes/<name>/README` -- an ordinary directory that happens to - // read like a mount -- and serve the tree's README over the real file. - var buf: [own_name_buf.len + 16]u8 = undefined; - const stem = std.fmt.bufPrint(&buf, "/9p/pardes/{s}", .{own_name_buf[0..own_name_len]}) catch return null; - var at: usize = 0; - while (std.mem.indexOfPos(u8, path, at, stem)) |found| : (at = found + 1) { - const rest = path[found + stem.len ..]; - if (rest.len != 0 and rest[0] != '/') continue; // a longer name that merely starts the same - if (rest.len <= 1) return "/"; - return rest[1..]; - } - return null; -} - -pub fn isOwnMount(path: []const u8) bool { - return ownMountSuffix(path) != null; -} - -test "a path inside this session's own mount is answered from the tree, not through the mount" { - noteOwnSocket("/run/user/1000/pardes-9p-demo.sock"); - defer own_name_len = 0; - // What the path names inside the tree, with the mount prefix taken off. - try std.testing.expectEqualStrings("index", ownMountSuffix("/mnt/9p/pardes/demo/index").?); - try std.testing.expectEqualStrings("pane/new", ownMountSuffix("/mnt/9p/pardes/demo/pane/new").?); - try std.testing.expectEqualStrings("/", ownMountSuffix("/mnt/9p/pardes/demo").?); - try std.testing.expectEqualStrings("/", ownMountSuffix("/mnt/9p/pardes/demo/").?); - // The registry posts the session there too, so that spelling counts. - try std.testing.expectEqualStrings("index", ownMountSuffix("/run/user/1000/9p/pardes/demo/index").?); - // Another session's mount is somebody else's to answer, and a name that - // merely STARTS with ours is a different name. - try std.testing.expect(ownMountSuffix("/mnt/9p/pardes/other/index") == null); - try std.testing.expect(ownMountSuffix("/mnt/9p/pardes/demo2/index") == null); - // An ordinary directory that reads like a mount is an ordinary directory. - // Serving the tree here would hand back the tree's README for the file on - // disk, and -- before `write` learned the same trick -- write the tree's - // bytes over it. - try std.testing.expect(ownMountSuffix("/home/goblin/src/pardes/demo/README") == null); - try std.testing.expect(ownMountSuffix("/home/goblin/src/pardes/demo.zig") == null); - // Relative paths never name a mount: they are resolved against a cwd first. - try std.testing.expect(ownMountSuffix("mnt/9p/pardes/demo/index") == null); - - // The syscall is the thing that never returns, so it is never made. - var out: [4096]u8 = undefined; - try std.testing.expect(resolveOs("/mnt/9p/pardes/demo/index", &out) == null); - - // A session with no listener has no name, so nothing is redirected. - noteOwnSocket("/tmp/not-a-pardes-socket"); - try std.testing.expect(ownMountSuffix("/mnt/9p/pardes/demo/index") == null); -} - -test "reading this session's own mount returns what the tree holds" { - const p = try pardes.Pardes.init(std.testing.allocator, .{ .tty_only = true, .cols = 80, .rows = 24 }); - defer p.deinit(); - noteOwnSocket("/run/user/1000/pardes-9p-selfread.sock"); - defer own_name_len = 0; - - const direct = try read(p, "/n/self/index"); - defer p.gpa.free(direct); - const mounted = try read(p, "/mnt/9p/pardes/selfread/index"); - defer p.gpa.free(mounted); - try std.testing.expectEqualStrings(direct, mounted); - try std.testing.expect(direct.len > 0); - - // A write goes to the same tree the read came from, not out through the - // mount: `/index` is 0400 there, and an OS write would instead try to - // create a file under a directory that does not exist. - try std.testing.expectError(error.ReadOnlyFilesystem, write(p, "/mnt/9p/pardes/selfread/index", "nope\n")); - try std.testing.expectError(error.ReadOnlyFilesystem, write(p, "/n/self/index", "nope\n")); -} - +/// Where a path really is on the host, and whether it is a directory. The +/// path may lie in a mount this editor serves, so the turn is given up for +/// the syscalls: another thread answers them. pub fn resolveOs(path: []const u8, out: *[4096]u8) ?Resolved { if (comptime !platform_has_fs) return null; - if (isOwnMount(path)) return null; var z: [4096]u8 = undefined; const path_z = std.fmt.bufPrintSentinel(&z, "{s}", .{path}, 0) catch return null; + pardes.turn.yield(); + defer pardes.turn.back(); const resolved = realpath(path_z, out) orelse return null; return .{ .path = std.mem.span(resolved), .dir = isDir(resolved) }; } @@ -818,13 +692,6 @@ pub fn read(p: *pardes.Pardes, path: []const u8) ![]u8 { } pub fn readLimit(p: *pardes.Pardes, path: []const u8, max_bytes: usize) ![]u8 { - // The same view, reached by a reader that never went through `resolve`. - // The rewritten path cannot contain the mount stem, so this recurses once. - if (ownMountSuffix(path)) |inner| { - var self_buf: [4096]u8 = undefined; - const internal = std.fmt.bufPrint(&self_buf, "/n/self/{s}", .{inner}) catch return error.FileTooLarge; - return readLimit(p, internal, max_bytes); - } const limit = @min(max_bytes, limits.max_file_bytes); if (std.mem.eql(u8, std.mem.trimEnd(u8, path, "/"), "/n")) { var out: std.Io.Writer.Allocating = .init(p.gpa); @@ -847,11 +714,8 @@ pub fn readLimit(p: *pardes.Pardes, path: []const u8, max_bytes: usize) ![]u8 { var native_buf: [4096]u8 = undefined; const native = resolveOs(remote_path, &native_buf) orelse return readFileLimit(p.gpa, path, limit); if (!native.dir) return readFileLimit(p.gpa, path, limit); - const saved_paths = p.fs.os_paths.items.len; - defer { - for (p.fs.os_paths.items[saved_paths..]) |temporary| p.gpa.free(temporary.path); - p.fs.os_paths.shrinkRetainingCapacity(saved_paths); - } + // The node stays in the table until no connection names it either + // (src/9p_io.zig, collectOs): a client may have walked to it too. const node = try osNode(p, native.path); return readNode(p, node, path, limit); } @@ -861,25 +725,10 @@ pub fn readLimit(p: *pardes.Pardes, path: []const u8, max_bytes: usize) ![]u8 { } for (p.fs.mounts.items) |mount| { if (!std.mem.eql(u8, name, mount.name)) continue; - if (ninep_io.Client.sameSession(mount.dial, p.fs.socket_path, p.fs.tcp_address, p.fs.quic_address)) { - const saved_paths = p.fs.os_paths.items.len; - defer { - for (p.fs.os_paths.items[saved_paths..]) |temporary| p.gpa.free(temporary.path); - p.fs.os_paths.shrinkRetainingCapacity(saved_paths); - } - var node = tree.root; - var directory = true; - var parts = std.mem.tokenizeScalar(u8, remote_path, '/'); - while (parts.next()) |part| { - if (!directory) return error.NotADirectory; - if (std.mem.eql(u8, part, ".")) continue; - const reply = tree.handle(p, .{ .tag = 0, .op = .lookup, .node = node, .data = part }); - if (reply.status != .ok) return error.FileNotFound; - node = reply.attr.node; - directory = reply.attr.dir; - } - return readNode(p, node, path, limit); - } + // A peer answers on its own time, and the peer may be this very + // editor: the turn goes out with the request. + pardes.turn.yield(); + defer pardes.turn.back(); return ninep_io.Client.readLimit(p.gpa, mount.dial, remote_path, path, limit); } return error.FileNotFound; @@ -940,68 +789,64 @@ fn readNode(p: *pardes.Pardes, initial_node: u64, path: []const u8, limit: usize } } +/// The editor's, between steps: a save or a dump the shell is performing, +/// never a step of the core. pub fn write(p: *pardes.Pardes, path: []const u8, bytes: []const u8) !void { - // The same view `readLimit` answers from the tree, so a write lands where - // the matching read came from -- and, just as importantly, never becomes - // an open(2) through a mount this loop is the one that answers. - if (ownMountSuffix(path)) |inner| { - var self_buf: [4096]u8 = undefined; - const internal = std.fmt.bufPrint(&self_buf, "/n/self/{s}", .{inner}) catch return error.PathTooLong; - return write(p, internal, bytes); - } if (std.mem.eql(u8, std.mem.trimEnd(u8, path, "/"), "/n")) return error.IsDirectory; - var self_path: ?[]const u8 = null; - if (std.mem.eql(u8, path, "/virtual")) self_path = ""; - if (std.mem.startsWith(u8, path, "/virtual/")) self_path = path[9..]; - if (std.mem.startsWith(u8, path, "/n/")) { - const explicit = path[3..]; - const cut = std.mem.indexOfScalar(u8, explicit, '/') orelse explicit.len; - const name = explicit[0..cut]; - const remote_path = if (cut < explicit.len) explicit[cut..] else "/"; - if (std.mem.eql(u8, name, "os")) return writeFile(remote_path, bytes); - if (std.mem.eql(u8, name, "self")) { - self_path = remote_path; - } else { - for (p.fs.mounts.items) |mount| { - if (!std.mem.eql(u8, name, mount.name)) continue; - if (ninep_io.Client.sameSession(mount.dial, p.fs.socket_path, p.fs.tcp_address, p.fs.quic_address)) { - var buf: [4096]u8 = undefined; - const local = try std.fmt.bufPrint(&buf, "/n{s}", .{remote_path}); - return write(p, local, bytes); - } - return ninep_io.Client.write(p.gpa, mount.dial, remote_path, bytes); - } - return error.FileNotFound; - } + if (std.mem.eql(u8, path, "/virtual")) return writeSelf(p, "", bytes); + if (std.mem.startsWith(u8, path, "/virtual/")) return writeSelf(p, path[9..], bytes); + if (!std.mem.startsWith(u8, path, "/n/")) return writeOut(p, null, path, bytes); + const explicit = path[3..]; + const cut = std.mem.indexOfScalar(u8, explicit, '/') orelse explicit.len; + const name = explicit[0..cut]; + const remote_path = if (cut < explicit.len) explicit[cut..] else "/"; + if (std.mem.eql(u8, name, "os")) return writeOut(p, null, remote_path, bytes); + if (std.mem.eql(u8, name, "self")) return writeSelf(p, remote_path, bytes); + for (p.fs.mounts.items) |mount| { + if (std.mem.eql(u8, name, mount.name)) return writeOut(p, mount.dial, remote_path, bytes); } - if (self_path) |name| { - var node = tree.resolveSelf(p, name) orelse return error.FileNotFound; - const attributes = tree.handle(p, .{ .tag = 0, .op = .getattr, .node = node }); - if (attributes.status != .ok) return error.FileNotFound; - if (attributes.attr.dir) return error.IsDirectory; - if (attributes.attr.mode & 0o200 == 0) return error.ReadOnlyFilesystem; - const opened = tree.handle(p, .{ .tag = 0, .op = .open, .node = node }); - if (opened.status != .ok) return error.OpenFailed; - if (opened.attr.node != 0) node = opened.attr.node; - defer _ = tree.handle(p, .{ .tag = 0, .op = .release, .node = node, .handle = opened.handle }); - var preserved: ?[]u8 = null; - defer if (preserved) |copy| p.gpa.free(copy); - const target = tree.Node.target(node); - if (target != null and target.? == .pane and target.?.pane.file == .body) { - preserved = try p.gpa.dupe(u8, bytes); - const trunc = tree.handle(p, .{ .tag = 0, .op = .setattr, .node = node, .truncate = true }); - if (trunc.status != .ok) return error.WriteFailed; - } - const contents = preserved orelse bytes; - var off: usize = 0; - while (off < contents.len) { - const reply = tree.handle(p, .{ .tag = 0, .op = .write, .node = node, .off = off, .data = contents[off..] }); - if (reply.status != .ok or reply.written == 0) return error.WriteFailed; - off += reply.written; - } - return; + return error.FileNotFound; +} + +/// The host, or the peer at `dial`, takes the bytes with the turn given up +/// entirely -- the peer may be this very editor, and a write it serves may +/// change any pane -- so they are copied first: `bytes` is usually a pane's +/// own text. +fn writeOut(p: *pardes.Pardes, dial: ?[]const u8, path: []const u8, bytes: []const u8) !void { + const copy = try p.gpa.dupe(u8, bytes); + defer p.gpa.free(copy); + pardes.turn.rest(); + defer pardes.turn.wake(); + if (dial) |d| return ninep_io.Client.write(p.gpa, d, path, copy); + return writeFile(path, copy); +} + +/// This editor's own tree, written in place. +fn writeSelf(p: *pardes.Pardes, name: []const u8, bytes: []const u8) !void { + var node = tree.resolveSelf(p, name) orelse return error.FileNotFound; + const attributes = tree.handle(p, .{ .tag = 0, .op = .getattr, .node = node }); + if (attributes.status != .ok) return error.FileNotFound; + if (attributes.attr.dir) return error.IsDirectory; + if (attributes.attr.mode & 0o200 == 0) return error.ReadOnlyFilesystem; + const opened = tree.handle(p, .{ .tag = 0, .op = .open, .node = node }); + if (opened.status != .ok) return error.OpenFailed; + if (opened.attr.node != 0) node = opened.attr.node; + defer _ = tree.handle(p, .{ .tag = 0, .op = .release, .node = node, .handle = opened.handle }); + var preserved: ?[]u8 = null; + defer if (preserved) |copy| p.gpa.free(copy); + const target = tree.Node.target(node); + if (target != null and target.? == .pane and target.?.pane.file == .body) { + preserved = try p.gpa.dupe(u8, bytes); + const trunc = tree.handle(p, .{ .tag = 0, .op = .setattr, .node = node, .truncate = true }); + if (trunc.status != .ok) return error.WriteFailed; + } + const contents = preserved orelse bytes; + var off: usize = 0; + while (off < contents.len) { + const reply = tree.handle(p, .{ .tag = 0, .op = .write, .node = node, .off = off, .data = contents[off..] }); + if (reply.status != .ok or reply.written == 0) return error.WriteFailed; + off += reply.written; } - return writeFile(path, bytes); } fn resolveEmbedded(word: []const u8, cwd: []const u8, scratch: *[4096]u8) ?Source { @@ -1069,6 +914,8 @@ pub fn find(arena: std.mem.Allocator, dir: []const u8, pat: []const u8, out: []u var hits: [find_max_hits][]const u8 = undefined; var hits_len: usize = 0; if (platform_has_fs) { + pardes.turn.yield(); + defer pardes.turn.back(); const io = std.Io.Threaded.global_single_threaded.io(); var root = try std.Io.Dir.cwd().openDir(io, dir, .{ .iterate = true }); defer root.close(io); @@ -1152,6 +999,9 @@ pub fn grep(arena: std.mem.Allocator, gpa: std.mem.Allocator, dir: []const u8, b const home = std.mem.trimEnd(u8, base, "/"); const files = try arena.alloc([]const u8, grep_max_files); var files_len: usize = 0; + // A walk and thousands of reads: the turn goes out for all of it. + pardes.turn.yield(); + defer pardes.turn.back(); { const io = std.Io.Threaded.global_single_threaded.io(); var root = try std.Io.Dir.cwd().openDir(io, dir, .{ .iterate = true }); @@ -1311,10 +1161,6 @@ test "Restore prefers default directory then falls back to original path" { fn readFileLimit(gpa: std.mem.Allocator, path: []const u8, limit: usize) ![]u8 { if (std.mem.indexOfScalar(u8, path, 0) != null) return error.OpenFailed; - // Same reason as `writeFile`: no core to answer from here, and the open is - // the call that never returns. Callers holding one use `readLimit`, which - // serves the tree instead. - if (isOwnMount(path)) return error.FileNotFound; if (!platform_has_fs or std.mem.startsWith(u8, path, "/virtual/")) { const archive_path = if (std.mem.startsWith(u8, path, "/virtual/")) path[9..] else path; var normalized_buf: [4096]u8 = undefined; @@ -1326,6 +1172,10 @@ fn readFileLimit(gpa: std.mem.Allocator, path: []const u8, limit: usize) ![]u8 { var pathbuf: [4096]u8 = undefined; const native = localPath(path) orelse return error.OpenFailed; const path_z = std.fmt.bufPrintSentinel(&pathbuf, "{s}", .{native}, 0) catch return error.PathTooLong; + // The path may be a mount this editor serves: the turn goes out with the + // syscalls, and the bytes come back into memory of the caller's own. + pardes.turn.yield(); + defer pardes.turn.back(); const fd = libc.open(path_z, .{ .ACCMODE = .RDONLY, .NONBLOCK = true }); if (fd < 0) return switch (libc.errno(fd)) { .ACCES, .PERM => error.PermissionDenied, diff --git a/src/gui/gui.zig b/src/gui/gui.zig index 3c2fd693..8da6f933 100644 --- a/src/gui/gui.zig +++ b/src/gui/gui.zig @@ -2422,6 +2422,7 @@ fn localSession( defer pardes.lsp.setStatusSink(null, null); while (!core.quit) { + pardes.turn.restoreSettled(); try core.pump(host); if (core.quit) break; // a session that ended does not restore into one if (core.takeAttach()) |req| { @@ -2465,7 +2466,7 @@ fn localSession( g.prepared_images.clearRetainingCapacity(); nc.native_images = true; nc.host = host; - if (fs) |f| f.reset(core); + if (fs) |f| f.reset(nc); core.deinit(); core = nc; shell.core = nc; @@ -2758,7 +2759,10 @@ const StdinFeed = struct { } var fds = [_]posix.pollfd{.{ .fd = 0, .events = posix.POLL.IN, .revents = 0 }}; - _ = posix.poll(&fds, 16) catch return out; + pardes.turn.rest(); + const polled = posix.poll(&fds, 16); + pardes.turn.wake(); + _ = polled catch return out; if ((fds[0].revents & posix.POLL.IN) == 0) { if ((fds[0].revents & (posix.POLL.HUP | posix.POLL.ERR)) != 0) out.eof = true; return out; @@ -3762,7 +3766,11 @@ fn waitInput(ctx: ?*anyopaque, timeout_ms: u32) void { if (s.gui) |g| { var sev = std.mem.zeroes(c.SDL_Event); const ms: c_int = if (timeout_ms != 0) @intCast(timeout_ms) else 16; - if (c.SDL_WaitEventTimeout(&sev, ms)) { + // The wait is the 9P connections' turn with the core. + pardes.turn.rest(); + const got = c.SDL_WaitEventTimeout(&sev, ms); + pardes.turn.wake(); + if (got) { dispatch(g, &in, &sev); while (c.SDL_PollEvent(&sev)) dispatch(g, &in, &sev); } @@ -3785,7 +3793,7 @@ fn pollFrame(ctx: ?*anyopaque) void { const s = shellOf(ctx); s.reconcilePtys(); const core = s.core; - if (s.fs) |f| if (f.tick(core).pending) s.queue.push(.fs_ready); + if (s.fs) |f| if (f.tick().pending) s.queue.push(.fs_ready); const g = s.gui orelse { if (core.settings.window_opacity_pending) { var unchanged_opacity: u8 = 100; @@ -3923,7 +3931,7 @@ fn gridPollFrame(ctx: ?*anyopaque) void { const s = shellOf(ctx); s.reconcilePtys(); if (s.fs) |f| { - if (f.tick(s.core).count != 0) s.saw_event = true; + if (f.tick().count != 0) s.saw_event = true; } pollCwds(s.core, s.ptys); if (s.core.animationActive()) { @@ -5077,6 +5085,13 @@ fn renderFrame( ); } } + // Body layers first: a notice chip over a tree-sitter context band + // is a tag layer, and instances paint in emission order. + for (surface.bodyLayers()) |*layer| { + if (layer.rows == 0) continue; + const batch_index = paintBatchAt(&paint_plan, layer.viewport.x, layer.viewport.y); + emitBodyLayer(g, instances, &cell_next[batch_index], layer, win_w, win_h, paint_plan.batches[batch_index].track, page, paint_plan.len == 1, false, topbar_pane_border_rgb); + } for (surface.tagLayers()) |*layer| { if (layer.cols == 0) continue; const batch_index = paintBatchAt(&paint_plan, layer.viewport.x, layer.viewport.y); @@ -5092,11 +5107,6 @@ fn renderFrame( emitTagLayer(g, instances, &cell_next[destination], layer, win_w, win_h, if (under) null else track, page, false, track.effect == .dissolve); } } - for (surface.bodyLayers()) |*layer| { - if (layer.rows == 0) continue; - const batch_index = paintBatchAt(&paint_plan, layer.viewport.x, layer.viewport.y); - emitBodyLayer(g, instances, &cell_next[batch_index], layer, win_w, win_h, paint_plan.batches[batch_index].track, page, paint_plan.len == 1, false, topbar_pane_border_rgb); - } for (paint_plan.batches[1..paint_plan.len], 1..) |batch, batch_index| { const track = batch.track.?; const under = track.effect == .vertical and track.phase == .opening; diff --git a/src/host_io.zig b/src/host_io.zig index a7babc74..72e149c9 100644 --- a/src/host_io.zig +++ b/src/host_io.zig @@ -955,6 +955,9 @@ pub fn forkShell( var cwd_buf: [4096]u8 = undefined; const cwd_z: ?[:0]const u8 = if (native_cwd.len == 0) null else dir: { const path = std.fmt.bufPrintSentinel(&cwd_buf, "{s}", .{native_cwd}, 0) catch return error.NameTooLong; + // A shell's directory may be inside a mount this editor serves. + pardes.turn.yield(); + defer pardes.turn.back(); const stat = try std.Io.Dir.cwd().statFile(std.Io.Threaded.global_single_threaded.io(), path, .{}); if (stat.kind != .directory) return error.NotDir; break :dir path; diff --git a/src/macos.zig b/src/macos.zig index d3b33c55..a7ee3164 100644 --- a/src/macos.zig +++ b/src/macos.zig @@ -773,6 +773,10 @@ fn initCore(runtime: ?*const Runtime, cols_arg: u16, rows_arg: u16) !void { for (&st.ptys, 0..) |*slot, id| if (slot.*) |*pt| startReader(st, pt, @intCast(id)); pardes.lsp.setStatusSink(st, lspStatusSink); + // From here the core is touched only inside the entry points above, + // each of which takes the turn for its own duration; between them the + // 9P connections have it. + pardes.turn.rest(); } fn wakeNinep(ctx: ?*anyopaque) void { @@ -781,6 +785,7 @@ fn wakeNinep(ctx: ?*anyopaque) void { } export fn pardes_deinit() void { + pardes.turn.wake(); const st = &(state orelse return); if (st.ninep) |listener| { listener.deinit(st.gpa); @@ -823,6 +828,8 @@ export fn pardes_deinit() void { } export fn pardes_should_quit() bool { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return true); return st.core.quit; } @@ -890,11 +897,15 @@ test "the animation clock spends real time, not callbacks" { } export fn pardes_animating() bool { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return false); return st.core.animationActive() or st.rotate_coasting or st.scroll_coasting; } export fn pardes_theme_bg() u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const th = if (state) |*st| st.core.theme() else &pardes.themes[0]; const bg = th.bg orelse return color_default; return @as(u32, bg[0]) << 16 | @as(u32, bg[1]) << 8 | bg[2]; @@ -913,6 +924,8 @@ export fn pardes_theme_bg() u32 { /// own (pardes_theme_bg == PARDES_COLOR_DEFAULT) drops the ground entirely, /// whatever percentage is set here. export fn pardes_window_opacity() u8 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 100); st.core.settings.window_opacity_pending = false; return st.core.settings.window_opacity; @@ -920,6 +933,8 @@ export fn pardes_window_opacity() u8 { /// Public AppKit backdrop strength, independent of background paint opacity. export fn pardes_window_blur() u8 { + pardes.turn.wake(); + defer pardes.turn.rest(); return if (state) |st| st.core.settings.window_blur else 0; } @@ -928,10 +943,14 @@ fn taglineFontPercent(core: ?*const pardes.Pardes) u8 { } export fn pardes_gui_tagline_font_percent() u8 { + pardes.turn.wake(); + defer pardes.turn.rest(); return taglineFontPercent(if (state) |*st| st.core else null); } export fn pardes_tagline_band_offset(row: u16, canvas_h: f32, cell_h: u32, tagline_h: u32) u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); return pardes.taglineBandOffset(row, canvas_h, cell_h, tagline_h, if (state) |*st| st.core.settings.workspace_tag else true); } @@ -940,16 +959,22 @@ export fn pardes_topbar_pane_border_px(cell_h: u32, tagline_h: u32) u32 { } export fn pardes_tagline_origin_col(col: u16, row: u16) f32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return @floatFromInt(col)); return pardes.taglineOriginColForFrame(st.core, col, row); } export fn pardes_grid_col_at(x: f32, row: u16, body_w: f32, tagline_w: f32) u16 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return pardes.gridColAt(null, x, row, body_w, tagline_w)); return pardes.gridColAt(st.core, x, row, body_w, tagline_w); } export fn pardes_topbar_pane_border_rgb() u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const rgb = pardes.config.gui_topbar_pane_border_rgb orelse fromTheme: { const st = state orelse return color_default; break :fromTheme st.core.chromeTheme().border; @@ -958,11 +983,15 @@ export fn pardes_topbar_pane_border_rgb() u32 { } export fn pardes_tag_active_bg() u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = state orelse return color_default; const rgb = st.core.chromeTheme().tag_active_bg; return @as(u32, rgb[0]) << 16 | @as(u32, rgb[1]) << 8 | rgb[2]; } export fn pardes_tagline_bg() u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = state orelse return color_default; const rgb = st.core.chromeTheme().tag_bg; return @as(u32, rgb[0]) << 16 | @as(u32, rgb[1]) << 8 | rgb[2]; @@ -1008,6 +1037,8 @@ test "the fallback preference order crosses the ABI intact and ends at the bound } export fn pardes_scene() Scene { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return .{}); const seconds = @as(f64, @floatFromInt(st.scene_ns)) / @as(f64, std.time.ns_per_s); return .{ @@ -1018,6 +1049,8 @@ export fn pardes_scene() Scene { } export fn pardes_postprocessor_unavailable() void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); st.core.disableSceneEffects(); st.core.settings.panel_transition = .off; @@ -1026,6 +1059,8 @@ export fn pardes_postprocessor_unavailable() void { } export fn pardes_panel_animation_failed() void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); st.core.abandonPanelAnimations(); } @@ -1210,6 +1245,8 @@ fn reloadWatchedFile(st: *State, pane: u8, announce: bool) bool { } export fn pardes_watch_changed(pane: u8, generation: u32) void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); if (pane >= watch_slot_count) return; if (st.file_watches.notify(pane, generation)) wake(st); @@ -1265,11 +1302,13 @@ fn drainInbox(st: *State) bool { } export fn pardes_tick() bool { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return false); st.inbox.wake_pending.store(false, .release); var did = drainInbox(st); if (st.ninep) |listener| { - const drained = listener.tick(st.core); + const drained = listener.tick(); did = did or drained.count != 0; if (drained.pending) wake(st); } @@ -1277,7 +1316,9 @@ export fn pardes_tick() bool { did = true; st.core.perform(effect); } + pardes.turn.settle(); if (restoreCore(st)) did = true; + pardes.turn.restoreSettled(); return did; } @@ -1309,7 +1350,7 @@ fn restoreCore(st: *State) bool { const generation = st.file_watches.stop(st.gpa, id); hostWatchFile(id, generation, null, 0); }; - if (st.ninep) |listener| listener.reset(st.core); + if (st.ninep) |listener| listener.reset(replacement); replacement.host = hostFor(st); st.core.deinit(); st.core = replacement; @@ -1485,6 +1526,8 @@ test "mac Restore cancels a PTY reader waiting for output capacity" { } export fn pardes_animation_tick() bool { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return false); const now: u64 = @intCast(@max(0, monotonicNs())); @@ -1546,6 +1589,8 @@ export fn pardes_animation_tick() bool { } export fn pardes_key(cp_arg: u32, text_ptr: ?[*]const u8, len: usize, mods: u32) void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); if (cp_arg > std.math.maxInt(u21)) return; const text: []const u8 = if (text_ptr) |p| p[0..len] else ""; @@ -1559,16 +1604,22 @@ export fn pardes_key(cp_arg: u32, text_ptr: ?[*]const u8, len: usize, mods: u32) } export fn pardes_paste(text_ptr: ?[*]const u8, len: usize) void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); const text: []const u8 = if (text_ptr) |p| p[0..len] else ""; st.core.update(.{ .paste = text }); } export fn pardes_mouse(button_arg: c_int, kind_arg: c_int, col: u16, row: u16, mods: u32) void { + pardes.turn.wake(); + defer pardes.turn.rest(); mouseWithHit(button_arg, kind_arg, col, row, mods, null, null); } export fn pardes_mouse_pixel(button_arg: c_int, kind_arg: c_int, col: u16, row: u16, mods: u32, x: f32, y: f32, bw: f32, bh: f32, tw: f32, th: f32) void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); const source = st.core.presentation.pointerFractional(st.core.screen_w, st.core.screen_h, x / @max(1, bw), y / @max(1, bh)); var hit: ?pardes.Mouse.BodyHit = null; @@ -1620,11 +1671,15 @@ fn mouseWithHit(button_arg: c_int, kind_arg: c_int, col: u16, row: u16, mods: u3 } export fn pardes_pointer_leave() void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); st.core.update(.pointer_leave); } export fn pardes_scroll(delta_rows: f32, delta_cols: f32, col: u16, row: u16) void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); st.scroll_col = col; st.scroll_row = row; @@ -1636,6 +1691,8 @@ export fn pardes_scroll(delta_rows: f32, delta_cols: f32, col: u16, row: u16) vo /// spinning dial. Also drops whatever the last gesture left banked, so the /// first millimetre of a new swipe cannot inherit a nearly-complete row. export fn pardes_scroll_begin() void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); st.scroll_lag = 0; st.scroll_lag_x = 0; @@ -1655,6 +1712,8 @@ export fn pardes_scroll_begin() void { /// momentum simply takes the gesture over, and the tail of that momentum is /// below the fling floor, so its end arms nothing. export fn pardes_scroll_end() void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); const last = st.scroll_last_ns; st.scroll_last_ns = 0; @@ -1671,6 +1730,8 @@ export fn pardes_scroll_end() void { } export fn pardes_rotate(degrees: f32) void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); if (degrees == 0) { st.rotate_lag = 0; @@ -1684,6 +1745,8 @@ export fn pardes_rotate(degrees: f32) void { } export fn pardes_rotate_end() void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); const last = st.rotate_last_ns; st.rotate_last_ns = 0; @@ -1726,6 +1789,8 @@ fn spendRotation(st: *State, degrees: f32) void { } export fn pardes_command(text_ptr: ?[*]const u8, len: usize) void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); const text: []const u8 = if (text_ptr) |p| p[0..len] else ""; if (text.len == 0) return; @@ -1733,6 +1798,8 @@ export fn pardes_command(text_ptr: ?[*]const u8, len: usize) void { } export fn pardes_resize(cols_arg: u16, rows_arg: u16, cell_w: u16, cell_h: u16) void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); const cols = @max(1, cols_arg); const rows = @max(1, rows_arg); @@ -1747,6 +1814,8 @@ export fn pardes_resize(cols_arg: u16, rows_arg: u16, cell_w: u16, cell_h: u16) } export fn pardes_frame() u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 0); const tz = tracy.zone(@src(), "pardes_frame"); defer tz.end(); @@ -1921,26 +1990,36 @@ fn collectImages(st: *State, surface: *const pardes.Surface) void { } export fn pardes_frame_images() u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 0); return @intCast(st.images_len); } export fn pardes_frame_image_list() ?[*]const Image { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return null); return if (st.images_len == 0) null else st.images.ptr; } export fn pardes_frame_panel_tracks() u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 0); return @intCast(st.panel_tracks_len); } export fn pardes_frame_panel_track_list() ?[*]const PanelTrack { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return null); return if (st.panel_tracks_len == 0) null else st.panel_tracks[0..].ptr; } export fn pardes_frame_presented(animated_panels: bool) bool { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return false); const was_animating = st.core.animationActive(); if (animated_panels) @@ -1951,11 +2030,15 @@ export fn pardes_frame_presented(animated_panels: bool) bool { } export fn pardes_frame_cells() ?[*]const Cell { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return null); return if (st.frame_len == 0) null else st.cells.ptr; } export fn pardes_frame_previous_cells() ?[*]const Cell { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return null); return if (st.panel_diff_len != st.frame_len or st.panel_diff_len == 0) null @@ -1964,6 +2047,8 @@ export fn pardes_frame_previous_cells() ?[*]const Cell { } export fn pardes_frame_changed_cells() ?[*]const u8 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return null); return if (st.panel_diff_len != st.frame_len or st.panel_diff_len == 0) null @@ -1972,16 +2057,22 @@ export fn pardes_frame_changed_cells() ?[*]const u8 { } export fn pardes_frame_cols() u16 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 0); return st.frame_cols; } export fn pardes_frame_rows() u16 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 0); return st.frame_rows; } export fn pardes_cursor_x() i32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return -1); if (tagCursorCovered(st.core)) return -1; return if (st.core.surface.cursor) |c| c.x else -1; @@ -2000,6 +2091,8 @@ export fn pardes_tag_layer_limit() u32 { } export fn pardes_tag_layer_value(index: u32, field: u32) u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 0); if (index >= pardes.MAX_TAG_LAYERS) return 0; const layer = &st.core.surface.tag_layers[index]; @@ -2025,6 +2118,8 @@ export fn pardes_tag_layer_value(index: u32, field: u32) u32 { } export fn pardes_tag_layer_cells(index: u32) ?[*]const Cell { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return null); if (index >= pardes.MAX_TAG_LAYERS) return null; const layer = &st.core.surface.tag_layers[index]; @@ -2036,6 +2131,8 @@ export fn pardes_tag_layer_cells(index: u32) ?[*]const Cell { } export fn pardes_body_layer_value(index: u32, field: u32) u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 0); if (index >= pardes.MAX_PANES) return 0; const layer = &st.core.surface.body_layers[index]; @@ -2056,6 +2153,8 @@ export fn pardes_body_layer_value(index: u32, field: u32) u32 { } export fn pardes_body_layer_cells(index: u32) ?[*]const Cell { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return null); if (index >= pardes.MAX_PANES) return null; const layer = &st.core.surface.body_layers[index]; @@ -2071,27 +2170,37 @@ export fn pardes_body_layer_limit() u32 { } export fn pardes_row_metrics(body_w: u16, body_h: u16, tagline_w: u16, tagline_h: u16) void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); st.core.update(.{ .resize = .{ .cols = st.core.screen_w, .rows = st.core.screen_h, .cell_pixels = st.core.cell_pixels, .row_metrics = .{ .body_w = @max(1, body_w), .body_h = @max(1, body_h), .tagline_w = @max(1, tagline_w), .tagline_h = @max(1, tagline_h) } } }); } export fn pardes_pointer_shape() u32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 0); return @intFromEnum(st.core.surface.pointer_shape); } export fn pardes_cursor_y() i32 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return -1); if (tagCursorCovered(st.core)) return -1; return if (st.core.surface.cursor) |c| c.y else -1; } export fn pardes_cursor_bar() bool { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return false); return if (st.core.surface.cursor) |c| c.bar else false; } export fn pardes_take_haptic() c_int { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return 0); return switch (st.core.takeHaptic()) { .none => 0, @@ -2103,6 +2212,8 @@ export fn pardes_take_haptic() c_int { var font_path_z: [4096:0]u8 = undefined; export fn pardes_font_take(size_hundredths: *u16) ?[*:0]const u8 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return null); if (comptime !pardes.font_picker) return null; const want = st.core.takeFontRequest() orelse return null; @@ -2118,6 +2229,8 @@ export fn pardes_font_observe( len: usize, point_hundredths: u16, ) bool { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return false); const ptr = effective_name orelse return false; if (len == 0 or len > 255 or point_hundredths == 0) return false; @@ -2129,6 +2242,8 @@ export fn pardes_font_ack( len: usize, point_hundredths: u16, ) bool { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return false); const ptr = effective_name orelse return false; if (len == 0 or len > 255 or point_hundredths == 0) return false; @@ -2136,6 +2251,8 @@ export fn pardes_font_ack( } export fn pardes_font_reject() void { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return); st.core.rejectFont(); } @@ -2143,6 +2260,8 @@ export fn pardes_font_reject() void { var active_path_z: [4096:0]u8 = undefined; export fn pardes_active_path() ?[*:0]const u8 { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return null); const path = activeFilePath(st) orelse return null; if (path.len == 0 or path.len >= active_path_z.len) return null; @@ -2152,6 +2271,8 @@ export fn pardes_active_path() ?[*:0]const u8 { } export fn pardes_active_dirty() bool { + pardes.turn.wake(); + defer pardes.turn.rest(); const st = &(state orelse return false); const pane = st.core.panes[st.core.active] orelse return false; const f = if (pane.file) |*x| x else return false; diff --git a/src/memory.zig b/src/memory.zig index f30e4392..e5394b4f 100644 --- a/src/memory.zig +++ b/src/memory.zig @@ -44,7 +44,48 @@ pub const Allocators = struct { pdf: Allocator, }; -var pardes_fallback: std.heap.StackFallbackAllocator(limits.arena.pardes) = undefined; +/// `std.heap.StackFallbackAllocator` with the buffer taken through its +/// lock-free interface. The core's allocator is used from two threads: a 9P +/// connection task answers a request while the editor's own thread is out +/// in a syscall that allocates (`readFileLimit`, `grep`; see `pardes.turn`). +fn SharedStackFallback(comptime size: usize) type { + return struct { + buffer: [size]u8 = undefined, + fallback_allocator: Allocator = undefined, + fixed: std.heap.FixedBufferAllocator = undefined, + + fn get(self: *@This()) Allocator { + self.fixed = .init(self.buffer[0..]); + return .{ .ptr = self, .vtable = &.{ .alloc = alloc, .resize = resize, .remap = remap, .free = free } }; + } + + fn alloc(ctx: *anyopaque, len: usize, alignment: std.mem.Alignment, ra: usize) ?[*]u8 { + const self: *@This() = @ptrCast(@alignCast(ctx)); + return self.fixed.threadSafeAllocator().rawAlloc(len, alignment, ra) orelse + self.fallback_allocator.rawAlloc(len, alignment, ra); + } + + fn resize(ctx: *anyopaque, buf: []u8, alignment: std.mem.Alignment, new_len: usize, ra: usize) bool { + const self: *@This() = @ptrCast(@alignCast(ctx)); + if (self.fixed.ownsPtr(buf.ptr)) return self.fixed.threadSafeAllocator().rawResize(buf, alignment, new_len, ra); + return self.fallback_allocator.rawResize(buf, alignment, new_len, ra); + } + + fn remap(ctx: *anyopaque, buf: []u8, alignment: std.mem.Alignment, new_len: usize, ra: usize) ?[*]u8 { + const self: *@This() = @ptrCast(@alignCast(ctx)); + if (self.fixed.ownsPtr(buf.ptr)) return self.fixed.threadSafeAllocator().rawRemap(buf, alignment, new_len, ra); + return self.fallback_allocator.rawRemap(buf, alignment, new_len, ra); + } + + fn free(ctx: *anyopaque, buf: []u8, alignment: std.mem.Alignment, ra: usize) void { + const self: *@This() = @ptrCast(@alignCast(ctx)); + if (self.fixed.ownsPtr(buf.ptr)) return self.fixed.threadSafeAllocator().rawFree(buf, alignment, ra); + return self.fallback_allocator.rawFree(buf, alignment, ra); + } + }; +} + +var pardes_fallback: SharedStackFallback(limits.arena.pardes) = .{}; var frame_fallback: std.heap.StackFallbackAllocator(limits.arena.frame) = undefined; var tree_sitter_fallback: std.heap.StackFallbackAllocator(limits.arena.tree_sitter) = undefined; var image_fallback: std.heap.StackFallbackAllocator(limits.arena.image) = undefined; @@ -61,7 +102,6 @@ var pdf_debug: Debug = .init; // One session at a time; free its allocations before deinit. Concurrent LSP workers use the caller's allocator. pub fn init(fallback: Allocator) Allocators { pardes_fallback.fallback_allocator = fallback; - pardes_fallback.get_called = if (std.debug.runtime_safety) false else {}; frame_fallback.fallback_allocator = fallback; frame_fallback.get_called = if (std.debug.runtime_safety) false else {}; tree_sitter_fallback.fallback_allocator = fallback; @@ -112,12 +152,12 @@ test "fixed allocators are separate, spill, and restart" { var allocs = init(std.testing.allocator); const core = try allocs.pardes.alloc(u8, 32); const frame = try allocs.frame.alloc(u8, 32); - try std.testing.expect(pardes_fallback.fixed_buffer_allocator.ownsPtr(core.ptr)); + try std.testing.expect(pardes_fallback.fixed.ownsPtr(core.ptr)); try std.testing.expect(frame_fallback.fixed_buffer_allocator.ownsPtr(frame.ptr)); try std.testing.expect(core.ptr != frame.ptr); const spill = try allocs.pardes.alloc(u8, limits.arena.pardes + 1); - try std.testing.expect(!pardes_fallback.fixed_buffer_allocator.ownsPtr(spill.ptr)); + try std.testing.expect(!pardes_fallback.fixed.ownsPtr(spill.ptr)); allocs.pardes.free(spill); allocs.frame.free(frame); allocs.pardes.free(core); @@ -127,7 +167,7 @@ test "fixed allocators are separate, spill, and restart" { defer deinit(); const restarted = try allocs.pardes.alloc(u8, 32); defer allocs.pardes.free(restarted); - try std.testing.expect(pardes_fallback.fixed_buffer_allocator.ownsPtr(restarted.ptr)); + try std.testing.expect(pardes_fallback.fixed.ownsPtr(restarted.ptr)); } test "memory limits preserve desktop capacities" { diff --git a/src/ninep/testing.zig b/src/ninep/testing.zig index 44329eee..6ea40956 100644 --- a/src/ninep/testing.zig +++ b/src/ninep/testing.zig @@ -30,13 +30,10 @@ pub const Answer = struct { }; pub fn call(p: *Pardes, req: Req) Answer { - p.update(.{ .fs_req = req }); var ans: Answer = .{}; + ans.reply = p.serveFs(req); + ans.bytes = p.fsPayload(ans.reply); while (p.nextEffect()) |e| switch (e) { - .fs_reply => |r| { - ans.reply = r; - ans.bytes = p.fsPayload(r); - }, .save_file, .save_text => ans.saved = true, .watch => |w| ans.watch = w.on, .write => |w| { diff --git a/src/ninep/tree.zig b/src/ninep/tree.zig index 5a0c1e11..b13b1472 100644 --- a/src/ninep/tree.zig +++ b/src/ninep/tree.zig @@ -68,6 +68,22 @@ pub const e_bad_event = "bad event syntax"; /// everywhere, as it is in acme (editors/acme/fsys.c, fsyscreate). pub const features: cloud9.fs.Features = .{ .remove = true }; +/// Whether a request must wait until no step is out in a syscall +/// (`pardes.turn`): what would change a pane such a step still points into. +/// Making a pane and rendering a screen do not -- every yield sits before +/// the mutation of its own step, so the layout and the surface are whole +/// under it -- and a step out there is never waiting on a write of its own, +/// so what parks here is never the syscall's own request. +pub fn needsQuiet(req: Req) bool { + return switch (req.op) { + .write, .setattr => true, + .release => req.remove, + else => false, + }; +} + +/// `needsQuiet`, plus the open that makes a pane: what costs a frame and a +/// wake of the editor. pub fn changesPane(req: Req) bool { return switch (req.op) { .write, .setattr => true, diff --git a/src/panes.zig b/src/panes.zig index f995986b..b3424385 100644 --- a/src/panes.zig +++ b/src/panes.zig @@ -3324,11 +3324,28 @@ pub const Image = struct { /// Construct an image pane from a path and optionally transferred dump bytes. /// `raw` must be image_gpa-owned and ownership transfers only on success. + /// Without them the file is read here, not at the first draw: a draw + /// must not go out into the host (`pardes.turn`), and the path may be a + /// mount this editor serves. A file that cannot be read draws blank. pub fn create(p: *pardes.Pardes, id: usize, path: []const u8, raw: []u8) !*pardes.Pane { + var bytes = raw; + var read_here = false; + if (bytes.len == 0) { + // The slot stays ours across the read: another request may + // make a pane meanwhile. + p.reserved_slots[id] = true; + defer p.reserved_slots[id] = false; + if (filesystem.read(p, path)) |from_disk| { + defer p.gpa.free(from_disk); + bytes = try p.image_gpa.dupe(u8, from_disk); + read_here = true; + } else |_| {} + } + errdefer if (read_here) p.image_gpa.free(bytes); const path_copy = try p.image_gpa.dupe(u8, path); errdefer p.image_gpa.free(path_copy); const pane = try p.newDocPane(id); - pane.image = .{ .path = path_copy, .raw = raw }; + pane.image = .{ .path = path_copy, .raw = bytes }; return pane; } @@ -3427,12 +3444,8 @@ pub const Image = struct { fn ensureDecoded(p: *pardes.Pardes, state: *State) void { if (state.tried) return; state.tried = true; - const bytes: []const u8 = if (state.raw.len > 0) - state.raw - else - filesystem.read(p, state.path) catch &.{}; - defer if (state.raw.len == 0) p.gpa.free(bytes); - if (image.decode(p.image_gpa, bytes)) |decoded| { + if (state.raw.len == 0) return; // nothing was readable at open + if (image.decode(p.image_gpa, state.raw)) |decoded| { state.rgba = decoded.rgba; state.iw = decoded.w; state.ih = decoded.h; @@ -3771,9 +3784,15 @@ pub const Pdf = struct { sections_output: ?SectionsOutput = null, outline_reveal_pending: ?pdf.OutlineInternalDestination = null, + /// Read whole and opened from memory (the bridge copies), never from + /// the file itself: MuPDF reads a file lazily at every page, and the + /// path may be a mount this editor serves, which only a read that + /// gives the turn up can come back from (`filesystem.readFile`). pub fn open(gpa: std.mem.Allocator, path: []const u8, page_one_based: usize) !@This() { const local = filesystem.localPath(path) orelse return error.NonLocalPath; - return initDocument(gpa, path, try Document.open(local), page_one_based); + const bytes = try filesystem.readFile(gpa, local); + defer gpa.free(bytes); + return initDocument(gpa, path, try Document.openBytes(bytes), page_one_based); } pub fn openBytes(gpa: std.mem.Allocator, path: []const u8, bytes: []const u8, page_one_based: usize) !@This() { @@ -5219,6 +5238,10 @@ pub const Pdf = struct { page_one_based: usize, ) !*pardes.Pane { if (comptime !enabled) return error.PdfDisabled; + // The slot stays ours across the read: another request may make a + // pane meanwhile. + core.reserved_slots[id] = true; + defer core.reserved_slots[id] = false; var state = if (filesystem.localPath(path) != null) try State.open(core.pdf_gpa, path, page_one_based) else virtual: { @@ -5903,9 +5926,11 @@ pub const Pdf = struct { null, }, .x = text_x, - .y = core.bodyTop(rect), + // Below the notice chips, like an image: a placed page is + // drawn after the cells and would paint a chip out. + .y = core.bodyTop(rect) + pane.notices.len, .w = text_width, - .h = rect.h -| pardes.BOX_H, + .h = (rect.h -| pardes.BOX_H) -| pane.notices.len, .rgba = placed.rgba, .iw = placed.width, // The texture contains the retained band, not the full page. diff --git a/src/pardes.zig b/src/pardes.zig index 0ac16162..cc57127f 100644 --- a/src/pardes.zig +++ b/src/pardes.zig @@ -39,6 +39,148 @@ pub const font_picker = platform == .gui or platform == .macos; pub const hosted = platform == .tty or platform == .gui or platform == .macos; +/// Whose turn it is with the core. The core is single-threaded: the editor's +/// thread has it, and lets go of it in two kinds of gap -- while it waits for +/// input (`rest`, `wake`), and while it is out in the host filesystem in the +/// middle of a step (`yield`, `back`). A 9P connection task takes the turn +/// in those gaps to answer a request (`take`, `give`), which is what lets the +/// editor read its own tree through a mount: the syscall goes out, another +/// thread answers it, the syscall comes back. A step of a connection task's +/// own -- a Look written to `look` -- goes out the same way. +/// +/// `out` counts the steps that are out. While one is, the core reads +/// consistently but that step still holds pointers into it, so nothing may +/// change it: a request that would is parked in the engine (src/9p_io.zig) +/// and retried when the turn is next given up with nothing out, which +/// `parked` remembers to arrange; and the editor's own wake waits for the +/// count to reach zero, since a step of the editor's may change anything. +/// 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. +/// +/// One per process, like the editor it belongs to. A process that never +/// starts it -- a test, the board, the browser -- has no second thread, and +/// every yield is a no-op. +pub var turn: Turn = .{}; + +pub const Turn = struct { + mutex: std.Io.Mutex = .init, + /// Set by `start`; a process that never starts the turn has none. + io: ?std.Io = null, + /// Steps out in a host syscall, from any thread. + out: u32 = 0, + /// A request that changes a pane was parked for want of quiet. + parked: bool = false, + wake_parked: ?*const fn (?*anyopaque) void = null, + wake_ctx: ?*anyopaque = null, + /// Bumped each time the editor has performed what the core asked of it, + /// so a request that asked for something -- a save, a shell -- can be + /// answered once it is done rather than merely queued, the way writing + /// `put` to acme's ctl returns after the file is written. + epoch: u64 = 0, + /// Bumped each time a shell has attempted the Restore the core asked for + /// (`takeRestore`), which lives after `pump` because it swaps the core + /// out from under it; a request that asked for one waits for this. + restores: u64 = 0, + /// Both of the above, and `out` reaching zero. + settled: std.Io.Condition = .init, + + /// The editor's thread takes the turn for good; it has it until `stop`. + pub fn start(t: *Turn, io: std.Io) void { + t.io = io; + t.mutex.lockUncancelable(io); + } + + pub fn stop(t: *Turn) void { + const io = t.io orelse return; + t.io = null; + t.parked = false; + t.mutex.unlock(io); + } + + /// The editor is about to wait for input: anyone may have the core. + pub fn rest(t: *Turn) void { + if (t.io == null) return; + t.release(); + } + + /// The editor takes the turn back -- once every step that is out in a + /// syscall has returned, since a step of the editor's may change what + /// those still point into. + pub fn wake(t: *Turn) void { + const io = t.io orelse return; + t.mutex.lockUncancelable(io); + while (t.out != 0) t.settled.waitUncancelable(io, &t.mutex); + } + + /// Whoever has the turn is about to block in the host in the middle of + /// a step. Nests: a directory listing that resolves each entry yields + /// inside a yield, and only the outermost one lets go. + pub fn yield(t: *Turn) void { + const io = t.io orelse return; + yield_depth += 1; + if (yield_depth != 1) return; + t.out += 1; + t.mutex.unlock(io); + } + + pub fn back(t: *Turn) void { + const io = t.io orelse return; + yield_depth -= 1; + if (yield_depth != 0) return; + t.mutex.lockUncancelable(io); + t.out -= 1; + if (t.out == 0) t.settled.broadcast(io); + } + + /// How many yields this thread is inside of. + threadlocal var yield_depth: u32 = 0; + + /// A connection task takes the turn; answers whether nothing is out + /// mid-step, which is what a request that changes a pane needs. + pub fn take(t: *Turn) bool { + t.mutex.lockUncancelable(t.io.?); + return t.out == 0; + } + + pub fn give(t: *Turn) void { + t.release(); + } + + /// The editor performed the core's effects. + pub fn settle(t: *Turn) void { + const io = t.io orelse return; + t.epoch +%= 1; + t.settled.broadcast(io); + } + + /// With the turn: waits, turn given up meanwhile, until the editor has + /// performed past `epoch`; stopping ends the wait early. + pub fn awaitSettled(t: *Turn, epoch: u64) void { + while (t.epoch == epoch) t.settled.wait(t.io.?, &t.mutex) catch return; + } + + /// The shell attempted the Restore the core asked for, whichever way it + /// went; called at the top of each loop step, when the previous step's + /// attempt is over. + pub fn restoreSettled(t: *Turn) void { + const io = t.io orelse return; + t.restores +%= 1; + t.settled.broadcast(io); + } + + pub fn awaitRestored(t: *Turn, restores: u64) void { + while (t.restores == restores) t.settled.wait(t.io.?, &t.mutex) catch return; + } + + fn release(t: *Turn) void { + const woken = t.out == 0 and t.parked; + if (woken) t.parked = false; + t.mutex.unlock(t.io.?); + if (woken) if (t.wake_parked) |f| f(t.wake_ctx); + } +}; + pub const can_attach = platform == .tty or platform == .gui; pub const terminal_panes = platform != .esp32p4; @@ -4754,7 +4896,6 @@ pub const Event = union(enum) { /// this must not clamp onto and preview the final grid cell. pointer_leave, tick, - fs_req: ctlfs.Req, }; pub const attach_name_max = 256; @@ -4792,7 +4933,6 @@ pub const Effect = union(enum) { /// Write the build-time theme ring below the per-user config directory. /// The pane receives the completion/error message from the native host. dump_themes: struct { pane: u8 }, - fs_reply: ctlfs.Reply, attach: struct { pane: u8, name: Buf(attach_name_max) }, detach: struct { pane: u8 }, quit, @@ -5910,11 +6050,6 @@ pub const Options = struct { ninep_name: []const u8 = "", ninep_tcp: ?[]const u8 = null, ninep_quic: ?[]const u8 = null, - ninep_identity: struct { - socket_path: []const u8 = "", - tcp_address: ?std.Io.net.IpAddress = null, - quic_address: ?std.Io.net.IpAddress = null, - } = .{}, mounts: []const filesystem.Mount = &.{}, config_dir: ?[]const u8 = null, startup_config_path: ?[]const u8 = null, @@ -6182,11 +6317,7 @@ pub const Pardes = struct { .pdf_gpa = pdf_gpa, .tree_sitter_gpa = tree_sitter_gpa, .opts = opts, - .fs = .{ - .socket_path = opts.ninep_identity.socket_path, - .tcp_address = opts.ninep_identity.tcp_address, - .quic_address = opts.ninep_identity.quic_address, - }, + .fs = .{}, .screen_w = opts.cols, .screen_h = opts.rows, .scratch = .init(gpa), @@ -6197,7 +6328,6 @@ pub const Pardes = struct { p.fs.started = ctlfs.events.now(); for (opts.mounts) |mount| try p.fs.mount(gpa, mount.name, mount.dial); p.opts.mounts = &.{}; - p.opts.ninep_identity = .{}; p.boot = Boot.of(opts); switch (p.boot) { .document => { @@ -6310,13 +6440,31 @@ pub const Pardes = struct { return @intCast(std.math.clamp(panes.File.displayWidth(text) + 2, 1, @as(usize, limit))); } + const Printed = struct { left: u16, dropped: usize }; + /// Prints `text` flush with the right edge of the band, and answers the - /// column it started at so a cursor can follow it. - fn printRight(s: *Surface, x: u16, row: u16, w: u16, text: []const u8, style: CellStyle) u16 { - const shown: u16 = @intCast(@min(@as(usize, w), panes.File.displayWidth(text))); - const left = x + w - shown; - _ = s.print(left, row, shown, text, style); - return left; + /// column it started at so a cursor can follow it. Text wider than the + /// band loses its HEAD: the end is the part that says something -- a + /// path's file name, an error's reason -- and `dropped` is how many + /// columns of it went. + fn printRight(s: *Surface, x: u16, row: u16, w: u16, text: []const u8, style: CellStyle) Printed { + const width = panes.File.displayWidth(text); + if (width <= w) { + const left: u16 = @intCast(x + w - width); + _ = s.print(left, row, @intCast(width), text, style); + return .{ .left = left, .dropped = 0 }; + } + // Whole glyphs only: a cut inside a wide one would leave the rest a + // column too wide and clip its LAST glyph instead of the first. + var dropped = width - w; + var start = panes.File.rawAtDisplay(text, dropped); + if (panes.File.rawDisplayCol(text, start) < dropped) { + dropped += 1; + start = panes.File.rawAtDisplay(text, dropped); + } + const left: u16 = @intCast(x + w - (width - dropped)); + _ = s.print(left, row, @intCast(width - dropped), text[start..], style); + return .{ .left = left, .dropped = dropped }; } /// The layout every single-pane boot starts from. @@ -6904,8 +7052,14 @@ pub const Pardes = struct { fn showMessage(p: *Pardes, id: usize, text: []const u8) void { if (id >= MAX_PANES) return; const pane = p.panes[id] orelse return; - pane.msg_len = @intCast(@min(text.len, pane.msg.len)); - @memcpy(pane.msg[0..pane.msg_len], text[0..pane.msg_len]); + // One row of printable text: a language server's multi-line report + // reads as one line, and the cut never leaves half a glyph behind. + var n: usize = @min(text.len, pane.msg.len); + if (n < text.len) { + while (n > 0 and text[n] & 0xC0 == 0x80) n -= 1; + } + for (text[0..n], pane.msg[0..n]) |c, *cell| cell.* = if (c < 0x20 or c == 0x7f) ' ' else c; + pane.msg_len = @intCast(n); } fn logMessage(p: *Pardes, id: usize, text: []const u8) void { @@ -7040,7 +7194,7 @@ pub const Pardes = struct { pub fn postEvent(p: *Pardes, ev: Event) void { switch (ev) { .key => |k| std.debug.assert(k.text.len == 0), - .output, .paste, .lsp_resp, .pipe_resp, .file_changed, .command, .fs_req => unreachable, + .output, .paste, .lsp_resp, .pipe_resp, .file_changed, .command => unreachable, else => {}, } if (p.in_len == p.in_q.len) return; @@ -7173,7 +7327,6 @@ pub const Pardes = struct { }, .theme_file => |t| if (v.watch_theme) |f| f(p.host.ctx, t.generation, t.on), .dump_themes => |d| if (v.dump_themes) |f| f(p.host.ctx, d.pane), - .fs_reply => {}, .attach => |a| { const name = a.name.slice(); @memcpy(p.attach_buf[0..name.len], name); @@ -7194,6 +7347,7 @@ pub const Pardes = struct { if (v.wait_input) |f| f(h.ctx, if (p.animationActive()) layout.Animation.frame_ms else 0); while (p.nextQueued()) |ev| p.update(ev); while (p.nextEffect()) |e| p.perform(e); + turn.settle(); // A quitting frame has already freed what it would draw. if (p.quit) return; if (v.poll_frame) |f| f(h.ctx); @@ -7208,24 +7362,35 @@ pub const Pardes = struct { p.needs_frame = false; } + /// One request of the control filesystem, from whichever thread has the + /// turn (`pardes.turn`). Not an event: it neither sweeps nor resets + /// anything the editor's own step may be in the middle of using, so it + /// can be answered while that step is out in a syscall. Only a request + /// that changes a pane costs a frame, which is why a round trip on a + /// local socket costs microseconds and not a vsync. + pub fn serveFs(p: *Pardes, req: ctlfs.Req) ctlfs.Reply { + if (!ctlfs.changesPane(req)) return ctlfs.handle(p, req); + p.needs_frame = true; + p.raw_hover_intent = false; + p.cancelLookHover(); + const reply = ctlfs.handle(p, req); + // The request was a whole step of its own, so it settles the way a + // step does: the cursor and scroll reconciled, the scripted panes + // told, and the panes it made announced to /log now rather than at + // the end of whatever the editor does next. + p.sync(); + p.fsReport(); + ctlfs.events.announce(p); + return reply; + } + pub fn update(p: *Pardes, ev: Event) void { - // Only two events cannot change the screen: a tick with nothing - // animating, and a filesystem request that only reads. The second is - // what 9P traffic overwhelmingly is, and answering it used to drag a - // whole vsync-blocked frame behind it -- which is why a round trip on - // a local socket cost a frame instead of a few microseconds. - p.needs_frame = p.needs_frame or switch (ev) { - .tick => false, - .fs_req => |req| ctlfs.changesPane(req), - else => true, - }; + // A tick with nothing animating is the one event that cannot change + // the screen. + p.needs_frame = p.needs_frame or ev != .tick; p.shell_rows.sweep(p.gpa); switch (ev) { .tick => {}, - .fs_req => |req| if (ctlfs.changesPane(req)) { - p.raw_hover_intent = false; - p.cancelLookHover(); - }, .mouse => |m| if (!(m.button == .none and m.kind == .motion)) { p.raw_hover_intent = false; p.cancelLookHover(); @@ -7324,7 +7489,6 @@ pub const Pardes = struct { }, .paste => |bytes| p.applyPaste(bytes), .command => |line| _ = p.executeBuiltinLine(p.active, line), - .fs_req => |r| p.emit(.{ .fs_reply = ctlfs.handle(p, r) }), .pinch => |scale| p.ov_pinch_scale = scale, .touch_scroll => |delta| p.ov_touch_scroll_delta = delta, .pointer_leave => { @@ -12915,11 +13079,7 @@ pub const Pardes = struct { .pdf_gpa = pdf_gpa, .tree_sitter_gpa = tree_sitter_gpa, .opts = opts, - .fs = .{ - .socket_path = opts.ninep_identity.socket_path, - .tcp_address = opts.ninep_identity.tcp_address, - .quic_address = opts.ninep_identity.quic_address, - }, + .fs = .{}, .screen_w = opts.cols, .screen_h = opts.rows, .scratch = .init(gpa), @@ -12934,13 +13094,9 @@ pub const Pardes = struct { p.cell_pixels = old.cell_pixels; p.row_metrics = old.row_metrics; p.native_images = old.native_images; - p.fs.socket_path = old.fs.socket_path; - p.fs.tcp_address = old.fs.tcp_address; - p.fs.quic_address = old.fs.quic_address; } for (opts.mounts) |mount| try p.fs.mount(gpa, mount.name, mount.dial); p.opts.mounts = &.{}; - p.opts.ninep_identity = .{}; var parsed = try dump.readZon(gpa, zon_bytes, "load"); defer parsed.deinit(); const st = parsed.value; @@ -13521,16 +13677,19 @@ pub const Pardes = struct { // right would put it past the pane, off the grid, and past what // the detached wire will encode -- which drops every frame for // as long as the prompt is up. - const left = printRight(s, cx, row, chip -| 1, text, msg_style); + const printed = printRight(s, cx, row, chip -| 1, text, msg_style); if (kind != .prompt or id != p.active) continue; const at = pane.promptAt() orelse continue; const prompt0 = (p.tagPrefix(pane) catch continue).len + at; const col = @as(usize, pane.tag_col); if (col < prompt0) continue; - // The cursor follows the text to wherever it landed. + // The cursor follows the text to wherever it landed; a caret + // in the part a narrow band dropped has nowhere to be. const prompt_col = panes.File.displayWidth(text[0..@min(col - prompt0, text.len)]); - if (left + prompt_col < tx + tw) - s.cursor = .{ .x = left + @as(u16, @intCast(prompt_col)), .y = row, .bar = pane.mode == .insert }; + if (prompt_col < printed.dropped) continue; + const caret = printed.left + (prompt_col - printed.dropped); + if (caret < tx + tw) + s.cursor = .{ .x = @intCast(caret), .y = row, .bar = pane.mode == .insert }; } } @@ -14018,12 +14177,19 @@ pub const Pardes = struct { const chip = r.x + r.w - cx; // Right aligned inside the chip, a blank cell short of its // edge: the same one the grid pass leaves for a prompt caret. + // Wider than the chip, the text loses its head, as on the grid. const room = p.tagCapacity(chip) -| 1; const shown = panes.File.displayWidth(text); - const pad = room -| shown; - const line = try arena.alloc(u8, pad + text.len); + var kept = text; + if (shown > room) { + var start = panes.File.rawAtDisplay(text, shown - room); + if (panes.File.rawDisplayCol(text, start) < shown - room) start = panes.File.rawAtDisplay(text, shown - room + 1); + kept = text[start..]; + } + const pad = room -| panes.File.displayWidth(kept); + const line = try arena.alloc(u8, pad + kept.len); @memset(line[0..pad], ' '); - @memcpy(line[pad..], text); + @memcpy(line[pad..], kept); try p.renderHeaderLayer(arena, NOTICE_LAYER_BASE + id * Pane.Notices.max + i, .notice, @intCast(id), .{ .x = cx, .y = first + @as(u16, @intCast(i)), @@ -14221,8 +14387,12 @@ pub const Pardes = struct { if (pane.hasPdf() and panes.Pdf.draw(p, pane, r, id, tx, tw)) return; if (pane.image) |*iv| { - const image_h = r.h -| BOX_H; - panes.Image.draw(p, iv, @intCast(id), pane.serial, tx, body_y, tw, image_h); + // Below the notice chips: a picture is drawn after the cells (the + // GUI's image pass, kitty's z-order), so a chip over it would be + // painted out. Text gets the overlay; a picture gives up the rows. + const shown = pane.notices.len; + const image_h = (r.h -| BOX_H) -| shown; + panes.Image.draw(p, iv, @intCast(id), pane.serial, tx, body_y + shown, tw, image_h); // thumbless, but the same one column as the real scrollbar below — // that is the whole point of drawing it, and like that one it runs // past the notice bands so the gutter has no notch in it diff --git a/src/tty/tty.zig b/src/tty/tty.zig index 23b7c09b..03ccd3e2 100644 --- a/src/tty/tty.zig +++ b/src/tty/tty.zig @@ -727,6 +727,7 @@ fn localSession( if (sh.inotify_fd >= 0) sh.watch_task = io.concurrent(watchFiles, .{ io, sh.inotify_fd, loop }) catch null; frames: while (!core.quit) { + pardes.turn.restoreSettled(); try core.pump(host); if (core.takeRestore()) |rp| blk: { const bytes = filesystem.readRestore(gpa, rp) catch |err| { @@ -760,7 +761,7 @@ fn localSession( _ = file_watch.applyThemeEffect(core, gpa, sh.inotify_fd, &sh.watches, 0, false, false); clearNativeImages(kitty_handles, vx, tty); nc.native_images = vx.caps.kitty_graphics; - if (fs) |f| f.reset(core); + if (fs) |f| f.reset(nc); core.deinit(); core = nc; sh.core = nc; @@ -871,15 +872,21 @@ const Shell = struct { fn waitInput(ctx: ?*anyopaque, timeout_ms: u32) void { const s = of(ctx); var batch: usize = 0; + // The wait is the 9P connections' turn with the core. if (timeout_ms == 0) { + pardes.turn.rest(); const first = s.loop.nextEvent() catch { + pardes.turn.wake(); s.core.quit = true; return s.reloadWatched(); }; + pardes.turn.wake(); batch = 1; if (s.apply(first)) return s.reloadWatched(); } else { + pardes.turn.rest(); std.Io.sleep(s.io, .fromMilliseconds(timeout_ms), .awake) catch {}; + pardes.turn.wake(); _ = s.apply(.tick); batch = 1; } @@ -994,7 +1001,7 @@ const Shell = struct { fn pollFrame(ctx: ?*anyopaque) void { const s = of(ctx); - if (s.fs) |f| if (f.tick(s.core).pending) { + if (s.fs) |f| if (f.tick().pending) { _ = s.loop.tryPostEvent(.fs_ready) catch {}; }; for (&s.ptys, 0..) |*slot, id| if (slot.*) |pt| { |
