summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--src/fs.zig176
1 files changed, 162 insertions, 14 deletions
diff --git a/src/fs.zig b/src/fs.zig
index 1641171..bd9cbd9 100644
--- a/src/fs.zig
+++ b/src/fs.zig
@@ -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);