diff options
Diffstat (limited to 'src/9p_io.zig')
| -rw-r--r-- | src/9p_io.zig | 123 |
1 files changed, 82 insertions, 41 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig index 7c366aa8..7b0976bc 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -228,10 +228,29 @@ pub const Listener = struct { 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) if (pardes.ctlfs.openOf(core, req)) |o| { - o.held = .{ .asker = conn, .req = req }; - }; + // keeps its ticket to answer it when what it waits on has something. + if (reply.status == .again and req.op == .read) { + // Only an open's record can hold a read; one that waits with + // none would sit parked until some unrelated write parks. + const o = pardes.ctlfs.openOf(core, req) orelse { + log.err("a read of node {x} waits with no open record to hold it", .{req.node}); + std.debug.assert(false); + return conn.reply(&reply, ""); + }; + // One read waits on an open at a time, as acme's window keeps + // one `eventx`; a second is refused rather than left parked + // where nothing would ever answer it. The first asked again by + // a retry is the same request, and holds its place. + if (o.held) |held| if (held.req.tag != req.tag) { + const other: *Runner.Conn = @ptrCast(@alignCast(held.asker)); + if (other.waiting(held.ticket)) + return conn.reply(&pardes.ctlfs.failText(req.tag, pardes.ctlfs.E.BUSY, pardes.ctlfs.e_in_use), ""); + }; + o.held = null; + const ticket = conn.hold(&reply) orelse return; + o.held = .{ .asker = conn, .req = req, .ticket = ticket }; + return; + } 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 @@ -253,10 +272,11 @@ pub const Listener = struct { /// 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. + /// retry of anything else parked there. Only while its ticket still + /// waits, asked in the same hold of the engine as the answer: Tflush + /// and clunk drop a park without a word to the backend, a `retry` takes + /// one out to ask again (and that asking answers it), and a record spent + /// on any of them would be lost to the reader that comes next. fn answerHeld(ctx: ?*anyopaque) void { if (comptime !supported) return; const l = of(ctx); @@ -267,27 +287,22 @@ pub const Listener = struct { for (&core.fs.opens) |*o| { const held = o.held 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) o.held = null; - conn.unlock(); - continue; - } - const reply = core.serveFs(held.req); - if (reply.status != .again) { - o.held = null; - conn.engine.reply(&reply, core.fsPayload(reply)); - } - conn.unlock(); - conn.flush(); + if (!conn.answerWith(held.ticket, Held{ .core = core, .req = held.req }, Held.make)) o.held = null; } } + /// A held read asked again, to make its answer while its park waits. + const Held = struct { + core: *pardes.Pardes, + req: pardes.ctlfs.Req, + + fn make(h: Held) ?Runner.Conn.Answer { + const reply = h.core.serveFs(h.req); + if (reply.status == .again) return null; + return .{ .reply = reply, .bytes = h.core.fsPayload(reply) }; + } + }; + /// The turn was given up quiet with a request parked: every connection /// retries what it parked. fn wakeParked(ctx: ?*anyopaque) void { @@ -1407,41 +1422,65 @@ test "a held read is answered when the log has news, and a flushed one spends no _ = try s.ask(.{ .write = .{ .fid = 1, .offset = 0, .data = "follow\n" } }, &remote); const heldNow = struct { - fn check(core: *pardes.Pardes) !void { + /// Waits until the core holds `n` reads. + fn count(core: *pardes.Pardes, n: usize) !void { const deadline = Client.nowMs() + Client.budget_ms; while (Client.nowMs() < deadline) { pardes.turn.wake(); - const any = for (core.fs.opens) |o| { - if (o.held != null) break true; - } else false; + var got: usize = 0; + for (core.fs.opens) |o| got += @intFromBool(o.held != null); pardes.turn.rest(); - if (any) return; + if (got == n) return; Client.nap(1); } return error.Timeout; } }; + const serial = p.panes[0].?.serial; + var event_path: [32]u8 = undefined; + var event_names: [3][]const u8 = .{ "pane", try std.fmt.bufPrint(&event_path, "{d}", .{serial}), "event" }; + _ = try s.ask(.{ .walk = .{ .fid = 0, .newfid = 2, .names = &event_names } }, &remote); + _ = try s.ask(.{ .open = .{ .fid = 2, .mode = ninep.oread } }, &remote); // 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. + // hearing of it. Its tag, asked again at once for a read of the pane's + // event, must not be answered with the log's next record, and that + // record must reach the log's next read. const flushed = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); try s.flush(); - try heldNow.check(p); + try heldNow.count(p, 1); _ = 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); + const reused = try s.cl.submit(.{ .read = .{ .fid = 2, .offset = 0, .count = 4096 } }); + try testing.expectEqual(flushed, reused); + try s.flush(); + try heldNow.count(p, 2); pardes.turn.wake(); p.setMessage(0, "after the flush"); pardes.turn.rest(); + try heldNow.count(p, 1); + pardes.turn.wake(); + p.fs.origin = 'K'; + _ = pardes.ctlfs.events.noteAction(p, 0, .body_exec, 0, 4, 0, "held"); + pardes.turn.rest(); + const clicked = try s.settle(); + try testing.expectEqual(reused, clicked.tag); + try testing.expectEqualStrings("KX0 4 0 4 held\n", clicked.result.read); 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. + // connection, with no retry of every parked request asked for; a + // second read on the same open meanwhile is refused, not orphaned. const waiting = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); try s.flush(); - try heldNow.check(p); + try heldNow.count(p, 1); + const second = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); + const refused = try s.settle(); + try testing.expectEqual(second, refused.tag); + try testing.expectEqualStrings(pardes.ctlfs.e_in_use, refused.result.fail); pardes.turn.wake(); p.setMessage(0, "while held"); const retry_all = pardes.turn.parked; @@ -1450,13 +1489,15 @@ test "a held read is answered when the log has news, and a flushed one spends no const answered = try s.settle(); try testing.expectEqual(waiting, answered.tag); try testing.expect(std.mem.indexOf(u8, answered.result.read, "while held") != null); + + // Clunked, the log's held read is interrupted by the engine and the + // record it would have had is spent on nobody. + _ = try s.cl.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 4096 } }); + try s.flush(); + try heldNow.count(p, 1); s.drop(1); - pardes.turn.wake(); - const left = for (p.fs.opens) |o| { - if (o.held != null) break true; - } else false; - pardes.turn.rest(); - try testing.expect(!left); + s.drop(2); + try heldNow.count(p, 0); } extern "c" fn setenv(name: [*:0]const u8, value: [*:0]const u8, overwrite: c_int) c_int; |
