summaryrefslogtreecommitdiff
path: root/9ns/src/bridge.zig
diff options
context:
space:
mode:
Diffstat (limited to '9ns/src/bridge.zig')
-rw-r--r--9ns/src/bridge.zig212
1 files changed, 199 insertions, 13 deletions
diff --git a/9ns/src/bridge.zig b/9ns/src/bridge.zig
index 3a072ec..089c3ca 100644
--- a/9ns/src/bridge.zig
+++ b/9ns/src/bridge.zig
@@ -4,6 +4,12 @@
//! Everything here is single-threaded and one request at a time. State is three
//! tables: inodes (nodeid → fid/qid, deduplicated by qid.path), open handles
//! (fh → fid plus a cached directory listing), and the reverse qid map.
+//!
+//! One request at a time does not mean deaf: while a 9P reply is outstanding
+//! the session polls the FUSE descriptor too (`nine.Interrupt`). A
+//! FUSE_INTERRUPT for the request being served becomes a Tflush, and if the
+//! server honours it the request fails with EINTR; anything else the kernel
+//! sends meanwhile is parked in a one-slot stash and served next.
const std = @import("std");
const cloud9 = @import("cloud9");
const fuse = @import("fuse.zig");
@@ -89,12 +95,22 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_
defer b.deinit();
b.req_buf = try gpa.alignedAlloc(u8, .@"8", request_buf_len);
+ b.spare_buf = try gpa.alignedAlloc(u8, .@"8", request_buf_len);
b.data_buf = try gpa.alloc(u8, max_write);
+ // Every read of the FUSE fd follows a poll; non-blocking makes sure a
+ // request the kernel withdrew in between cannot park us in read(2) while
+ // a 9P reply is due.
+ fuse.setNonblocking(fuse_fd) catch return error.FuseIo;
+
// Abandon any pending 9P reply once the child is gone (stop_fd readable),
// including the initial root stat below: a silent server must not pin us.
session.stop_fd = stop_fd;
defer session.stop_fd = -1;
+ // And watch the FUSE fd meanwhile: INTERRUPTs become Tflush, other
+ // requests (INIT arrives during the root stat) wait in the stash.
+ session.interrupt = b.interruptSource();
+ defer session.interrupt = null;
// Node 1 is the root; its qid comes from a stat so lookups resolving back to
// it (e.g. via a walk) dedupe onto node 1.
@@ -115,6 +131,17 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_
.{ .fd = stop_fd, .events = linux.POLL.IN, .revents = 0 },
};
while (true) {
+ // A request that arrived while a 9P reply was outstanding goes first.
+ // It lives in the spare buffer; swap so that the spare is free again
+ // for anything that arrives while this one is being served.
+ if (b.stash) |req| {
+ b.stash = null;
+ std.mem.swap([]align(8) u8, &b.req_buf, &b.spare_buf);
+ if (!try b.dispatch(req)) return;
+ continue;
+ }
+ if (b.fuse_gone) return;
+ if (b.fuse_fail) |e| return e;
pfds[0].revents = 0;
pfds[1].revents = 0;
const rc = linux.poll(&pfds, pfds.len, -1);
@@ -128,7 +155,8 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_
return;
}
if (pfds[0].revents == 0) continue;
- const req = (fuse.readRequest(fuse_fd, b.req_buf) catch |e| switch (e) {
+ const req = (fuse.readRequestOnce(fuse_fd, b.req_buf) catch |e| switch (e) {
+ error.Retry => continue,
error.Protocol => return error.FuseProtocol,
else => return error.FuseIo,
}) orelse {
@@ -145,7 +173,23 @@ const Bridge = struct {
nine: *nine.Session,
opts: Options,
req_buf: []align(8) u8 = &.{},
+ /// Second request buffer: what the interrupt poll reads into. Holds the
+ /// stashed request until `serve` swaps it in.
+ spare_buf: []align(8) u8 = &.{},
data_buf: []u8 = &.{},
+ /// A non-INTERRUPT request read while a 9P reply was outstanding (its body
+ /// points into `spare_buf`). While it is set the FUSE fd is not polled
+ /// during waits, so a second one cannot arrive.
+ stash: ?fuse.Request = null,
+ /// `unique` of the FUSE request being served, if any: the only one an
+ /// INTERRUPT may cancel.
+ cur_unique: ?u64 = null,
+ /// An INTERRUPT for `cur_unique` was consumed: chunked loops stop early
+ /// even when the flushed reply won the race. Reset per request.
+ interrupted: bool = false,
+ /// The FUSE fd reported ENODEV / a failure while the session was waiting.
+ fuse_gone: bool = false,
+ fuse_fail: ?error{ FuseIo, FuseProtocol } = null,
inodes: std.AutoHashMapUnmanaged(u64, Inode) = .empty,
by_qid: std.AutoHashMapUnmanaged(u64, u64) = .empty,
handles: std.AutoHashMapUnmanaged(u64, Handle) = .empty,
@@ -165,9 +209,65 @@ const Bridge = struct {
b.inodes.deinit(b.gpa);
b.by_qid.deinit(b.gpa);
if (b.req_buf.len != 0) b.gpa.free(b.req_buf);
+ if (b.spare_buf.len != 0) b.gpa.free(b.spare_buf);
if (b.data_buf.len != 0) b.gpa.free(b.data_buf);
}
+ // -- interrupt source (polled by nine.Session while a reply is outstanding) -----
+
+ fn interruptSource(b: *Bridge) nine.Interrupt {
+ return .{ .ctx = b, .watch = interruptWatch, .onReadable = interruptReadable, .armed = interruptArmed };
+ }
+
+ /// Poll the FUSE fd only while the stash has room: with it full a second
+ /// request would have nowhere to go.
+ fn interruptWatch(ctx: *anyopaque) i32 {
+ const b: *Bridge = @ptrCast(@alignCast(ctx));
+ return if (b.stash == null and !b.fuse_gone and b.fuse_fail == null) b.fuse_fd else -1;
+ }
+
+ /// Reads the request the kernel has ready. An INTERRUPT for the request in
+ /// flight asks the session to flush it; one for any other request is
+ /// dropped (the kernel expects no reply); anything else is stashed.
+ fn interruptReadable(ctx: *anyopaque) nine.Session.Error!bool {
+ const b: *Bridge = @ptrCast(@alignCast(ctx));
+ const req = (fuse.readRequestOnce(b.fuse_fd, b.spare_buf) catch |e| switch (e) {
+ error.Retry => return false,
+ error.Protocol => {
+ b.fuse_fail = error.FuseProtocol;
+ return error.Stopped;
+ },
+ else => {
+ b.fuse_fail = error.FuseIo;
+ return error.Stopped;
+ },
+ }) orelse {
+ b.trace("fuse fd reports ENODEV while a 9P reply is outstanding", .{});
+ b.fuse_gone = true;
+ return error.Stopped;
+ };
+ const h = req.header;
+ if (h.op() == .interrupt) {
+ const in = fuse.body(fuse.InterruptIn, req) catch return false;
+ const cur = b.cur_unique orelse std.math.maxInt(u64);
+ if (in.unique == cur) {
+ b.trace("<- interrupt for unique={d} (in flight): sending Tflush", .{in.unique});
+ b.interrupted = true;
+ return true;
+ }
+ b.trace("<- interrupt for unique={d} (not in flight; ignored)", .{in.unique});
+ return false;
+ }
+ b.trace("<- {s} unique={d} nodeid={d} stashed while a 9P reply is outstanding", .{ opName(h.op()), h.unique, h.nodeid });
+ b.stash = req;
+ return false;
+ }
+
+ fn interruptArmed(ctx: *anyopaque) bool {
+ const b: *Bridge = @ptrCast(@alignCast(ctx));
+ return b.interrupted;
+ }
+
fn trace(b: *const Bridge, comptime fmt: []const u8, args: anytype) void {
if (b.opts.debug) std.debug.print("9ns: " ++ fmt ++ "\n", args);
}
@@ -188,7 +288,12 @@ const Bridge = struct {
b.reply(h.unique, &.{}) catch {};
return false;
}
+ b.cur_unique = h.unique;
+ b.interrupted = false;
+ defer b.cur_unique = null;
b.handle(req) catch |e| {
+ if (e == error.Stopped and b.fuse_gone) return false;
+ if (e == error.Stopped and b.fuse_fail != null) return b.fuse_fail.?;
const code: linux.E = switch (e) {
error.Nine => b.last_err,
error.BadRequest => .INVAL,
@@ -201,6 +306,7 @@ const Bridge = struct {
error.TooLarge => .NAMETOOLONG,
error.BadDir => .IO,
error.Closed, error.Protocol, error.Io, error.Stopped => .IO,
+ error.Interrupted => .INTR,
error.FuseIo => return error.FuseIo,
};
if (wants_reply) try b.replyError(h.unique, code);
@@ -565,15 +671,15 @@ const Bridge = struct {
while (true) {
// A server that ignores the offset would otherwise feed us forever.
if (offset >= max_dir_bytes) return error.BadDir;
+ // A flushed read whose reply still won the race: the listing is
+ // incomplete either way, so stop here rather than read on.
+ if (b.interrupted) return error.Interrupted;
const n = try b.read(fid, offset, b.data_buf);
if (n == 0) break;
try parseDirRecords(b.gpa, b.data_buf[0..n], &list);
offset += n;
}
- // Entries carrying the root's own qid.path get the root's ino (1), as GETATTR would report it.
- for (list.entries.items[2..]) |*e| if (e.ino == b.root_path) {
- e.ino = fuse.root_id;
- };
+ for (list.entries.items[2..]) |*e| e.ino = inoFromPath(e.ino);
return list;
}
@@ -586,21 +692,29 @@ const Bridge = struct {
}
fn walkName(b: *Bridge, fid: u32, name: []const u8) nine.Session.Error!u32 {
- const newfid = b.nine.allocFid();
- _ = b.nine.walk(fid, newfid, &.{name}) catch |e| {
- b.nine.freeFid(newfid);
- return b.nineErr("walk", fid, e);
- };
+ const newfid = try b.walkTo(fid, &.{name});
b.trace(" 9p walk fid={d} newfid={d} name={s} -> ok", .{ fid, newfid, name });
return newfid;
}
fn clone(b: *Bridge, fid: u32) nine.Session.Error!u32 {
- const newfid = b.nine.clone(fid) catch |e| return b.nineErr("clone", fid, e);
+ const newfid = try b.walkTo(fid, &.{});
b.trace(" 9p walk fid={d} newfid={d} (clone) -> ok", .{ fid, newfid });
return newfid;
}
+ /// allocFid + walk. On Rerror the new fid was never bound; after an
+ /// interruption the server may or may not have bound it (the Rflush
+ /// tells us only that no reply is coming), so it is clunked to be sure.
+ fn walkTo(b: *Bridge, fid: u32, names: []const []const u8) nine.Session.Error!u32 {
+ const newfid = b.nine.allocFid();
+ _ = b.nine.walk(fid, newfid, names) catch |e| {
+ if (e == error.Interrupted) b.clunkQuiet(newfid) else b.nine.freeFid(newfid);
+ return b.nineErr("walk", fid, e);
+ };
+ return newfid;
+ }
+
fn open9(b: *Bridge, fid: u32, mode: u8) nine.Session.Error!nine.Session.Open {
const o = b.nine.open(fid, mode) catch |e| return b.nineErr("open", fid, e);
b.trace(" 9p open fid={d} mode={d} -> iounit={d}", .{ fid, mode, o.iounit });
@@ -674,12 +788,18 @@ const Bridge = struct {
// -- attrs ---------------------------------------------------------------------------
+ /// The inode number reported to the kernel is the 9P qid.path, for the root
+ /// too: FUSE only needs the root's *nodeid* to be 1, and a server may hand
+ /// qid.path 1 to some other file (Pardes gives it to /self), which would
+ /// otherwise make `find` see a directory cycle. qid.path 0 maps to a
+ /// sentinel because inode 0 is treated as invalid by much of userland.
fn inoOf(b: *const Bridge, nodeid: u64, qid: cloud9.Qid) u64 {
- return if (nodeid == fuse.root_id or qid.path == b.root_path) fuse.root_id else qid.path;
+ _ = b;
+ _ = nodeid;
+ return inoFromPath(qid.path);
}
fn inoOfNode(b: *const Bridge, nodeid: u64) u64 {
- if (nodeid == fuse.root_id) return fuse.root_id;
const ino = b.inodes.get(nodeid) orelse return nodeid;
return b.inoOf(nodeid, ino.qid);
}
@@ -700,6 +820,11 @@ const Bridge = struct {
// -- pure helpers (unit-tested) ------------------------------------------------------------
/// Attr from a 9P Stat: DMDIR → S_IFDIR else S_IFREG, low 9 permission bits kept.
+/// qid.path → inode number; 0 becomes a sentinel (inode 0 reads as "invalid" to many tools).
+pub fn inoFromPath(path: u64) u64 {
+ return if (path == 0) std.math.maxInt(u64) - 1 else path;
+}
+
pub fn attrFromStat(st: cloud9.Stat, ino: u64, uid: u32, gid: u32) fuse.Attr {
const ftype: u32 = if (st.mode & cloud9.dmdir != 0) fuse.S_IFDIR else fuse.S_IFREG;
return .{
@@ -964,6 +1089,67 @@ test "dirent names the kernel would reject are dropped from listings" {
try testing.expectEqualStrings("also", list.entries.items[1].name);
}
+test "interrupt source: INTERRUPT in flight → Tflush, other INTERRUPTs ignored, requests stashed" {
+ var pi = try nine.PipeInterrupt.init(0); // only its fake FUSE fd and inject() are used
+ defer pi.deinit();
+ var fs: nine.FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush, .on_hang_inject = &pi };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ var b: Bridge = .{ .gpa = testing.allocator, .fuse_fd = pi.read_end, .nine = s, .opts = .{ .uid = 0, .gid = 0 } };
+ defer b.deinit();
+ b.spare_buf = try testing.allocator.alignedAlloc(u8, .@"8", request_buf_len);
+ s.interrupt = b.interruptSource();
+ defer s.interrupt = null;
+ try testing.expectEqual(pi.read_end, Bridge.interruptWatch(&b));
+
+ // Serving unique 7. An INTERRUPT for 6 is already queued (ignored); the
+ // server fires the one for 7 (pi.unique) once the read at offset 0 hangs.
+ var buf: [100]u8 = undefined;
+ b.cur_unique = 7;
+ pi.unique = 7;
+ try pi.inject(6);
+ try testing.expectError(error.Interrupted, b.read(1, 0, &buf));
+ try testing.expect(b.interrupted);
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ try testing.expectEqual(fs.hung_tag.load(.seq_cst), fs.flush_oldtag.load(.seq_cst));
+ try testing.expect(b.stash == null);
+ // The session is intact: a clunk-style cleanup rpc and a further read work.
+ b.interrupted = false;
+ b.cur_unique = 8;
+ try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40]));
+ try testing.expectEqual(@as(usize, 0), s.client.pending());
+
+ // A FORGET arriving during a wait is stashed, and the fd is then not watched.
+ var wire: [48]u8 = undefined;
+ const hdr = fuse.InHeader{ .len = 48, .opcode = @intFromEnum(fuse.Opcode.forget), .unique = 99, .nodeid = 5, .uid = 0, .gid = 0, .pid = 0, .total_extlen = 0, .padding = 0 };
+ @memcpy(wire[0..40], std.mem.asBytes(&hdr));
+ @memcpy(wire[40..48], std.mem.asBytes(&fuse.ForgetIn{ .nlookup = 1 }));
+ try testing.expectEqual(@as(usize, 48), linux.write(pi.write_end, &wire, wire.len));
+ fs.read_delay_ns = 30 * std.time.ns_per_ms;
+ b.cur_unique = 9;
+ try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40]));
+ try testing.expect(!b.interrupted);
+ const stashed = b.stash orelse return error.TestUnexpectedResult;
+ try testing.expectEqual(fuse.Opcode.forget, stashed.header.op());
+ try testing.expectEqual(@as(u64, 5), stashed.header.nodeid);
+ try testing.expectEqual(@as(u64, 1), (try fuse.body(fuse.ForgetIn, stashed)).nlookup);
+ try testing.expectEqual(@as(i32, -1), Bridge.interruptWatch(&b));
+ // With the stash full an INTERRUPT is not even looked at.
+ try pi.inject(9);
+ try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40]));
+ try testing.expect(!b.interrupted);
+ b.stash = null;
+ try testing.expectEqual(pi.read_end, Bridge.interruptWatch(&b));
+ // Once the stash is served the queued INTERRUPT is consumed (and ignored:
+ // its request is not the one in flight any more).
+ b.cur_unique = 10;
+ fs.read_delay_ns = 30 * std.time.ns_per_ms;
+ try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40]));
+ try testing.expect(!b.interrupted);
+ try testing.expect(b.stash == null);
+}
+
test "DirList frees its names" {
var list: DirList = .{};
try list.entries.append(testing.allocator, .{ .name = try testing.allocator.dupe(u8, "x"), .ino = 1, .dtype = fuse.DT_REG });