diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-27 20:09:02 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-10-01 00:12:14 -0300 |
| commit | 28c814aa5cfea23e8950ef007916ffc5d089288d (patch) | |
| tree | 6e8c017c2140e1e7e8d0f3ee584a7f4aee1050f8 /src | |
| parent | c1be5b6b7c11f5dc17fd221c8d01b5693d0ffd0b (diff) | |
| download | pardes-28c814aa5cfea23e8950ef007916ffc5d089288d.tar.gz pardes-28c814aa5cfea23e8950ef007916ffc5d089288d.zip | |
Hold a read that has nothing yet and answer it when its file has news
A following log, event, pty/data or a pty/run before its answer used to
answer .again and wait for a wakeAll, which only the parked-write path
asked for, so the band-aid had every queue push set turn.parked. Now the
core keeps such a read (ctlfs.hold) and, as the turn is given up after
anything that queued a record, ran a command out or closed a pane,
answers it on its own connection, the way factotum answers the log reads
it keeps and acme an event read. Only a read the engine still holds
parked is answered, because cloud9 tells the backend nothing of a
Tflush, so a flushed read spends no record.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
Diffstat (limited to 'src')
| -rw-r--r-- | src/9p_io.zig | 145 | ||||
| -rw-r--r-- | src/fs.zig | 6 | ||||
| -rw-r--r-- | src/ninep/events.zig | 10 | ||||
| -rw-r--r-- | src/ninep/pty.zig | 18 | ||||
| -rw-r--r-- | src/ninep/tree.zig | 31 | ||||
| -rw-r--r-- | src/pardes.zig | 5 |
6 files changed, 201 insertions, 14 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig index 7e09fb27..6182ca75 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -227,6 +227,9 @@ pub const Listener = struct { const restores = pardes.turn.restores; const reply = core.serveFs(req); if (req.op == .release and quiet) l.collectOs(); + // A read with nothing yet stays parked in the engine, and the core + // keeps it to answer when what it waits on has something. + if (reply.status == .again and req.op == .read) pardes.ctlfs.hold(core, req, conn); 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 @@ -246,6 +249,43 @@ pub const Listener = struct { conn.reply(&reply, ""); } + /// With the turn, as it is given up: answers each read the core holds + /// whose file now has something, on the connection that asked, with no + /// retry of anything else parked there. Only a read its engine still + /// holds parked is answered: Tflush drops one without a word to the + /// backend, and a `retry` takes one out to ask again, and a record spent + /// on either would be lost to the reader that comes next. + fn answerHeld(ctx: ?*anyopaque) void { + if (comptime !supported) return; + const l = of(ctx); + if (l.stopping.load(.acquire)) return; + const core = l.core; + if (!core.fs.news) return; + core.fs.news = false; + for (&core.fs.held) |*slot| { + const held = slot.* orelse continue; + const conn: *Runner.Conn = @ptrCast(@alignCast(held.asker)); + conn.lock(); + const parked: ?bool = for (conn.engine.slots) |sl| { + if (sl.used and sl.req.tag == held.req.tag) break sl.parked; + } else null; + // Gone (flushed, answered, hung up) it is forgotten; out being + // asked again, the asking answers it. + if (parked != true) { + if (parked == null) slot.* = null; + conn.unlock(); + continue; + } + const reply = core.serveFs(held.req); + if (reply.status != .again) { + slot.* = null; + conn.engine.reply(&reply, core.fsPayload(reply)); + } + conn.unlock(); + conn.flush(); + } + } + /// The turn was given up quiet with a request parked: every connection /// retries what it parked. fn wakeParked(ctx: ?*anyopaque) void { @@ -641,6 +681,7 @@ pub const Listener = struct { pardes.turn.wake(); } pardes.turn.wake_parked = null; + pardes.turn.answer_held = null; pardes.turn.stop(); if (l.path_len != 0) { var z: [sun_path_len:0]u8 = undefined; @@ -672,6 +713,7 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes, named: [ l.* = .{ .io = io, .core = core }; pardes.turn.start(io); pardes.turn.wake_parked = Listener.wakeParked; + pardes.turn.answer_held = Listener.answerHeld; pardes.turn.wake_ctx = l; l.runner.init(.{ .io = io, @@ -1312,6 +1354,109 @@ test "a change waits while the editor is out mid-step, a read does not, and the try testing.expectEqualStrings("after\n", pane.file.?.content); } +test "a held read is answered when the log has news, and a flushed one spends nothing" { + if (comptime !supported) return error.SkipZigTest; + const gpa = testing.allocator; + var directory: [64:0]u8 = undefined; + _ = try std.fmt.bufPrintSentinel(&directory, "/tmp/pardes-held-XXXXXX", .{}, 0); + if (mkdtemp(&directory) == null) return error.TempDirectoryFailed; + 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(); + _ = try p.setTestFile("held\n"); + p.update(.tick); // the pane is logged, before the log is opened + while (p.nextEffect()) |_| {} + const l = listen(testing.io, gpa, p, "held", "", null, null) orelse return error.ListenFailed; + defer { + l.reset(p); + l.deinit(gpa); + } + // This thread is the editor's and plays the client too, resting while + // it does so that the runner's tasks can take the turn and answer. + const s = try gpa.create(Client.Session); + defer gpa.destroy(s); + s.* = .{ .fd = -1, .deadline = Client.nowMs() + 3 * Client.budget_ms, .display_path = "/log" }; + var sock_buf: [sun_path_len]u8 = undefined; + pardes.turn.rest(); + defer pardes.turn.wake(); + s.fd = try Client.connect(try Client.resolve(&sock_buf, l.path()), s.deadline); + defer _ = libc.close(s.fd); + s.cl = .init(.{ .in = &s.in, .out = &s.out }); + var remote: Client.RemoteError = .{}; + _ = try s.ask(.{ .version = .{} }, &remote); + _ = try s.ask(.{ .attach = .{ .fid = 0, .uname = "held" } }, &remote); + _ = try s.ask(.{ .walk = .{ .fid = 0, .newfid = 1, .names = &.{"log"} } }, &remote); + _ = try s.ask(.{ .open = .{ .fid = 1, .mode = ninep.ordwr } }, &remote); + var frozen: u64 = 0; + while (true) { + const n = (try s.ask(.{ .read = .{ .fid = 1, .offset = frozen, .count = 4096 } }, &remote)).read.len; + if (n == 0) break; + frozen += n; + } + _ = try s.ask(.{ .write = .{ .fid = 1, .offset = 0, .data = "follow\n" } }, &remote); + + const heldNow = struct { + fn check(core: *pardes.Pardes) !void { + const deadline = Client.nowMs() + Client.budget_ms; + while (Client.nowMs() < deadline) { + pardes.turn.wake(); + const any = for (core.fs.held) |slot| { + if (slot != null) break true; + } else false; + pardes.turn.rest(); + if (any) return; + Client.nap(1); + } + return error.Timeout; + } + }; + + // Flushed while held, a read is gone from the engine without the core + // hearing of it; the record logged next must reach the read after it. + const flushed = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); + try s.flush(); + try heldNow.check(p); + _ = try s.cl.submit(.{ .flush = .{ .oldtag = flushed } }); + const interrupted = try s.settle(); + try testing.expectEqual(flushed, interrupted.tag); + try testing.expect(interrupted.result == .fail); + try testing.expect((try s.settle()).result == .flush); + pardes.turn.wake(); + p.setMessage(0, "after the flush"); + pardes.turn.rest(); + try testing.expect(std.mem.indexOf(u8, (try s.ask(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }, &remote)).read, "after the flush") != null); + + // Held, a read is answered by the record that arrives, on its own + // connection, with no retry of every parked request asked for. + const waiting = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); + try s.flush(); + try heldNow.check(p); + pardes.turn.wake(); + p.setMessage(0, "while held"); + const retry_all = pardes.turn.parked; + pardes.turn.rest(); + try testing.expect(!retry_all); + const answered = try s.settle(); + try testing.expectEqual(waiting, answered.tag); + try testing.expect(std.mem.indexOf(u8, answered.result.read, "while held") != null); + s.drop(1); + pardes.turn.wake(); + const left = for (p.fs.held) |slot| { + if (slot != null) break true; + } else false; + pardes.turn.rest(); + try testing.expect(!left); +} + 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; @@ -1243,6 +1243,12 @@ pub const Namespace = struct { panes: [MAX_PANES]tree.pane.State = @splat(.{}), listeners: u16 = 0, origin: u8 = 'K', + /// Reads that found nothing yet, one per open that can wait: the log's + /// and a run's opens, and one event and one pty/data reader a pane. + held: [tree.screen.snapshot_slots + tree.pty.run_slots + 2 * MAX_PANES]?tree.Held = @splat(null), + /// Something a held read may be waiting on changed since they were last + /// answered: a record queued, a run answered, a pane gone. + news: bool = false, /// The editor-wide event ring: panes made, renamed, saved and closed, and /// what the editor said. Recorded whether or not anyone reads /log. log: tree.events.Queue = .{ .cap = limits.log_bytes }, diff --git a/src/ninep/events.zig b/src/ninep/events.zig index bd5fecf3..6b8c744e 100644 --- a/src/ninep/events.zig +++ b/src/ninep/events.zig @@ -50,10 +50,6 @@ pub const Queue = struct { q.buf.shrinkRetainingCapacity(q.buf.items.len - 4); return; }; - // A read parked on this queue is retried only when the turn is next - // given up with `parked` set; without it a follower sleeps until some - // unrelated request happens to park. - pardes.turn.parked = true; } pub fn peek(q: *const Queue) ?[]const u8 { @@ -105,7 +101,7 @@ pub fn pending(q: *const Queue) u64 { return record.len; } -/// One record per read; `.again` parks the read until a record arrives. +/// One record per read; `.again` holds the read until a record arrives. pub fn readQueue(p: *Pardes, req: Req, q: *Queue) Reply { const record = q.peek() orelse return .{ .tag = req.tag, .status = .again }; if (req.size < record.len) return Reply.fail(req.tag, E.INVAL); @@ -129,6 +125,7 @@ pub fn noteInstall(p: *Pardes, id: usize) void { pub fn noteRetire(p: *Pardes, id: usize, pane: *Pane) void { if (id >= MAX_PANES) return; tree.pty.shellGone(p, id); + p.fs.news = true; // a read held on its event or pty/data hears it went if (p.fs.panes[id].unannounced) { p.fs.panes[id].unannounced = false; return; @@ -174,6 +171,7 @@ fn pushLog(p: *Pardes, record: []u8) void { while (end > 0 and record[end] & 0xC0 == 0x80) end -= 1; record[end] = '\n'; p.fs.log.push(p.gpa, record[0 .. end + 1]); + p.fs.news = true; } /// An open freezes the ring's text, so `cat log` answers what happened lately @@ -384,6 +382,7 @@ pub fn noteAction( var buf: [max_record_text + 64]u8 = undefined; const record = formatRecord(&buf, p.fs.origin, action, q0, q1, flag, text); p.fs.panes[id].events.push(p.gpa, record); + p.fs.news = true; return true; } @@ -397,6 +396,7 @@ pub fn notePtyOutput(p: *Pardes, id: usize, bytes: []const u8) void { pf.pty_out.push(p.gpa, bytes[off..][0..n]); off += n; } + p.fs.news = true; } const EventRecord = struct { action: Action, q0: u32, q1: u32 }; diff --git a/src/ninep/pty.zig b/src/ninep/pty.zig index a689d87f..2a38f1f0 100644 --- a/src/ninep/pty.zig +++ b/src/ninep/pty.zig @@ -162,10 +162,10 @@ pub fn openRun(p: *Pardes, req: Req, serial: u32) Reply { return Reply.fail(req.tag, E.NFILE); } -fn answer(slot: *Run, comptime fmt: []const u8, args: anytype) void { +fn answer(p: *Pardes, slot: *Run, comptime fmt: []const u8, args: anytype) void { slot.len = @intCast((std.fmt.bufPrint(&slot.answer, fmt ++ "\n", args) catch unreachable).len); slot.phase = .done; - pardes.turn.parked = true; // wake the read waiting on it + p.fs.news = true; // the read held on it can be answered } pub fn writeRun(p: *Pardes, req: Req, id: usize, pane: *Pane) Reply { @@ -180,7 +180,7 @@ pub fn writeRun(p: *Pardes, req: Req, id: usize, pane: *Pane) Reply { const state = pane.terminal orelse return tree.failText(req.tag, E.INVAL, e_bad_line); const marks = &state.stream.handler; if (pf.unmarked) { - answer(slot, "error no prompt marks", .{}); + answer(p, slot, "error no prompt marks", .{}); } else if (pf.run != null or marks.phase != .input or !pardes.panes.Terminal.promptInputEmpty(pane) or p.hostTtyTaken(id)) { @@ -189,7 +189,7 @@ pub fn writeRun(p: *Pardes, req: Req, id: usize, pane: *Pane) Reply { // the middle of it. The phase is pardes's own marks: a nested // shell's prompt (ssh, a shell with its own integration) looks like // an empty prompt to ghostty but is not the shell this run knows. - answer(slot, "busy", .{}); + answer(p, slot, "busy", .{}); } else { slot.want = marks.started +% 1; slot.prompts = marks.prompts; @@ -285,7 +285,7 @@ pub fn noteMarks(p: *Pardes, id: usize, pane: *Pane) void { .kept_line => {}, } if (!empty) return; - answer(slot, "error not run", .{}); + answer(p, slot, "error not run", .{}); pf.run = null; return; } @@ -305,11 +305,11 @@ pub fn noteMarks(p: *Pardes, id: usize, pane: *Pane) void { // The header is the whole first line, so a count there can never be // mistaken for output; `cut` with no count: its start scrolled away. if (printed == null or scrolled_out) - answer(slot, "exit {d} cut", .{status}) + answer(p, slot, "exit {d} cut", .{status}) else if (cut > 0) - answer(slot, "exit {d} cut {d}", .{ status, cut }) + answer(p, slot, "exit {d} cut {d}", .{ status, cut }) else - answer(slot, "exit {d}", .{status}); + answer(p, slot, "exit {d}", .{status}); if (keep.len > 0) { const nl = @intFromBool(keep[keep.len - 1] != '\n'); if (p.gpa.alloc(u8, keep.len + nl)) |owned| { @@ -326,7 +326,7 @@ pub fn noteMarks(p: *Pardes, id: usize, pane: *Pane) void { pub fn shellGone(p: *Pardes, id: usize) void { const pf = &p.fs.panes[id]; const idx = pf.run orelse return; - answer(&p.fs.runs[idx], "error shell gone", .{}); + answer(p, &p.fs.runs[idx], "error shell gone", .{}); pf.run = null; } diff --git a/src/ninep/tree.zig b/src/ninep/tree.zig index e8bcad2f..3ebea669 100644 --- a/src/ninep/tree.zig +++ b/src/ninep/tree.zig @@ -107,6 +107,32 @@ pub fn changesPane(req: Req) bool { pub const out_reserve = 4 * 1024; +/// A read that found nothing yet (`.again`), kept to be answered when what it +/// waits on has something, the way factotum keeps its log's waiting reads +/// and answers them on append (security/auth/factotum/log.c:4-52) and acme +/// an event read (editors/acme/xfid.c:994). `asker` is the connection it +/// came on, which only the listener knows (src/9p_io.zig, `answerHeld`). +pub const Held = struct { asker: *anyopaque, req: Req }; + +/// Keeps a read that answered `.again`. An open waits with one read at a +/// time, as acme's window keeps one `eventx`: a newer read on the same open +/// takes the place of the older. +pub fn hold(p: *Pardes, req: Req, asker: *anyopaque) void { + var free: ?*?Held = null; + for (&p.fs.held) |*slot| { + const held = slot.* orelse { + if (free == null) free = slot; + continue; + }; + if (held.req.node == req.node and held.req.handle == req.handle) { + slot.* = .{ .asker = asker, .req = req }; + return; + } + } + // One slot per open that can wait, so a free one is always there. + if (free) |slot| slot.* = .{ .asker = asker, .req = req }; +} + // ---- nodes ---- pub const TopFile = enum(u4) { @@ -559,6 +585,11 @@ fn remove(p: *Pardes, req: Req) Reply { } fn release(p: *Pardes, req: Req) Reply { + // A read held on this open goes with it: a hangup pays its releases + // before the connection's slot is reused, so none outlives its asker. + for (&p.fs.held) |*slot| if (slot.*) |held| { + if (held.req.node == req.node and held.req.handle == req.handle) slot.* = null; + }; // The handle's bookkeeping runs whether or not the removal is allowed. const done = releaseHandle(p, req); return if (req.remove) remove(p, req) else done; diff --git a/src/pardes.zig b/src/pardes.zig index 29a9ab3a..a38eb837 100644 --- a/src/pardes.zig +++ b/src/pardes.zig @@ -72,6 +72,8 @@ pub const Turn = struct { /// A request that changes a pane was parked for want of quiet. parked: bool = false, wake_parked: ?*const fn (?*anyopaque) void = null, + /// Answers the reads the core holds, with the turn, as it is given up. + answer_held: ?*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 @@ -174,6 +176,9 @@ pub const Turn = struct { } fn release(t: *Turn) void { + // Whatever the holder just did may be what a held read waits for, + // and answering it takes the core, so before letting go. + if (t.answer_held) |f| f(t.wake_ctx); const woken = t.out == 0 and t.parked; if (woken) t.parked = false; t.mutex.unlock(t.io.?); |
