summaryrefslogtreecommitdiff
path: root/src/9p_io.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/9p_io.zig')
-rw-r--r--src/9p_io.zig123
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;