diff options
Diffstat (limited to '9ns/src/bridge.zig')
| -rw-r--r-- | 9ns/src/bridge.zig | 212 |
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 }); |
