summaryrefslogtreecommitdiff
path: root/src/9p_io.zig
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/9p_io.zig
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/9p_io.zig')
-rw-r--r--src/9p_io.zig145
1 files changed, 145 insertions, 0 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;