diff options
Diffstat (limited to 'src/9p_io.zig')
| -rw-r--r-- | src/9p_io.zig | 145 |
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; |
