diff options
Diffstat (limited to 'src/fs.zig')
| -rw-r--r-- | src/fs.zig | 176 |
1 files changed, 162 insertions, 14 deletions
@@ -49,8 +49,11 @@ pub const Op = enum(u8) { readdir, }; -/// `again` parks a read or write until the backend answers the same tag -/// later (or the engine retries it); any other op answering `again` fails. +/// `again` parks the request until the engine retries it (or, for a read or +/// write, until the backend answers the same tag later). A read, readdir, +/// write, open, clunk, remove or truncating wstat can park; a walk, attach, +/// stat, create or renaming wstat cannot, because their names live in the +/// input frame that parking lets go of, and those answer an error instead. pub const Status = enum(u8) { ok, again, @@ -391,6 +394,20 @@ pub fn Server(comptime Backend: type, comptime opts: Options) type { count: u32 = 0, req: Backend.Req = undefined, data: [park_data_max]u8 = undefined, + /// What a parked open needs to become a job again; a parked + /// wstat is always a truncation to zero, and a clunk or remove + /// keeps nothing but its fid. + omode: u8 = 0, + step: u8 = 0, + + /// Whether the slot holds a whole job rather than a read or write + /// the engine answers from the slot itself. + fn holdsJob(sl: *const Slot) bool { + return switch (sl.kind) { + .open, .wstat, .clunk, .remove => true, + else => false, + }; + } }; const Job = struct { @@ -424,6 +441,9 @@ pub fn Server(comptime Backend: type, comptime opts: Options) type { forgetting: bool = false, /// The reply went out; only releases remain. closing: bool = false, + /// This is a parked job asked again this retry round: parked + /// once more, it sits the round out like a retried read does. + retried: bool = false, /// A create's permissions, a wstat's changes. cperm: u32 = 0, set: Set = .{}, @@ -659,6 +679,9 @@ pub fn Server(comptime Backend: type, comptime opts: Options) type { var best: ?usize = null; for (&s.slots, 0..) |*sl, i| { if (!sl.used or !sl.parked or sl.retried) continue; + // A parked job goes back into the job slot, so it waits its + // turn for that. + if (sl.holdsJob() and s.job.kind != .none) continue; if (best == null or sl.seq < s.slots[best.?].seq) best = i; } const i = best orelse { @@ -669,9 +692,30 @@ pub fn Server(comptime Backend: type, comptime opts: Options) type { for (&s.slots) |*sl| sl.retried = false; return null; } - s.slots[i].retried = true; - s.slots[i].parked = false; - return s.slots[i].req; + const sl = &s.slots[i]; + if (sl.holdsJob()) return s.resumeJob(i); + sl.retried = true; + sl.parked = false; + return sl.req; + } + + /// Takes a parked job out of its slot and asks its request again; + /// the reply then goes through `jobReply` like any other. + fn resumeJob(s: *Self, i: usize) Backend.Req { + const sl = &s.slots[i]; + assert(s.job.kind == .none); + s.job = .{ + .kind = sl.kind, + .tag = sl.tag, + .fid = sl.fid, + .omode = sl.omode, + .step = sl.step, + .set = if (sl.kind == .wstat) .{ .length = true } else .{}, + .retried = true, + }; + const req = sl.req; + sl.* = .{}; + return s.ask(req); } /// The next backend operation, or null while one is outstanding, the @@ -1161,6 +1205,9 @@ pub fn Server(comptime Backend: type, comptime opts: Options) type { j.forgetting = false; return; } + // Before a clunk lets its fid go: a parked release keeps the fid + // until the release is really paid. + if (r.status == .again) return s.parkJob(); if (j.kind == .clunk or j.kind == .remove) { s.dropFid(j.fid); if (j.kind == .clunk) @@ -1174,7 +1221,6 @@ pub fn Server(comptime Backend: type, comptime opts: Options) type { s.finishJob(); return; } - if (r.status == .again) return s.parkJob(); if (r.status == .err) { if (j.kind == .walk) return s.stopWalk(if (r.ename.len != 0) r.ename else errString(r.errno)); s.failReply(j.tag, r); @@ -1310,13 +1356,18 @@ pub fn Server(comptime Backend: type, comptime opts: Options) type { fn parkJob(s: *Self) void { const j = &s.job; - switch (j.kind) { - .read, .readdir, .write => {}, - else => { - s.fail(j.tag, e_again); - s.finishJob(); - return; - }, + // A wstat parks only as the truncation to zero a Linux client + // sends for O_TRUNC: a rename's name is in the frame, and a mode + // or mtime would need keeping. + const parkable = switch (j.kind) { + .read, .readdir, .write, .open, .clunk, .remove => true, + .wstat => j.set.length and j.length == 0 and !j.set.name and !j.set.mode and !j.set.mtime, + else => false, + }; + if (!parkable) { + s.fail(j.tag, e_again); + s.finishJob(); + return; } const i = s.freeSlot() orelse { s.fail(j.tag, e_again); @@ -1327,12 +1378,15 @@ pub fn Server(comptime Backend: type, comptime opts: Options) type { sl.* = .{ .used = true, .parked = true, + .retried = j.retried, .seq = s.tick(), .tag = j.tag, .kind = j.kind, .fid = j.fid, .count = j.count, .req = j.req, + .omode = j.omode, + .step = j.step, }; if (j.req.data.len != 0) { if (j.req.data.len > park_data_max) { @@ -1354,6 +1408,16 @@ pub fn Server(comptime Backend: type, comptime opts: Options) type { sl.parked = true; return; } + // A parked job is answered as a job. One that cannot become the + // job right now stays parked; the retry asks it again. + if (sl.holdsJob()) { + if (s.job.kind != .none) { + sl.parked = true; + return; + } + _ = s.resumeJob(i); + return s.jobReply(r, bytes); + } if (r.status == .err) { s.failReply(sl.tag, r); sl.* = .{}; @@ -1477,6 +1541,10 @@ const StubFs = struct { filler: [1024]u8 = @splat('x'), event: ?[]const u8 = null, park_writes: bool = false, + park_opens: bool = false, + park_setattr: bool = false, + park_releases: bool = false, + park_lookups: bool = false, releases: u32 = 0, calls: u32 = 0, writes: [128]u8 = undefined, @@ -1534,6 +1602,7 @@ const StubFs = struct { const i = find(req.node) orelse return fail; switch (req.op) { .lookup => { + if (st.park_lookups) return .{ .reply = .{ .tag = req.tag, .status = .again } }; if (std.mem.eql(u8, req.data, "..")) return .{ .reply = .{ .tag = req.tag, .attr = st.attrOf(tree[find(tree[i].parent).?]) } }; for (tree) |e| { @@ -1545,11 +1614,16 @@ const StubFs = struct { }, .getattr => return .{ .reply = .{ .tag = req.tag, .attr = st.attrOf(tree[i]) } }, .setattr => { + if (st.park_setattr) return .{ .reply = .{ .tag = req.tag, .status = .again } }; if (req.truncate and req.node == body_node) st.body = ""; return .{ .reply = .{ .tag = req.tag, .attr = st.attrOf(tree[i]) } }; }, - .open => return .{ .reply = .{ .tag = req.tag, .handle = 7 } }, + .open => { + if (st.park_opens) return .{ .reply = .{ .tag = req.tag, .status = .again } }; + return .{ .reply = .{ .tag = req.tag, .handle = 7 } }; + }, .release => { + if (st.park_releases) return .{ .reply = .{ .tag = req.tag, .status = .again } }; st.releases += 1; return .{ .reply = .{ .tag = req.tag } }; }, @@ -2088,6 +2162,80 @@ test "fs server: a blocked read parks, and the connection keeps working" { try testing.expectEqualStrings(e_again, got.msg.rerror.ename); } +test "fs server: a parked open, truncate and clunk complete on retry, and a walk cannot park" { + var h: Harness = .{}; + try h.handshake(4096); + + // An open the backend is not ready for: parked, and the connection keeps + // answering everything else meanwhile. + _ = try h.walkTo(5, 1, &.{ "1", "body" }); + h.fsys.park_opens = true; + try h.send(6, .{ .topen = .{ .fid = 1, .mode = oread } }); + try h.quiet(); + _ = try h.walkTo(7, 2, &.{ "1", "tag" }); + try h.send(8, .{ .tstat = .{ .fid = 2 } }); + var got = try h.reap(); + try testing.expectEqualStrings("tag", got.msg.rstat.stat.name); + try h.quiet(); + h.fsys.park_opens = false; + h.pump(); + got = try h.reap(); + try testing.expectEqual(@as(u16, 6), got.tag); + try testing.expect(got.msg == .ropen); + try h.send(9, .{ .tread = .{ .fid = 1, .offset = 0, .count = 64 } }); + got = try h.reap(); + try testing.expectEqualStrings("hello, body\n", got.msg.rread.data); + + // A truncating open parks at its truncate step and resumes there: the + // truncate still happens before the open. + _ = try h.walkTo(10, 3, &.{ "1", "body" }); + h.fsys.park_setattr = true; + try h.send(11, .{ .topen = .{ .fid = 3, .mode = owrite | c9.otrunc } }); + try h.quiet(); + try testing.expectEqualStrings("hello, body\n", h.fsys.body); + h.fsys.park_setattr = false; + h.pump(); + got = try h.reap(); + try testing.expectEqual(@as(u16, 11), got.tag); + try testing.expect(got.msg == .ropen); + try testing.expectEqualStrings("", h.fsys.body); + + // A clunk of an open fid parks at its release and pays it on retry. + const paid = h.fsys.releases; + h.fsys.park_releases = true; + try h.send(12, .{ .tclunk = .{ .fid = 1 } }); + try h.quiet(); + try testing.expectEqual(paid, h.fsys.releases); + h.fsys.park_releases = false; + h.pump(); + got = try h.reap(); + try testing.expectEqual(@as(u16, 12), got.tag); + try testing.expect(got.msg == .rclunk); + try testing.expectEqual(paid + 1, h.fsys.releases); + + // A walk keeps its names in the input frame, which parking would let go + // of, so it is refused rather than parked. + h.fsys.park_lookups = true; + try h.send(13, .{ .twalk = .{ .fid = 0, .newfid = 4, .nwname = 1, .wname = wnames(&.{"index"}) } }); + got = try h.reap(); + try testing.expectEqualStrings(e_again, got.msg.rerror.ename); + h.fsys.park_lookups = false; + + // A flush reaches a parked job the way it reaches a parked read. + h.fsys.park_opens = true; + try h.send(14, .{ .topen = .{ .fid = 2, .mode = oread } }); + try h.quiet(); + try h.send(15, .{ .tflush = .{ .oldtag = 14 } }); + got = try h.reap(); + try testing.expectEqual(@as(u16, 14), got.tag); + try testing.expectEqualStrings(e_interrupted, got.msg.rerror.ename); + got = try h.reap(); + try testing.expect(got.msg == .rflush); + h.fsys.park_opens = false; + h.pump(); + try h.quiet(); +} + test "fs server: the reply queue is a FIFO that survives a partial write" { var h: Harness = .{}; try h.handshake(4096); |
