From 3e9f8805f293f622bb885cf849b5ce47dc062ad1 Mon Sep 17 00:00:00 2001 From: Gabriel Schneider Date: Sat, 19 Sep 2026 23:55:47 -0300 Subject: 9ns: --name and /mnt/9p/ mounts, qid.path as inode number, interrupts as Tflush - --name NAME (default derived from the transport: socket basename, tcp-IP-PORT, spawned command, fdN) mounts at /mnt/9p/; --mount still overrides. ensureMountpoint walks down and creates missing components, shadowing the deepest unwritable ancestor. NINE_MOUNT is the only exported variable. - The inode number reported to the kernel is the 9P qid.path for every node, root included; a server handing qid.path 1 to a file (Pardes /self) no longer collides with the root. - FUSE_INTERRUPT for the request in flight becomes Tflush; a blocked read returns EINTR when the server answers the flush, chunked transfers return short counts, other requests arriving meanwhile are stashed and served next. Servers ignoring Tflush still block until they answer. - 9ns-test now covers nine/bridge/fuse; new adv_bridge_interrupt suite (28); 9ns-itest grows to 88 checks. Co-Authored-By: Claude Fable 5.1 --- 9ns/src/bridge.zig | 212 +++++++++++++++++++++++++++++++++++++++++++++++++---- 1 file changed, 199 insertions(+), 13 deletions(-) (limited to '9ns/src/bridge.zig') 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 }); -- cgit v1.3