summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-27 20:09:02 -0300
committerGabriel Schneider <[email protected]>2026-10-01 00:12:14 -0300
commit28c814aa5cfea23e8950ef007916ffc5d089288d (patch)
tree6e8c017c2140e1e7e8d0f3ee584a7f4aee1050f8 /src
parentc1be5b6b7c11f5dc17fd221c8d01b5693d0ffd0b (diff)
downloadpardes-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.zig145
-rw-r--r--src/fs.zig6
-rw-r--r--src/ninep/events.zig10
-rw-r--r--src/ninep/pty.zig18
-rw-r--r--src/ninep/tree.zig31
-rw-r--r--src/pardes.zig5
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;
diff --git a/src/fs.zig b/src/fs.zig
index 8e2c52a1..50a446dc 100644
--- a/src/fs.zig
+++ b/src/fs.zig
@@ -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.?);