//! FUSE ↔ 9P2000 translation: the request loop that turns kernel FUSE requests //! into synchronous 9P calls on a `nine.Session` and sends the replies back. //! //! 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 post = cloud9.post; const fuse = @import("fuse.zig"); const nine = @import("nine.zig"); const linux = std.os.linux; pub const Options = struct { /// Reported owner of every file. uid: u32, gid: u32, /// attr/entry cache validity (0 = none). attr_timeout_ns: u64 = 1_000_000_000, /// FOPEN_DIRECT_IO on every regular file. direct_io: bool = true, /// Trace every request, reply and 9P call to stderr. debug: bool = false, }; /// Largest single READ/WRITE payload we accept from the kernel. pub const max_write: u32 = 1 << 20; /// Upper bound on the raw bytes of one directory listing (about a million entries); /// past it the listing fails with EIO instead of eating memory. pub const max_dir_bytes: u64 = 64 << 20; /// The kernel refuses dirents longer than this (FUSE_NAME_MAX) with EIO. pub const max_name_len: usize = 1024; /// Request buffer: `max_write` plus room for the header and the largest in-struct. pub const request_buf_len: usize = max_write + 4096; const Inode = struct { fid: u32, qid: cloud9.Qid, nlookup: u64, /// nodeid of the directory this inode was looked up in (root: itself). Used for "..". parent: u64, }; pub const Entry = struct { name: []u8, ino: u64, dtype: u32 }; pub const DirList = struct { entries: std.ArrayList(Entry) = .empty, pub fn deinit(d: *DirList, gpa: std.mem.Allocator) void { for (d.entries.items) |e| gpa.free(e.name); d.entries.deinit(gpa); } }; const Handle = struct { fid: u32, nodeid: u64, dir: ?DirList, }; /// Errors a request handler may surface. Policy failures are ordinary errors /// that the dispatcher maps to an errno; `FuseIo` means the kernel side is broken. const HandlerError = nine.Session.Error || error{ BadRequest, NoEntry, BadHandle, Exdev, Perm, NotSup, /// A directory listing the server sent could not be parsed (EIO, not fatal). BadDir, FuseIo, }; /// Runs until the FUSE fd reports ENODEV, a DESTROY arrives, or `stop_fd` /// becomes readable (also while a 9P reply is outstanding). Returns /// `error.Closed` if the 9P server went away. pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_fid: u32, stop_fd: i32, opts: Options) !void { var effective = opts; // With page caching on, a nonzero attr cache lets the kernel trust a stale // (often zero) size and truncate reads: 9P sizes are authoritative and change // under us. direct_io ignores the cached size, so the cache is safe only there. if (!effective.direct_io) effective.attr_timeout_ns = 0; var b: Bridge = .{ .gpa = gpa, .fuse_fd = fuse_fd, .nine = session, .opts = effective, }; 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. var root_qid: cloud9.Qid = .{ .type = cloud9.qtdir, .version = 0, .path = 0 }; if (b.stat(root_fid)) |st| { root_qid = st.qid; b.root_path = st.qid.path; try b.by_qid.put(gpa, st.qid.path, b.root_id); } else |e| switch (e) { error.Nine => {}, error.Stopped => return, else => return error.Closed, } try b.inodes.put(gpa, b.root_id, .{ .fid = root_fid, .qid = root_qid, .nlookup = 1, .parent = b.root_id }); var pfds = [_]linux.pollfd{ .{ .fd = fuse_fd, .events = linux.POLL.IN, .revents = 0 }, .{ .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); switch (linux.errno(rc)) { .SUCCESS => {}, .INTR, .AGAIN => continue, else => return error.Io, } if (pfds[1].revents != 0) { b.trace("stop_fd readable; leaving serve loop", .{}); return; } if (pfds[0].revents == 0) continue; 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 { b.trace("fuse fd reports ENODEV; unmounted", .{}); return; }; if (!try b.dispatch(req)) return; } } const Bridge = struct { gpa: std.mem.Allocator, fuse_fd: i32, nine: *nine.Session, opts: Options, /// Node id of this bridge's root inode. `fuse.root_id` (1) in /// single-connection mode; in mntgen mode the mount's `root_node`, /// which carries the mount index in the top bits (see `serveMntgen`). root_id: u64 = fuse.root_id, /// Mixed into reported inode numbers so two servers that hand out the /// same qid.path (two ramfs instances, say) cannot share an inode number /// inside one mntgen mount. 0 in single-connection mode. ino_xor: u64 = 0, 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 (0 = none): the only /// one an INTERRUPT may cancel. Atomic: in mntgen mode the dispatcher /// thread scans it to route FUSE_INTERRUPTs to the right server. cur_unique: std.atomic.Value(u64) = .init(0), /// 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, next_node: u64 = 2, next_fh: u64 = 1, /// qid.path of the root, reported as ino 1 wherever it shows up. root_path: u64 = 0, /// errno of the most recent Rerror that a handler did not swallow. Kept here /// because `Session.rpc` clears its ename on every call, and error paths /// clunk (an rpc) before the dispatcher maps the failure to an errno. last_err: linux.E = .IO, fn deinit(b: *Bridge) void { var it = b.handles.valueIterator(); while (it.next()) |h| if (h.dir) |*d| d.deinit(b.gpa); b.handles.deinit(b.gpa); 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.load(.seq_cst); if (cur != 0 and 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); } // -- dispatch -------------------------------------------------------------------- /// Handles one request. Returns false when the loop should stop (DESTROY). /// Fatal errors (dead 9P session, broken FUSE fd) propagate. fn dispatch(b: *Bridge, req: fuse.Request) !bool { const h = req.header; const op = h.op(); b.trace("<- {s} unique={d} nodeid={d} len={d} (fids={d} inodes={d} handles={d})", .{ opName(op), h.unique, h.nodeid, h.len, b.nine.fidsInUse(), b.inodes.count(), b.handles.count() }); const wants_reply = switch (op) { .forget, .batch_forget, .interrupt => false, else => true, }; if (op == .destroy) { b.reply(h.unique, &.{}) catch {}; return false; } b.cur_unique.store(h.unique, .seq_cst); b.interrupted = false; defer b.cur_unique.store(0, .seq_cst); 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, error.NoEntry => .NOENT, error.BadHandle => .BADF, error.Exdev => .XDEV, error.Perm => .PERM, error.NotSup => .NOSYS, error.OutOfMemory => .NOMEM, 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); switch (e) { error.Closed, error.Protocol, error.Io => return error.Closed, error.Stopped => return false, // the child is gone; the mount is being torn down else => {}, } }; return true; } fn handle(b: *Bridge, req: fuse.Request) HandlerError!void { const u = req.header.unique; switch (req.header.op()) { .init => { const in = try body(fuse.InitIn, req); const out = fuse.initReply(in, max_write); try b.reply(u, &.{std.mem.asBytes(&out)}); }, .lookup => { const name = try nameAfter(void, req); const entry = try b.lookupEntry(req.header.nodeid, name); try b.reply(u, &.{std.mem.asBytes(&entry)}); }, .forget => { const in = try body(fuse.ForgetIn, req); try b.forget(req.header.nodeid, in.nlookup); }, .batch_forget => { const in = try body(fuse.BatchForgetIn, req); const rest = req.body[@sizeOf(fuse.BatchForgetIn)..]; const count: usize = in.count; if (rest.len < count * @sizeOf(fuse.ForgetOne)) return error.BadRequest; for (0..count) |i| { const one = std.mem.bytesToValue(fuse.ForgetOne, rest[i * @sizeOf(fuse.ForgetOne) ..][0..@sizeOf(fuse.ForgetOne)]); try b.forget(one.nodeid, one.nlookup); } }, .getattr => { const ino = b.inodes.get(req.header.nodeid) orelse return error.NoEntry; const st = try b.stat(ino.fid); const out = b.attrOut(st, b.inoOf(req.header.nodeid, ino.qid)); try b.reply(u, &.{std.mem.asBytes(&out)}); }, .setattr => try b.setattr(req), .open => try b.openFile(req, false), .opendir => try b.openFile(req, true), .read => { const in = try body(fuse.ReadIn, req); const h = b.handles.get(in.fh) orelse return error.BadHandle; const want: usize = @min(in.size, max_write); const n = try b.read(h.fid, in.offset, b.data_buf[0..want]); try b.reply(u, &.{b.data_buf[0..n]}); }, .write => { const in = try body(fuse.WriteIn, req); const h = b.handles.get(in.fh) orelse return error.BadHandle; const rest = req.body[@sizeOf(fuse.WriteIn)..]; if (rest.len < in.size) return error.BadRequest; const n = try b.write(h.fid, in.offset, rest[0..in.size]); const out = fuse.WriteOut{ .size = @intCast(n) }; try b.reply(u, &.{std.mem.asBytes(&out)}); }, .readdir => try b.readdir(req), .release, .releasedir => { const in = try body(fuse.ReleaseIn, req); const kv = b.handles.fetchRemove(in.fh) orelse return error.BadHandle; var h = kv.value; if (h.dir) |*d| d.deinit(b.gpa); try b.clunk(h.fid); try b.reply(u, &.{}); }, .flush, .fsync, .fsyncdir => try b.reply(u, &.{}), .create => try b.create(req), .mkdir => { const in = try body(fuse.MkdirIn, req); const name = try nameAfter(fuse.MkdirIn, req); const parent = b.inodes.get(req.header.nodeid) orelse return error.NoEntry; const fid = try b.clone(parent.fid); _ = b.create9(fid, name, cloud9.dmdir | (in.mode & 0o777), cloud9.oread) catch |e| { b.clunkQuiet(fid); return e; }; try b.clunk(fid); const entry = try b.lookupEntry(req.header.nodeid, name); try b.reply(u, &.{std.mem.asBytes(&entry)}); }, .unlink, .rmdir => { const name = try nameAfter(void, req); const parent = b.inodes.get(req.header.nodeid) orelse return error.NoEntry; const tmp = try b.walkName(parent.fid, name); try b.remove(tmp); try b.reply(u, &.{}); }, .rename => { const in = try body(fuse.RenameIn, req); const old = try nameAfter(fuse.RenameIn, req); const new = try secondName(req, old, @sizeOf(fuse.RenameIn)); try b.rename(req.header.nodeid, in.newdir, old, new, 0); try b.reply(u, &.{}); }, .rename2 => { const in = try body(fuse.Rename2In, req); const old = try nameAfter(fuse.Rename2In, req); const new = try secondName(req, old, @sizeOf(fuse.Rename2In)); try b.rename(req.header.nodeid, in.newdir, old, new, in.flags); try b.reply(u, &.{}); }, .statfs => { const out = fuse.StatfsOut{ .st = .{ .bsize = 4096, .namelen = 255, .frsize = 4096 } }; try b.reply(u, &.{std.mem.asBytes(&out)}); }, .interrupt => {}, .destroy => unreachable, // handled in dispatch .access => return error.NotSup, else => return error.NotSup, } } // -- handlers ---------------------------------------------------------------------- /// walk(parent → new fid, [name]) + stat, deduplicated by qid.path. Bumps nlookup. fn lookupEntry(b: *Bridge, parent_id: u64, name: []const u8) HandlerError!fuse.EntryOut { const parent = b.inodes.get(parent_id) orelse return error.NoEntry; const newfid = try b.walkName(parent.fid, name); const st = b.stat(newfid) catch |e| { b.clunkQuiet(newfid); return e; }; const qid = st.qid; var nodeid: u64 = undefined; if (b.by_qid.get(qid.path)) |existing| { // A directory and a file sharing a qid.path (a server bug) must not // share a node: the kernel would mark the inode bad, and for the // root that is fatal for the whole mount. const merge = if (b.inodes.getPtr(existing)) |ino| (ino.qid.type & cloud9.qtdir) == (qid.type & cloud9.qtdir) else false; if (merge) { const ino = b.inodes.getPtr(existing).?; ino.nlookup += 1; ino.qid = qid; nodeid = existing; if (existing == b.root_id) { b.clunkQuiet(newfid); } else { // Keep the fresh fid (it is bound to the current file at this // name) and retire the older one. const stale = ino.fid; ino.fid = newfid; b.clunkQuiet(stale); } } else { // Stale reverse entry, or a type clash: bind a fresh node to it. nodeid = try b.newInode(newfid, qid, parent_id); } } else { nodeid = try b.newInode(newfid, qid, parent_id); } var out = fuse.EntryOut{ .nodeid = nodeid, .generation = 0, .attr = b.attrFrom(st, b.inoOf(nodeid, qid)), }; out.entry_valid = b.opts.attr_timeout_ns / 1_000_000_000; out.entry_valid_nsec = @intCast(b.opts.attr_timeout_ns % 1_000_000_000); out.attr_valid = out.entry_valid; out.attr_valid_nsec = out.entry_valid_nsec; return out; } fn newInode(b: *Bridge, fid: u32, qid: cloud9.Qid, parent: u64) HandlerError!u64 { const nodeid = b.next_node; b.inodes.put(b.gpa, nodeid, .{ .fid = fid, .qid = qid, .nlookup = 1, .parent = parent }) catch |e| { b.clunkQuiet(fid); return e; }; b.by_qid.put(b.gpa, qid.path, nodeid) catch |e| { _ = b.inodes.remove(nodeid); b.clunkQuiet(fid); return e; }; b.next_node += 1; return nodeid; } fn forget(b: *Bridge, nodeid: u64, n: u64) HandlerError!void { if (nodeid == b.root_id) return; const ino = b.inodes.getPtr(nodeid) orelse return; if (ino.nlookup > n) { ino.nlookup -= n; return; } const fid = ino.fid; const path = ino.qid.path; _ = b.inodes.remove(nodeid); if (b.by_qid.get(path)) |mapped| { if (mapped == nodeid) _ = b.by_qid.remove(path); } b.clunk(fid) catch |e| switch (e) { error.Nine => {}, else => return e, }; } fn setattr(b: *Bridge, req: fuse.Request) HandlerError!void { const in = try body(fuse.SetattrIn, req); const ino = b.inodes.get(req.header.nodeid) orelse return error.NoEntry; const old = try b.stat(ino.fid); const old_mode = old.mode; var st = nine.dontcare; var changed = false; if (in.valid & fuse.FATTR_UID != 0 and in.uid != b.opts.uid) return error.Perm; if (in.valid & fuse.FATTR_GID != 0 and in.gid != b.opts.gid) return error.Perm; if (in.valid & fuse.FATTR_SIZE != 0) { st.length = in.size; changed = true; } if (in.valid & fuse.FATTR_MODE != 0) { st.mode = (old_mode & ~@as(u32, 0o777)) | (in.mode & 0o777); changed = true; } if (in.valid & fuse.FATTR_MTIME_NOW != 0) { st.mtime = nowSeconds(); changed = true; } else if (in.valid & fuse.FATTR_MTIME != 0) { st.mtime = @truncate(in.mtime); changed = true; } if (changed) try b.wstat(ino.fid, st); const fresh = try b.stat(ino.fid); const out = b.attrOut(fresh, b.inoOf(req.header.nodeid, ino.qid)); try b.reply(req.header.unique, &.{std.mem.asBytes(&out)}); } fn openFile(b: *Bridge, req: fuse.Request, is_dir: bool) HandlerError!void { const in = try body(fuse.OpenIn, req); const ino = b.inodes.get(req.header.nodeid) orelse return error.NoEntry; const mode: u8 = if (is_dir) cloud9.oread else openMode(in.flags); const fid = try b.clone(ino.fid); _ = b.open9(fid, mode) catch |e| { b.clunkQuiet(fid); return e; }; const fh = try b.newHandle(fid, req.header.nodeid); const out = fuse.OpenOut{ .fh = fh, .open_flags = if (!is_dir and b.opts.direct_io) fuse.FOPEN_DIRECT_IO else 0, }; try b.reply(req.header.unique, &.{std.mem.asBytes(&out)}); } fn newHandle(b: *Bridge, fid: u32, nodeid: u64) HandlerError!u64 { const fh = b.next_fh; b.handles.put(b.gpa, fh, .{ .fid = fid, .nodeid = nodeid, .dir = null }) catch |e| { b.clunkQuiet(fid); return e; }; b.next_fh += 1; return fh; } fn create(b: *Bridge, req: fuse.Request) HandlerError!void { const in = try body(fuse.CreateIn, req); const name = try nameAfter(fuse.CreateIn, req); const parent = b.inodes.get(req.header.nodeid) orelse return error.NoEntry; // The created fid becomes the open file. const fid = try b.clone(parent.fid); _ = b.create9(fid, name, in.mode & 0o777, openMode(in.flags)) catch |e| { b.clunkQuiet(fid); return e; }; const entry = b.lookupEntry(req.header.nodeid, name) catch |e| { b.clunkQuiet(fid); return e; }; const fh = try b.newHandle(fid, entry.nodeid); const oo = fuse.OpenOut{ .fh = fh, .open_flags = if (b.opts.direct_io) fuse.FOPEN_DIRECT_IO else 0, }; try b.reply(req.header.unique, &.{ std.mem.asBytes(&entry), std.mem.asBytes(&oo) }); } fn rename(b: *Bridge, parent_id: u64, newdir: u64, old: []const u8, new: []const u8, flags: u32) HandlerError!void { if (newdir != parent_id) return error.Exdev; const rf: linux.RENAME = @bitCast(flags); if (rf.EXCHANGE or rf.WHITEOUT) return error.BadRequest; const parent = b.inodes.get(parent_id) orelse return error.NoEntry; const tmp = try b.walkName(parent.fid, old); defer b.clunkQuiet(tmp); var st = nine.dontcare; st.name = new; b.wstat(tmp, st) catch |e| { // 9P2000 rename never replaces an existing name; POSIX rename does. if (e != error.Nine or rf.NOREPLACE or b.nine.errno() != .EXIST) return e; try b.renameOver(parent.fid, tmp, new); }; } /// Replace `new` with the file behind `src`. An (empty) directory target is /// removed first: it holds no data and the VFS already ruled out mismatched /// types. A file target is parked under a temporary name so that a failing /// second rename can put it back instead of having destroyed it. fn renameOver(b: *Bridge, parent_fid: u32, src: u32, new: []const u8) HandlerError!void { const victim = try b.walkName(parent_fid, new); const vst = b.stat(victim) catch |e| { b.clunkQuiet(victim); return e; }; var st = nine.dontcare; st.name = new; if (vst.mode & cloud9.dmdir != 0) { b.trace(" rename target is a directory; removing it and retrying", .{}); try b.remove(victim); return b.wstat(src, st); } var park_buf: [48]u8 = undefined; const park = std.fmt.bufPrint(&park_buf, ".9ns-rename-{x}", .{randomU64()}) catch unreachable; b.trace(" rename target exists; parking it as {s} and retrying", .{park}); var pst = nine.dontcare; pst.name = park; b.wstat(victim, pst) catch |e| { b.clunkQuiet(victim); return e; }; b.wstat(src, st) catch |e| { b.trace(" rename still failed; restoring the target", .{}); const saved = b.last_err; b.wstat(victim, st) catch {}; b.last_err = saved; b.clunkQuiet(victim); return e; }; b.remove(victim) catch b.trace(" could not remove the parked target {s}", .{park}); } fn readdir(b: *Bridge, req: fuse.Request) HandlerError!void { const in = try body(fuse.ReadIn, req); const h = b.handles.getPtr(in.fh) orelse return error.BadHandle; if (in.offset == 0 or h.dir == null) { if (h.dir) |*d| d.deinit(b.gpa); h.dir = null; h.dir = try b.loadDir(h.fid, h.nodeid); } const dir = &h.dir.?; const size: usize = @min(in.size, max_write); const used = packDirents(dir.entries.items, in.offset, b.data_buf[0..size]); try b.reply(req.header.unique, &.{b.data_buf[0..used]}); } /// Reads the whole directory and builds its listing, "." and ".." first. fn loadDir(b: *Bridge, fid: u32, nodeid: u64) HandlerError!DirList { var list: DirList = .{}; errdefer list.deinit(b.gpa); const self_ino = b.inoOfNode(nodeid); const parent_ino = if (b.inodes.get(nodeid)) |ino| b.inoOfNode(ino.parent) else self_ino; try list.entries.append(b.gpa, .{ .name = try b.gpa.dupe(u8, "."), .ino = self_ino, .dtype = fuse.DT_DIR }); try list.entries.append(b.gpa, .{ .name = try b.gpa.dupe(u8, ".."), .ino = parent_ino, .dtype = fuse.DT_DIR }); var offset: u64 = 0; 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; } for (list.entries.items[2..]) |*e| e.ino = inoFromPath(e.ino); return list; } // -- 9P wrappers (tracing) ----------------------------------------------------------- fn stat(b: *Bridge, fid: u32) nine.Session.Error!cloud9.Stat { const st = b.nine.stat(fid) catch |e| return b.nineErr("stat", fid, e); b.trace(" 9p stat fid={d} -> name={s} mode={o} len={d} qid={x}", .{ fid, st.name, st.mode, st.length, st.qid.path }); return st; } fn walkName(b: *Bridge, fid: u32, name: []const u8) nine.Session.Error!u32 { 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 = 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 }); return o; } fn create9(b: *Bridge, fid: u32, name: []const u8, perm: u32, mode: u8) nine.Session.Error!nine.Session.Open { const o = b.nine.create(fid, name, perm, mode) catch |e| return b.nineErr("create", fid, e); b.trace(" 9p create fid={d} name={s} perm={o} mode={d} -> iounit={d}", .{ fid, name, perm, mode, o.iounit }); return o; } fn read(b: *Bridge, fid: u32, offset: u64, buf: []u8) nine.Session.Error!usize { const n = b.nine.read(fid, offset, buf) catch |e| return b.nineErr("read", fid, e); b.trace(" 9p read fid={d} offset={d} count={d} -> {d}", .{ fid, offset, buf.len, n }); return n; } fn write(b: *Bridge, fid: u32, offset: u64, data: []const u8) nine.Session.Error!usize { const n = b.nine.write(fid, offset, data) catch |e| return b.nineErr("write", fid, e); b.trace(" 9p write fid={d} offset={d} count={d} -> {d}", .{ fid, offset, data.len, n }); return n; } fn wstat(b: *Bridge, fid: u32, st: cloud9.Stat) nine.Session.Error!void { b.nine.wstat(fid, st) catch |e| return b.nineErr("wstat", fid, e); b.trace(" 9p wstat fid={d} name={s} mode={x} len={x} mtime={x} -> ok", .{ fid, st.name, st.mode, st.length, st.mtime }); } fn clunk(b: *Bridge, fid: u32) nine.Session.Error!void { b.nine.clunk(fid) catch |e| return b.nineErr("clunk", fid, e); b.trace(" 9p clunk fid={d} -> ok", .{fid}); } /// Best-effort clunk during error unwinding; a dead session surfaces on the /// next call. Does not disturb the errno of the failure being unwound. fn clunkQuiet(b: *Bridge, fid: u32) void { const saved = b.last_err; defer b.last_err = saved; b.clunk(fid) catch {}; } fn remove(b: *Bridge, fid: u32) nine.Session.Error!void { b.nine.remove(fid) catch |e| return b.nineErr("remove", fid, e); b.trace(" 9p remove fid={d} -> ok", .{fid}); } fn nineErr(b: *Bridge, what: []const u8, fid: u32, e: nine.Session.Error) nine.Session.Error { if (e == error.Nine) { b.last_err = b.nine.errno(); b.trace(" 9p {s} fid={d} -> Rerror \"{s}\" ({s})", .{ what, fid, b.nine.ename[0..b.nine.ename_len], @tagName(b.nine.errno()) }); } else { b.trace(" 9p {s} fid={d} -> {s}", .{ what, fid, @errorName(e) }); } return e; } // -- FUSE wrappers (tracing) -------------------------------------------------------- fn reply(b: *Bridge, unique: u64, payloads: []const []const u8) error{FuseIo}!void { var total: usize = 0; for (payloads) |p| total += p.len; b.trace("-> unique={d} ok ({d} bytes)", .{ unique, total }); fuse.reply(b.fuse_fd, unique, payloads) catch return error.FuseIo; } fn replyError(b: *Bridge, unique: u64, code: linux.E) error{FuseIo}!void { b.trace("-> unique={d} error E{s}", .{ unique, @tagName(code) }); fuse.replyError(b.fuse_fd, unique, code) catch return error.FuseIo; } // -- 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 { _ = nodeid; return inoFromPath(qid.path) ^ b.ino_xor; } fn inoOfNode(b: *const Bridge, nodeid: u64) u64 { const ino = b.inodes.get(nodeid) orelse return nodeid; return b.inoOf(nodeid, ino.qid); } fn attrFrom(b: *const Bridge, st: cloud9.Stat, ino: u64) fuse.Attr { return attrFromStat(st, ino, b.opts.uid, b.opts.gid); } fn attrOut(b: *const Bridge, st: cloud9.Stat, ino: u64) fuse.AttrOut { return .{ .attr_valid = b.opts.attr_timeout_ns / 1_000_000_000, .attr_valid_nsec = @intCast(b.opts.attr_timeout_ns % 1_000_000_000), .attr = b.attrFrom(st, ino), }; } }; // =========================================================================== // mntgen: many servers behind one FUSE mount // =========================================================================== // // `9ns --mntgen` serves one FUSE mount whose synthetic root lists the posted // 9P services in `$XDG_RUNTIME_DIR/9p` (cloud9.post's registry; no connection // is made to list). A walk into a name dials that server lazily and starts a // per-server worker thread running the ordinary bridge translation above. // // Node ids carry the server in the top bits: a request for `(index, local)` // is routed by `index` to that server's mount. Indexes are ordinals, never // reused, so a stale kernel-side inode of a dead server can never be // conflated with a fresh inode of its replacement. A dead server answers EIO // on its whole subtree until the next walk into its name re-dials it (a new // mount, a new index); nothing reconnects eagerly. /// Bit position of the mount index inside a FUSE node id; the low bits are /// one server's bridge node ids, the top bits name the server (0 = the /// synthetic root itself). pub const mount_shift: u6 = 32; /// Per-server node ids live below 2^32 (a bridge never reuses one). pub const mount_node_mask: u64 = (1 << mount_shift) - 1; /// Mount indexes are ordinals and are never reused; this bounds how many /// distinct dials one 9ns process serves in its lifetime. pub const max_mounts: usize = 4096; /// Concurrent OPENDIRs of the synthetic root (each snapshots the registry). pub const max_root_dirs: usize = 64; /// The staged buffer handed to `post.posted` for one root listing. pub const stage_len: usize = 8192; /// Registry entries that are directories are served like the root itself /// (a synthetic directory mirroring the real one, dialing the sockets /// found inside); this bounds the synthetic directories one 9ns serves. pub const max_synth_dirs: usize = 64; /// How deep those registry subdirectories nest. pub const max_synth_depth: u8 = 8; /// The node-id index reserved for synthetic registry subdirectories; /// mounts use ordinals below 4096, so this never collides with one. pub const synth_index: u32 = 0xFFFF_FFFF; /// The node id a server's subtree lives under: `index` in the top bits, /// `local` (1 = that server's 9P root) below. pub fn mountNode(index: u32, local: u64) u64 { return @as(u64, index) << mount_shift | local; } /// The server a kernel request's node id belongs to. pub fn mountIndex(nodeid: u64) u32 { return @intCast(nodeid >> mount_shift); } /// Deterministic inode number for a synthetic-root entry (a posted name): /// FNV-1a of the name, so readdir inos are stable across calls. pub fn nameIno(name: []const u8) u64 { var h: u64 = 0xcbf29ce484222325; for (name) |c| { h ^= c; h *%= 0x100000001b3; } return h; } pub const MntgenOptions = struct { /// Used only on the dispatcher (main) thread: registry listing and /// dial. The worker threads never call io (their locks use the /// uncancelable futex paths, which are thread-safe globals). io: std.Io, /// Environment block (post.Env); XDG_RUNTIME_DIR names the registry. env: post.Env, uname: []const u8, aname: []const u8 = "", /// Maximum 9P message size to request per dial. msize: u32 = 131072, }; /// One registry subdirectory served as a synthetic directory: the real /// path it mirrors, the registry-relative key its children dial under, /// its slot (which names its node id) and its depth from the registry. const SynthDir = struct { slot: usize = 0, parent: u64 = 0, depth: u8 = 0, path_buf: [post.sun_path_len]u8 = undefined, path_len: u16 = 0, rel_buf: [post.sun_path_len]u8 = undefined, rel_len: u16 = 0, fn path(sd: *const SynthDir) [:0]const u8 { return sd.path_buf[0..sd.path_len :0]; } fn rel(sd: *const SynthDir) []const u8 { return sd.rel_buf[0..sd.rel_len]; } fn nodeid(sd: *const SynthDir) u64 { return mountNode(synth_index, sd.slot + 1); } }; /// One dialed server: its 9P session, its bridge state and its worker /// thread, plus the queue the dispatcher feeds requests through. const Mount = struct { gpa: std.mem.Allocator, io: std.Io, debug: bool, name: []u8, index: u32, /// This server's root node id (index in the top bits, 1 below). root_node: u64, /// Attr of the server's 9P root, cached from the dial-time stat: the /// dispatcher answers LOOKUP of the name from it without an rpc (the /// session belongs to the worker thread). root_attr: fuse.Attr, session: *nine.Session, b: *Bridge, mutex: std.Io.Mutex = .init, cond: std.Io.Condition = .init, queue: std.ArrayList([]align(8) u8) = .empty, /// The 9P connection died: the subtree answers EIO; a walk into the /// name re-dials as a new mount. Set by the worker, read by all. dead: std.atomic.Value(bool) = .init(false), /// Draining and exiting (child gone, FUSE device gone or DESTROY). /// Under `mutex`. stopping: bool = false, /// The dispatcher drops a FUSE_INTERRUPT's target unique in here; the /// worker's session polls it while a 9P reply is outstanding. int_pipe: [2]i32, thread: std.Thread, }; /// Runs the mntgen dispatcher on the calling thread until the FUSE fd /// reports ENODEV, a DESTROY arrives, or `stop_fd` becomes readable (the /// program exited; it is also what unblocks every worker's 9P wait). The /// synthetic root (node 1) is served here; everything else is routed by the /// node id's top bits to the owning server's worker. pub fn serveMntgen(gpa: std.mem.Allocator, fuse_fd: i32, stop_fd: i32, mo: MntgenOptions, opts: Options) !void { var effective = opts; if (!effective.direct_io) effective.attr_timeout_ns = 0; // same reasoning as `serve` var mg: Mntgen = .{ .gpa = gpa, .io = mo.io, .fuse_fd = fuse_fd, .stop_fd = stop_fd, .mo = mo, .opts = effective, .mounts = try gpa.alloc(?*Mount, max_mounts), }; defer mg.deinit(); @memset(mg.mounts, null); mg.req_buf = try gpa.alignedAlloc(u8, .@"8", request_buf_len); mg.data_buf = try gpa.alloc(u8, max_write); mg.stage = try gpa.alignedAlloc(u8, comptime std.mem.Alignment.fromByteUnits(@alignOf(usize)), stage_len); // Every read of the FUSE fd follows a poll (see `serve`). fuse.setNonblocking(fuse_fd) catch return error.FuseIo; var pfds = [_]linux.pollfd{ .{ .fd = fuse_fd, .events = linux.POLL.IN, .revents = 0 }, .{ .fd = stop_fd, .events = linux.POLL.IN, .revents = 0 }, }; while (true) { pfds[0].revents = 0; pfds[1].revents = 0; const rc = linux.poll(&pfds, pfds.len, -1); switch (linux.errno(rc)) { .SUCCESS => {}, .INTR, .AGAIN => continue, else => return error.Io, } if (pfds[1].revents != 0) { mg.trace("stop_fd readable; leaving mntgen loop", .{}); return; } if (pfds[0].revents == 0) continue; const req = (fuse.readRequestOnce(fuse_fd, mg.req_buf) catch |e| switch (e) { error.Retry => continue, error.Protocol => return error.FuseProtocol, else => return error.FuseIo, }) orelse { mg.trace("fuse fd reports ENODEV; unmounted", .{}); return; }; if (!try mg.route(req)) return; } } const Mntgen = struct { gpa: std.mem.Allocator, io: std.Io, fuse_fd: i32, stop_fd: i32, mo: MntgenOptions, opts: Options, req_buf: []align(8) u8 = &.{}, data_buf: []u8 = &.{}, stage: []align(@alignOf(usize)) u8 = &.{}, /// Slot per mount ordinal; `mounts[index]`. Mutated only by the /// dispatcher thread; workers are reached through their queue. mounts: []?*Mount, /// Next mount ordinal to hand out. Starts at 1: an index-0 mount would /// make that server's root node id collide with the synthetic root /// (node 1), and `route` sends every nodeid below 2^32 to the root /// handler anyway. next_index: u32 = 1, /// Open directory handles of the synthetic root (dispatcher-owned). root_dirs: [max_root_dirs]?*DirList = @splat(null), /// Synthetic registry subdirectories (dispatcher-owned), by slot. synths: [max_synth_dirs]?*SynthDir = @splat(null), fn deinit(mg: *Mntgen) void { // Wake every worker, then join: at this point the child is gone (or // the FUSE device is), so `stop_fd` readable makes any in-flight // 9P rpc fail with error.Stopped and each worker exits promptly. for (mg.mounts) |slot| { const m = slot orelse continue; m.mutex.lockUncancelable(mg.io); m.stopping = true; m.mutex.unlock(mg.io); m.cond.signal(mg.io); } for (mg.mounts) |slot| { const m = slot orelse continue; m.thread.join(); for (m.queue.items) |buf| mg.gpa.free(buf); m.queue.deinit(mg.gpa); _ = linux.close(m.int_pipe[0]); _ = linux.close(m.int_pipe[1]); m.b.deinit(); m.session.deinit(); mg.gpa.free(m.name); mg.gpa.destroy(m.b); mg.gpa.destroy(m.session); mg.gpa.destroy(m); } for (&mg.root_dirs) |*slot| { if (slot.*) |list| { list.deinit(mg.gpa); mg.gpa.destroy(list); slot.* = null; } } for (&mg.synths) |*slot| { if (slot.*) |sd| { mg.gpa.destroy(sd); slot.* = null; } } if (mg.req_buf.len != 0) mg.gpa.free(mg.req_buf); if (mg.data_buf.len != 0) mg.gpa.free(mg.data_buf); if (mg.stage.len != 0) mg.gpa.free(mg.stage); mg.gpa.free(mg.mounts); } fn trace(mg: *const Mntgen, comptime fmt: []const u8, args: anytype) void { if (mg.opts.debug) std.debug.print("9ns: " ++ fmt ++ "\n", args); } fn reply(mg: *Mntgen, unique: u64, payloads: []const []const u8) error{FuseIo}!void { var total: usize = 0; for (payloads) |p| total += p.len; mg.trace("-> unique={d} ok ({d} bytes)", .{ unique, total }); fuse.reply(mg.fuse_fd, unique, payloads) catch return error.FuseIo; } fn replyError(mg: *Mntgen, unique: u64, code: linux.E) error{FuseIo}!void { mg.trace("-> unique={d} error E{s}", .{ unique, @tagName(code) }); fuse.replyError(mg.fuse_fd, unique, code) catch return error.FuseIo; } /// Handles one kernel request. Returns false when the loop should stop /// (DESTROY). Fatal FUSE-device errors propagate. fn route(mg: *Mntgen, req: fuse.Request) !bool { const h = req.header; const op = h.op(); mg.trace("<- {s} unique={d} nodeid={d} len={d}", .{ opName(op), h.unique, h.nodeid, h.len }); switch (op) { .init => { const in = fuse.body(fuse.InitIn, req) catch { mg.replyError(h.unique, .INVAL) catch {}; return true; }; const out = fuse.initReply(in, max_write); try mg.reply(h.unique, &.{std.mem.asBytes(&out)}); return true; }, .destroy => { mg.reply(h.unique, &.{}) catch {}; mg.trace("DESTROY; unmounting", .{}); return false; }, .interrupt => { const in = fuse.body(fuse.InterruptIn, req) catch return true; mg.routeInterrupt(in.unique); return true; }, // A BATCH_FORGET carries entries for many owners at once and // cannot be routed by its header nodeid (the kernel sends 0); // see `distributeForgets`. .batch_forget => { mg.distributeForgets(req); return true; }, else => {}, } if (h.nodeid == fuse.root_id) { try mg.handleRoot(req); return true; } if (mountIndex(h.nodeid) == synth_index) { try mg.handleSynthDir(req); return true; } const wants_reply = switch (op) { .forget, .batch_forget => false, else => true, }; const idx = mountIndex(h.nodeid); const m = if (idx < max_mounts) mg.mounts[idx] else null; if (m != null and !m.?.dead.load(.seq_cst)) { mg.enqueue(m.?, mg.req_buf[0..h.len]); return true; } // A retired mount (its server died and a later walk re-dialed under // a new index) or an unknown node: the dead subtree answers EIO, a // FORGET is simply dropped. mg.trace(" nodeid={d} has no live mount (index {d})", .{ h.nodeid, idx }); if (wants_reply) mg.replyError(h.unique, .IO) catch {}; return true; } /// Copies the raw request bytes and hands them to the mount's worker. /// Never blocks on the worker; allocation or a racing death reject the /// request with an errno reply instead. fn enqueue(mg: *Mntgen, m: *Mount, raw: []const u8) void { const wants_reply = switch (@as(fuse.Opcode, @enumFromInt(std.mem.readInt(u32, raw[4..8], .little)))) { .forget, .batch_forget => false, else => true, }; const unique = std.mem.readInt(u64, raw[8..16], .little); const buf = mg.gpa.alignedAlloc(u8, .@"8", raw.len) catch { if (wants_reply) mg.replyError(unique, .NOMEM) catch {}; return; }; @memcpy(buf, raw); var reject: ?linux.E = null; m.mutex.lockUncancelable(mg.io); if (m.dead.load(.seq_cst) or m.stopping) { reject = .IO; } else if (m.queue.append(mg.gpa, buf)) |_| { m.cond.signal(mg.io); } else |_| { reject = .NOMEM; } m.mutex.unlock(mg.io); if (reject) |code| { mg.gpa.free(buf); if (wants_reply) mg.replyError(unique, code) catch {}; } } /// FUSE_BATCH_FORGET carries (nodeid, nlookup) entries for many owners /// at once — synthetic-root, synthetic-subdirectory and per-mount /// nodeids can all appear in one batch, and the kernel puts 0 in the /// header nodeid. Routing such a batch like an ordinary request would /// hand every entry to one wrong owner and drop the rest: a mount's /// bridge would leak the fid behind each dropped entry forever, and /// a dropped synthetic-subdirectory entry would leak its slot until /// every subdirectory lookup answers EIO. So the dispatcher keeps /// what it owns — the root holds nothing, a synth slot is freed here — /// and forwards each mount-owned entry to its worker as a plain /// single FUSE_FORGET. fn distributeForgets(mg: *Mntgen, req: fuse.Request) void { const in = fuse.body(fuse.BatchForgetIn, req) catch return; const rest = req.body[@sizeOf(fuse.BatchForgetIn)..]; const count: usize = in.count; if (rest.len < count * @sizeOf(fuse.ForgetOne)) return; for (0..count) |i| { const one = std.mem.bytesToValue(fuse.ForgetOne, rest[i * @sizeOf(fuse.ForgetOne) ..][0..@sizeOf(fuse.ForgetOne)]); mg.forgetOne(one.nodeid, one.nlookup, req.header.unique); } } /// Forgets one node by id, whatever owns it: the root holds nothing, /// a synthetic subdirectory's slot goes back to the pool, and a /// mount-owned node is forwarded to its worker (which clunks the fid /// behind it). Unknown or dead mounts drop the entry, like a single /// FORGET routed by `route`. fn forgetOne(mg: *Mntgen, nodeid: u64, nlookup: u64, unique: u64) void { if (nodeid == fuse.root_id) return; const idx = mountIndex(nodeid); if (idx == synth_index) { const local = nodeid & mount_node_mask; if (local == 0 or local > max_synth_dirs) return; const slot: usize = @intCast(local - 1); if (mg.synths[slot]) |sd| { mg.trace(" synthetic directory slot {d} forgotten", .{slot}); mg.gpa.destroy(sd); mg.synths[slot] = null; } return; } const m = if (idx < max_mounts) mg.mounts[idx] else null; if (m == null or m.?.dead.load(.seq_cst)) return; // A synthesized single FORGET (no reply is expected for one, so // the unique is only bookkeeping). var buf: [@sizeOf(fuse.InHeader) + @sizeOf(fuse.ForgetIn)]u8 = undefined; const hdr = fuse.InHeader{ .len = @sizeOf(fuse.InHeader) + @sizeOf(fuse.ForgetIn), .opcode = @intFromEnum(fuse.Opcode.forget), .unique = unique, .nodeid = nodeid, .uid = 0, .gid = 0, .pid = 0, .total_extlen = 0, .padding = 0, }; @memcpy(buf[0..@sizeOf(fuse.InHeader)], std.mem.asBytes(&hdr)); @memcpy(buf[@sizeOf(fuse.InHeader)..], std.mem.asBytes(&fuse.ForgetIn{ .nlookup = nlookup })); mg.enqueue(m.?, &buf); } /// A FUSE_INTERRUPT names the request it wants cancelled; FUSE uniques /// are unique across the whole connection, so the mount whose bridge is /// currently serving that unique gets the packet and its session turns /// it into a Tflush (see the worker's interrupt source below). fn routeInterrupt(mg: *Mntgen, target: u64) void { for (mg.mounts) |slot| { const m = slot orelse continue; if (m.dead.load(.seq_cst)) continue; if (m.b.cur_unique.load(.seq_cst) == target) { mg.trace(" interrupt for unique={d}: forwarding to '{s}'", .{ target, m.name }); var packet: [8]u8 = undefined; std.mem.writeInt(u64, &packet, target, .little); // The pipe is small and nonblocking; a dropped packet only // means one interrupt missed its window (the kernel does // not retry INTERRUPTs, but the child's exit ends the // session through stop_fd regardless). _ = linux.write(m.int_pipe[1], &packet, packet.len); return; } } mg.trace(" interrupt for unique={d} (not in flight; ignored)", .{target}); } // -- the synthetic root (node 1) ----------------------------------------- fn rootAttr(mg: *const Mntgen) fuse.Attr { // Read-only like /srv: services are posted and unposted by their // servers, not created and removed through the mount. return .{ .ino = fuse.root_id, .mode = fuse.S_IFDIR | 0o555, .nlink = 2, .uid = mg.opts.uid, .gid = mg.opts.gid, .blksize = 4096, }; } fn rootEntryOut(mg: *Mntgen, nodeid: u64, attr: fuse.Attr) fuse.EntryOut { // No entry caching for synthetic-root entries: a walk re-LOOKUPs the // name, which is what notices a dead server and re-dials it. No // invalidation machinery needed, and nothing to invalidate. _ = mg; return .{ .nodeid = nodeid, .generation = 0, .attr = attr }; } fn handleRoot(mg: *Mntgen, req: fuse.Request) error{FuseIo}!void { const u = req.header.unique; switch (req.header.op()) { .getattr => { const out = fuse.AttrOut{ .attr = mg.rootAttr() }; try mg.reply(u, &.{std.mem.asBytes(&out)}); }, .lookup => try mg.rootLookup(req), .opendir => { var fh: ?usize = null; for (&mg.root_dirs, 0..) |*slot, i| { if (slot.* == null) { fh = i; break; } } const slot = fh orelse return mg.replyError(u, .MFILE); const list = mg.rootListing() catch { return mg.replyError(u, .IO); }; mg.root_dirs[slot] = list; const out = fuse.OpenOut{ .fh = slot }; try mg.reply(u, &.{std.mem.asBytes(&out)}); }, .readdir => { const in = fuse.body(fuse.ReadIn, req) catch return mg.replyError(u, .BADF); if (in.fh >= max_root_dirs) return mg.replyError(u, .BADF); const list = mg.root_dirs[@intCast(in.fh)] orelse return mg.replyError(u, .BADF); const size: usize = @min(in.size, max_write); const used = packDirents(list.entries.items, in.offset, mg.data_buf[0..size]); try mg.reply(u, &.{mg.data_buf[0..used]}); }, .release, .releasedir => { const in = fuse.body(fuse.ReleaseIn, req) catch return mg.replyError(u, .BADF); if (in.fh < max_root_dirs) { if (mg.root_dirs[@intCast(in.fh)]) |list| { list.deinit(mg.gpa); mg.gpa.destroy(list); mg.root_dirs[@intCast(in.fh)] = null; } } try mg.reply(u, &.{}); }, .statfs => { const out = fuse.StatfsOut{ .st = .{ .bsize = 4096, .namelen = 255, .frsize = 4096 } }; try mg.reply(u, &.{std.mem.asBytes(&out)}); }, .flush, .fsync, .fsyncdir => try mg.reply(u, &.{}), .forget, .batch_forget => {}, .access => try mg.replyError(u, .NOSYS), // Capability probes (xattrs, statx with STATX_ALL) must read as // "not supported", like the ordinary bridge's answer for them: // EPERM makes `ls -l` and plain `stat` blame the mount root // itself with "Operation not permitted", while the kernel // caches ENOSYS as "no xattrs / no extra attrs here" and falls // back to the GETATTR data. .setxattr, .getxattr, .listxattr, .removexattr, .statx => try mg.replyError(u, .NOSYS), // The root is synthetic and read-only: services are managed by // their servers (cloud9.post's post/unpost), not through files. else => try mg.replyError(u, .PERM), } } fn rootLookup(mg: *Mntgen, req: fuse.Request) error{FuseIo}!void { const u = req.header.unique; const name = fuse.nameAfter(void, req) catch return mg.replyError(u, .INVAL); if (std.mem.eql(u8, name, ".") or std.mem.eql(u8, name, "..")) { const out = mg.rootEntryOut(fuse.root_id, mg.rootAttr()); return mg.reply(u, &.{std.mem.asBytes(&out)}); } if (!post.legalName(name)) return mg.replyError(u, .NOENT); if (findMount(mg.mounts, name)) |m| { // Live: answer from the dial-time snapshot. No 9P rpc (the // session belongs to the worker thread); no connection made. mg.trace(" lookup '{s}': mount {d} already live", .{ name, m.index }); const out = mg.rootEntryOut(m.root_node, m.root_attr); return mg.reply(u, &.{std.mem.asBytes(&out)}); } // What the entry is decides what a walk into it becomes: a socket // dials (the original behavior), a directory is served like the // root itself (its sockets dial on walk, its directories recurse), // anything else answers EIO. var path_buf: [post.sun_path_len]u8 = undefined; const entry_path = post.registryPath(mg.mo.env, name, &path_buf) catch return mg.replyError(u, .NOENT); const st = std.Io.Dir.statFile(.cwd(), mg.mo.io, entry_path, .{}) catch { mg.trace(" lookup '{s}': nothing posted under that name", .{name}); return mg.replyError(u, .NOENT); }; switch (st.kind) { .directory => { const sd = mg.newSynth(entry_path, name, 1, fuse.root_id) catch |e| { mg.trace(" lookup '{s}': no synthetic slot: {t}", .{ name, e }); return mg.replyError(u, .IO); }; const node = sd.nodeid(); const out = mg.rootEntryOut(node, mg.synthAttr(node)); return mg.reply(u, &.{std.mem.asBytes(&out)}); }, .unix_domain_socket => {}, else => { mg.trace(" lookup '{s}': registry entry is not a socket", .{name}); return mg.replyError(u, .IO); }, } const m = mg.dialMount(name) catch |e| switch (e) { error.Stale => { mg.trace(" lookup '{s}': registry entry is stale (no server behind it)", .{name}); return mg.replyError(u, .IO); }, else => { mg.trace(" lookup '{s}': dial failed: {t}", .{ name, e }); return mg.replyError(u, .IO); }, }; mg.trace(" lookup '{s}': dialed as mount {d}", .{ name, m.index }); const out = mg.rootEntryOut(m.root_node, m.root_attr); try mg.reply(u, &.{std.mem.asBytes(&out)}); } /// One OPENDIR of the synthetic root: `.` and `..` plus every registry /// entry whose name the kernel would accept, snapshotted for the life /// of the handle (a fresh OPENDIR sees fresh posts). fn rootListing(mg: *Mntgen) !*DirList { const list = try mg.gpa.create(DirList); errdefer mg.gpa.destroy(list); list.* = .{}; errdefer list.deinit(mg.gpa); try list.entries.append(mg.gpa, .{ .name = try mg.gpa.dupe(u8, "."), .ino = fuse.root_id, .dtype = fuse.DT_DIR }); try list.entries.append(mg.gpa, .{ .name = try mg.gpa.dupe(u8, ".."), .ino = fuse.root_id, .dtype = fuse.DT_DIR }); var names = post.posted(mg.mo.io, mg.mo.env, mg.stage) catch |e| { mg.trace(" registry listing failed: {t}", .{e}); return error.Registry; }; while (names.next()) |name| { if (!validDirentName(name)) continue; // raw entries; dial filters further try list.entries.append(mg.gpa, .{ .name = try mg.gpa.dupe(u8, name), .ino = nameIno(name), .dtype = fuse.DT_DIR }); } return list; } // -- synthetic registry subdirectories ------------------------------------ /// A directory entry in the registry (or in one of its /// subdirectories) is served like the root: a synthetic directory /// listing the real one, whose sockets dial on walk and whose /// directories recurse. Served on the dispatcher thread, like the /// root. fn handleSynthDir(mg: *Mntgen, req: fuse.Request) error{FuseIo}!void { const u = req.header.unique; const local = req.header.nodeid & mount_node_mask; if (local == 0 or local > max_synth_dirs) return mg.replyError(u, .IO); const slot: usize = @intCast(local - 1); const sd = mg.synths[slot] orelse return mg.replyError(u, .IO); switch (req.header.op()) { .forget, .batch_forget => { // The kernel dropped the dentry; the slot goes with it. mg.trace(" synthetic directory slot {d} forgotten", .{slot}); mg.gpa.destroy(sd); mg.synths[slot] = null; }, .getattr => { const out = fuse.AttrOut{ .attr = mg.synthAttr(sd.nodeid()) }; try mg.reply(u, &.{std.mem.asBytes(&out)}); }, .lookup => try mg.synthLookup(sd, req), .opendir => { var fh: ?usize = null; for (&mg.root_dirs, 0..) |*dir_slot, i| { if (dir_slot.* == null) { fh = i; break; } } const dir_slot = fh orelse return mg.replyError(u, .MFILE); const list = mg.synthListing(sd) catch { return mg.replyError(u, .IO); }; mg.root_dirs[dir_slot] = list; const out = fuse.OpenOut{ .fh = dir_slot }; try mg.reply(u, &.{std.mem.asBytes(&out)}); }, .readdir => { const in = fuse.body(fuse.ReadIn, req) catch return mg.replyError(u, .BADF); if (in.fh >= max_root_dirs) return mg.replyError(u, .BADF); const list = mg.root_dirs[@intCast(in.fh)] orelse return mg.replyError(u, .BADF); const size: usize = @min(in.size, max_write); const used = packDirents(list.entries.items, in.offset, mg.data_buf[0..size]); try mg.reply(u, &.{mg.data_buf[0..used]}); }, .release, .releasedir => { const in = fuse.body(fuse.ReleaseIn, req) catch return mg.replyError(u, .BADF); if (in.fh < max_root_dirs) { if (mg.root_dirs[@intCast(in.fh)]) |list| { list.deinit(mg.gpa); mg.gpa.destroy(list); mg.root_dirs[@intCast(in.fh)] = null; } } try mg.reply(u, &.{}); }, .statfs => { const out = fuse.StatfsOut{ .st = .{ .bsize = 4096, .namelen = 255, .frsize = 4096 } }; try mg.reply(u, &.{std.mem.asBytes(&out)}); }, .flush, .fsync, .fsyncdir => try mg.reply(u, &.{}), // Capability probes read as "not supported", like the root's. .access, .setxattr, .getxattr, .listxattr, .removexattr, .statx => try mg.replyError(u, .NOSYS), else => try mg.replyError(u, .PERM), } } fn synthAttr(mg: *const Mntgen, nodeid: u64) fuse.Attr { // Read-only like the root and like /srv. return .{ .ino = nodeid, .mode = fuse.S_IFDIR | 0o555, .nlink = 2, .uid = mg.opts.uid, .gid = mg.opts.gid, .blksize = 4096, }; } fn synthLookup(mg: *Mntgen, sd: *SynthDir, req: fuse.Request) error{FuseIo}!void { const u = req.header.unique; const name = fuse.nameAfter(void, req) catch return mg.replyError(u, .INVAL); if (std.mem.eql(u8, name, ".")) { const out = mg.rootEntryOut(sd.nodeid(), mg.synthAttr(sd.nodeid())); return mg.reply(u, &.{std.mem.asBytes(&out)}); } // The kernel resolves ".." from its own dentry tree and a LOOKUP of // it has never been observed, but it must not alias the directory // onto itself either: answer with the parent's node id (the root's // for a top-level subdirectory — its attr is the same shape). if (std.mem.eql(u8, name, "..")) { const out = mg.rootEntryOut(sd.parent, mg.synthAttr(sd.parent)); return mg.reply(u, &.{std.mem.asBytes(&out)}); } if (!post.legalName(name)) return mg.replyError(u, .NOENT); // The registry-relative key this child dials under (a mount's // name, for findMount). var key_buf: [post.sun_path_len]u8 = undefined; const key = std.fmt.bufPrint(&key_buf, "{s}/{s}", .{ sd.rel(), name }) catch return mg.replyError(u, .NOTNAM); var path_buf: [post.sun_path_len]u8 = undefined; const child = std.fmt.bufPrintSentinel(&path_buf, "{s}/{s}", .{ sd.path(), name }, 0) catch return mg.replyError(u, .NOTNAM); if (findMount(mg.mounts, key)) |m| { mg.trace(" lookup '{s}': mount {d} already live", .{ key, m.index }); const out = mg.rootEntryOut(m.root_node, m.root_attr); return mg.reply(u, &.{std.mem.asBytes(&out)}); } const st = std.Io.Dir.statFile(.cwd(), mg.mo.io, child, .{}) catch |e| switch (e) { error.FileNotFound => { mg.trace(" lookup '{s}': no entry", .{key}); return mg.replyError(u, .NOENT); }, else => { mg.trace(" lookup '{s}': stat failed: {t}", .{ key, e }); return mg.replyError(u, .IO); }, }; switch (st.kind) { .directory => { const child_sd = mg.newSynth(child, key, sd.depth + 1, sd.nodeid()) catch |e| { mg.trace(" lookup '{s}': no synthetic slot: {t}", .{ key, e }); return mg.replyError(u, .IO); }; const node = child_sd.nodeid(); const out = mg.rootEntryOut(node, mg.synthAttr(node)); return mg.reply(u, &.{std.mem.asBytes(&out)}); }, .unix_domain_socket => { const m = mg.dialMountAt(child, key) catch |e| { mg.trace(" lookup '{s}': dial failed: {t}", .{ key, e }); return mg.replyError(u, .IO); }; mg.trace(" lookup '{s}': dialed as mount {d}", .{ key, m.index }); const out = mg.rootEntryOut(m.root_node, m.root_attr); try mg.reply(u, &.{std.mem.asBytes(&out)}); }, // Not a service and not a directory: the entry answers EIO on // walk, like a plain file in the registry itself. else => { mg.trace(" lookup '{s}': entry is not a socket or directory", .{key}); return mg.replyError(u, .IO); }, } } /// One OPENDIR of a synthetic subdirectory: `.` and `..` plus the /// real directory's entries, snapshotted for the life of the handle /// (a fresh OPENDIR sees fresh entries), like the root. fn synthListing(mg: *Mntgen, sd: *SynthDir) !*DirList { const list = try mg.gpa.create(DirList); errdefer mg.gpa.destroy(list); list.* = .{}; errdefer list.deinit(mg.gpa); const node = sd.nodeid(); try list.entries.append(mg.gpa, .{ .name = try mg.gpa.dupe(u8, "."), .ino = node, .dtype = fuse.DT_DIR }); try list.entries.append(mg.gpa, .{ .name = try mg.gpa.dupe(u8, ".."), .ino = sd.parent, .dtype = fuse.DT_DIR }); var names = post.postedDir(mg.mo.io, sd.path(), mg.stage) catch |e| { mg.trace(" directory listing failed: {t}", .{e}); return error.Registry; }; while (names.next()) |name| { if (!validDirentName(name)) continue; try list.entries.append(mg.gpa, .{ .name = try mg.gpa.dupe(u8, name), .ino = nameIno(name), .dtype = fuse.DT_DIR }); } return list; } /// Allocates a synthetic directory node mirroring `path`, keyed by /// the registry-relative `key`, at `depth` under `parent`'s node id. fn newSynth(mg: *Mntgen, path: [:0]const u8, key: []const u8, depth: u8, parent: u64) !*SynthDir { if (depth > max_synth_depth) return error.TooDeep; var slot: ?usize = null; for (&mg.synths, 0..) |*s, i| { if (s.* == null) { slot = i; break; } } const i = slot orelse return error.TooMany; const sd = try mg.gpa.create(SynthDir); errdefer mg.gpa.destroy(sd); sd.* = .{ .slot = i, .parent = parent, .depth = depth }; if (path.len + 1 > sd.path_buf.len) return error.NameTooLong; @memcpy(sd.path_buf[0..path.len], path); sd.path_buf[path.len] = 0; sd.path_len = @intCast(path.len); if (key.len > sd.rel_buf.len) return error.NameTooLong; @memcpy(sd.rel_buf[0..key.len], key); sd.rel_len = @intCast(key.len); mg.synths[i] = sd; return sd; } // -- dialing --------------------------------------------------------------- /// Dials `name` out of the registry, attaches, stats the server root and /// starts its worker thread. Runs on the dispatcher thread, in service /// of the LOOKUP that triggered it (so a hung server delays that walk, /// like it would delay any 9P client). fn dialMount(mg: *Mntgen, name: []const u8) !*Mount { const stream = try post.dial(mg.mo.io, mg.mo.env, name); return mg.mountStream(stream, name); } /// Dials the socket at `path` (a registry subdirectory entry) and /// mounts it under `key`, the registry-relative path — the /// subdirectory analogue of `dialMount`. fn dialMountAt(mg: *Mntgen, path: [:0]const u8, key: []const u8) !*Mount { const stream = try post.dialPath(mg.mo.io, path); return mg.mountStream(stream, key); } /// The shared dial tail: session, attach, stat, bridge and worker. fn mountStream(mg: *Mntgen, stream: std.Io.net.Stream, key: []const u8) !*Mount { if (mg.next_index >= max_mounts) return error.TooManyMounts; const index: u32 = mg.next_index; const fd: i32 = @intCast(stream.socket.handle); // The session below does blocking I/O: make sure a dial that left // the descriptor nonblocking cannot spin its read loop on EAGAIN, // and keep the descriptor out of any future exec (the running // program already forked, but hygiene is free). _ = linux.fcntl(fd, linux.F.SETFL, 0); _ = linux.fcntl(fd, linux.F.SETFD, linux.FD_CLOEXEC); const session = mg.gpa.create(nine.Session) catch |e| { _ = linux.close(fd); return e; }; errdefer mg.gpa.destroy(session); // The child's exit must unblock the whole dial — Tversion included, // not just the attach/stat after it: a silent server must not pin // the dispatcher past the program. session.* = nine.Session.connectWatched(mg.gpa, .{ .fd = fd }, mg.mo.msize, mg.stop_fd) catch |e| { _ = linux.close(fd); return e; }; errdefer session.deinit(); _ = try session.attach(0, mg.mo.uname, mg.mo.aname); const st = try session.stat(0); const b = try mg.gpa.create(Bridge); errdefer mg.gpa.destroy(b); const root_node = mountNode(index, fuse.root_id); b.* = .{ .gpa = mg.gpa, .fuse_fd = mg.fuse_fd, .nine = session, .opts = mg.opts, .root_id = root_node, .ino_xor = @as(u64, index + 1) << 48, }; errdefer b.deinit(); b.data_buf = try mg.gpa.alloc(u8, max_write); // The kernel-visible root of this server's subtree: global node id // (index included), fid 0, its qid from the stat above. From here on // the bridge is an ordinary single-server bridge, just with node // ids that already carry the index. try b.inodes.put(mg.gpa, root_node, .{ .fid = 0, .qid = st.qid, .nlookup = 1, .parent = root_node }); try b.by_qid.put(mg.gpa, st.qid.path, root_node); b.root_path = st.qid.path; b.next_node = root_node + 1; var pipes: [2]i32 = undefined; if (linux.errno(linux.pipe2(&pipes, .{ .CLOEXEC = true, .NONBLOCK = true })) != .SUCCESS) { return error.SystemResources; } errdefer { _ = linux.close(pipes[0]); _ = linux.close(pipes[1]); } const m = try mg.gpa.create(Mount); errdefer mg.gpa.destroy(m); m.* = .{ .gpa = mg.gpa, .io = mg.io, .debug = mg.opts.debug, .name = try mg.gpa.dupe(u8, key), .index = index, .root_node = root_node, .root_attr = attrFromStat(st, inoFromPath(st.qid.path) ^ b.ino_xor, mg.opts.uid, mg.opts.gid), .session = session, .b = b, .int_pipe = pipes, .thread = undefined, }; errdefer mg.gpa.free(m.name); // The interrupt source must be installed before the thread starts. session.interrupt = mountInterrupt(m); m.thread = std.Thread.spawn(.{}, workerMain, .{m}) catch |e| { session.interrupt = null; return e; }; mg.mounts[index] = m; mg.next_index += 1; return m; } }; /// The first live mount posted under `name`, skipping dead ones (a walk into /// a name whose server died dials it afresh rather than reuse the corpse). fn findMount(mounts: []const ?*Mount, name: []const u8) ?*Mount { for (mounts) |slot| { const m = slot orelse continue; if (m.dead.load(.seq_cst)) continue; if (std.mem.eql(u8, m.name, name)) return m; } return null; } /// A server's worker: pops copied requests off its queue and runs them /// through the ordinary bridge dispatch, one at a time (same concurrency /// contract as single-connection 9ns). The dispatcher keeps reading /dev/fuse /// meanwhile, so a slow server never blocks the other names. fn workerMain(m: *Mount) void { while (true) { m.mutex.lockUncancelable(m.io); while (m.queue.items.len == 0 and !m.stopping and !m.dead.load(.seq_cst)) { m.cond.waitUncancelable(m.io, &m.mutex); } const buf: ?[]align(8) u8 = if (m.queue.items.len != 0) m.queue.orderedRemove(0) else null; m.mutex.unlock(m.io); if (buf) |bytes| { defer m.gpa.free(bytes); serveQueued(m, bytes); continue; } break; // empty, and stopping or dead } } fn serveQueued(m: *Mount, bytes: []align(8) u8) void { const header = std.mem.bytesToValue(fuse.InHeader, bytes[0..@sizeOf(fuse.InHeader)]); const req = fuse.Request{ .header = header, .body = bytes[@sizeOf(fuse.InHeader)..header.len] }; if (m.debug) std.debug.print("9ns: [{s}] <- {s} unique={d} nodeid={d}\n", .{ m.name, opName(req.header.op()), req.header.unique, req.header.nodeid }); const keep_going = m.b.dispatch(req) catch { // The 9P session (or the FUSE device) died mid-request: dispatch has // already replied EIO for this one. Everything still queued answers // EIO just as fast, and no new request is routed here again. m.dead.store(true, .seq_cst); if (m.debug) std.debug.print("9ns: [{s}] server connection lost; subtree now answers EIO\n", .{m.name}); return; }; if (!keep_going) { // DESTROY or the child exited (error.Stopped): drain and exit. m.mutex.lockUncancelable(m.io); m.stopping = true; m.mutex.unlock(m.io); } } // -- a mount's interrupt source ------------------------------------------------- // // The worker's session polls `int_pipe[0]` while a 9P reply is outstanding; // the dispatcher writes the interrupted request's unique into it. Unlike the // single-connection source, no FUSE reading happens here — the dispatcher // owns /dev/fuse. fn mountInterrupt(m: *Mount) nine.Interrupt { return .{ .ctx = m, .watch = mountWatch, .onReadable = mountOnReadable, .armed = mountArmed }; } fn mountWatch(ctx: *anyopaque) i32 { const m: *Mount = @ptrCast(@alignCast(ctx)); return m.int_pipe[0]; } fn mountOnReadable(ctx: *anyopaque) nine.Session.Error!bool { const m: *Mount = @ptrCast(@alignCast(ctx)); var packet: [8]u8 = undefined; var hit = false; while (true) { const rc = linux.read(m.int_pipe[0], &packet, packet.len); switch (linux.errno(rc)) { .SUCCESS => { // Pipe writes of 8 bytes are atomic; a short read cannot happen. const target = std.mem.readInt(u64, &packet, .little); const cur = m.b.cur_unique.load(.seq_cst); if (cur != 0 and target == cur) hit = true; }, .INTR => continue, .AGAIN => break, // drained else => return error.Io, } } if (hit) { if (m.debug) std.debug.print("9ns: [{s}] interrupt for unique={d} (in flight): sending Tflush\n", .{ m.name, m.b.cur_unique.load(.seq_cst) }); m.b.interrupted = true; return true; } return false; } fn mountArmed(ctx: *anyopaque) bool { const m: *Mount = @ptrCast(@alignCast(ctx)); return m.b.interrupted; } // -- 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 .{ .ino = ino, // The kernel marks an inode bad when size > LLONG_MAX; clamp hostile lengths. .size = @min(st.length, std.math.maxInt(i64)), // Saturating: a hostile length of 2^64-1 must not overflow. .blocks = st.length / 512 + @intFromBool(st.length % 512 != 0), .atime = st.atime, .mtime = st.mtime, .ctime = st.mtime, .mode = ftype | (st.mode & 0o777), .nlink = 1, .uid = uid, .gid = gid, .blksize = 4096, }; } /// Kernel open(2) flags → 9P open mode. O_APPEND has no 9P equivalent and is ignored. pub fn openMode(flags: u32) u8 { const o: linux.O = @bitCast(flags); var mode: u8 = switch (o.ACCMODE) { .RDONLY => cloud9.oread, .WRONLY => cloud9.owrite, .RDWR => cloud9.ordwr, }; if (o.TRUNC) mode |= cloud9.otrunc; return mode; } /// Parses consecutive 9P directory records (2-byte size + Stat) and appends entries. pub fn parseDirRecords(gpa: std.mem.Allocator, bytes: []const u8, list: *DirList) error{ OutOfMemory, BadDir }!void { var pos: usize = 0; while (pos < bytes.len) { if (bytes.len - pos < 2) return error.BadDir; const size: usize = std.mem.readInt(u16, bytes[pos..][0..2], .little); if (bytes.len - pos < 2 + size) return error.BadDir; const st = cloud9.Stat.decode(bytes[pos..][0 .. 2 + size]) catch return error.BadDir; pos += 2 + size; // The kernel rejects a whole READDIR reply (EIO) over one bad name, and // "." and ".." are synthesised by loadDir: drop such records instead. if (!validDirentName(st.name)) continue; const name = try gpa.dupe(u8, st.name); errdefer gpa.free(name); try list.entries.append(gpa, .{ .name = name, .ino = st.qid.path, .dtype = if (st.mode & cloud9.dmdir != 0) fuse.DT_DIR else fuse.DT_REG, }); } } /// A name the kernel will accept in a dirent and that does not duplicate the synthetic "." / "..". pub fn validDirentName(name: []const u8) bool { if (name.len == 0 or name.len > max_name_len) return false; if (std.mem.indexOfAny(u8, name, "/\x00") != null) return false; if (std.mem.eql(u8, name, ".") or std.mem.eql(u8, name, "..")) return false; return true; } /// Packs dirents from `entries[offset..]` into `buf`; each record's `off` is its index + 1. /// Returns the number of bytes used. pub fn packDirents(entries: []const Entry, offset: u64, buf: []u8) usize { var used: usize = 0; var i: usize = @intCast(@min(offset, entries.len)); while (i < entries.len) : (i += 1) { const e = entries[i]; if (!fuse.addDirent(buf, &used, e.ino, @as(u64, i) + 1, e.dtype, e.name)) break; } return used; } fn randomU64() u64 { var bytes: [8]u8 = undefined; if (linux.errno(linux.getrandom(&bytes, bytes.len, 0)) == .SUCCESS) return std.mem.readInt(u64, &bytes, .little); var ts: linux.timespec = undefined; _ = linux.clock_gettime(.MONOTONIC, &ts); return @as(u64, @bitCast(ts.nsec)) ^ (@as(u64, @bitCast(ts.sec)) << 32); } fn nowSeconds() u32 { var ts: linux.timespec = undefined; if (linux.errno(linux.clock_gettime(.REALTIME, &ts)) != .SUCCESS) return 0; return @intCast(@as(u64, @intCast(ts.sec)) & 0xFFFF_FFFF); } fn opName(op: fuse.Opcode) []const u8 { return switch (op) { _ => "unknown", else => @tagName(op), }; } // Thin adapters so fuse.zig's parse errors become HandlerError.BadRequest. fn body(comptime T: type, req: fuse.Request) error{BadRequest}!*const T { return fuse.body(T, req) catch error.BadRequest; } fn nameAfter(comptime T: type, req: fuse.Request) error{BadRequest}![]const u8 { return fuse.nameAfter(T, req) catch error.BadRequest; } fn secondName(req: fuse.Request, first: []const u8, offset: usize) error{BadRequest}![]const u8 { return fuse.secondName(req, first, offset) catch error.BadRequest; } // -- tests ------------------------------------------------------------------------------ const testing = std.testing; test { // Force semantic analysis of `serve` and the whole dispatch path, which no // unit test can exercise without a FUSE mount. testing.refAllDecls(@This()); } fn testStat(name: []const u8, mode: u32, length: u64, path: u64) cloud9.Stat { return .{ .type = 0, .dev = 0, .qid = .{ .type = if (mode & cloud9.dmdir != 0) cloud9.qtdir else 0, .version = 0, .path = path }, .mode = mode, .atime = 100, .mtime = 200, .length = length, .name = name, .uid = "u", .gid = "g", .muid = "u", }; } test "attr mapping: DMDIR → S_IFDIR|perm, length → size/blocks" { const d = attrFromStat(testStat("d", cloud9.dmdir | 0o755, 0, 9), 9, 1000, 1001); try testing.expectEqual(fuse.S_IFDIR | 0o755, d.mode); try testing.expectEqual(@as(u64, 9), d.ino); try testing.expectEqual(@as(u64, 0), d.size); try testing.expectEqual(@as(u64, 0), d.blocks); try testing.expectEqual(@as(u32, 1000), d.uid); try testing.expectEqual(@as(u32, 1001), d.gid); try testing.expectEqual(@as(u32, 1), d.nlink); const f = attrFromStat(testStat("f", 0o640 | cloud9.dmappend, 1025, 4), 4, 0, 0); try testing.expectEqual(fuse.S_IFREG | 0o640, f.mode); // dmappend bit not leaked try testing.expectEqual(@as(u64, 1025), f.size); try testing.expectEqual(@as(u64, 3), f.blocks); try testing.expectEqual(@as(u32, 4096), f.blksize); try testing.expectEqual(@as(u64, 100), f.atime); try testing.expectEqual(@as(u64, 200), f.mtime); try testing.expectEqual(@as(u64, 200), f.ctime); try testing.expectEqual(@as(u64, 1), attrFromStat(testStat("f", 0o600, 512, 4), 4, 0, 0).blocks); try testing.expectEqual(@as(u64, 2), attrFromStat(testStat("f", 0o600, 513, 4), 4, 0, 0).blocks); } test "open flag → 9P mode mapping" { const rdonly: u32 = @bitCast(linux.O{ .ACCMODE = .RDONLY }); const wronly: u32 = @bitCast(linux.O{ .ACCMODE = .WRONLY }); const rdwr: u32 = @bitCast(linux.O{ .ACCMODE = .RDWR }); const trunc: u32 = @bitCast(linux.O{ .TRUNC = true }); const append: u32 = @bitCast(linux.O{ .APPEND = true }); const creat: u32 = @bitCast(linux.O{ .CREAT = true }); try testing.expectEqual(cloud9.oread, openMode(rdonly)); try testing.expectEqual(cloud9.owrite, openMode(wronly)); try testing.expectEqual(cloud9.ordwr, openMode(rdwr)); try testing.expectEqual(cloud9.owrite | cloud9.otrunc, openMode(wronly | trunc)); try testing.expectEqual(cloud9.ordwr | cloud9.otrunc, openMode(rdwr | trunc | creat)); try testing.expectEqual(cloud9.owrite, openMode(wronly | append)); // O_APPEND ignored } test "dirlist parsing from two hand-encoded Stat records" { var buf: [512]u8 = undefined; const a = try cloud9.Stat.encode(testStat("alpha", 0o644, 10, 0x11), &buf); const bb = try cloud9.Stat.encode(testStat("beta", cloud9.dmdir | 0o755, 0, 0x22), buf[a.len..]); const bytes = buf[0 .. a.len + bb.len]; // Sanity: the record is prefixed by its own 2-byte size. try testing.expectEqual(a.len - 2, std.mem.readInt(u16, bytes[0..2], .little)); var list: DirList = .{}; defer list.deinit(testing.allocator); try parseDirRecords(testing.allocator, bytes, &list); try testing.expectEqual(@as(usize, 2), list.entries.items.len); try testing.expectEqualStrings("alpha", list.entries.items[0].name); try testing.expectEqual(@as(u64, 0x11), list.entries.items[0].ino); try testing.expectEqual(fuse.DT_REG, list.entries.items[0].dtype); try testing.expectEqualStrings("beta", list.entries.items[1].name); try testing.expectEqual(@as(u64, 0x22), list.entries.items[1].ino); try testing.expectEqual(fuse.DT_DIR, list.entries.items[1].dtype); // Truncated input is a protocol error and leaves earlier entries intact. try testing.expectError(error.BadDir, parseDirRecords(testing.allocator, bytes[0 .. bytes.len - 1], &list)); try testing.expectEqual(@as(usize, 3), list.entries.items.len); } test "readdir packing and offset resumption" { const names = [_][]const u8{ ".", "..", "one", "two", "three" }; var entries: [names.len]Entry = undefined; for (&entries, names, 0..) |*e, n, i| e.* = .{ .name = @constCast(n), .ino = 100 + i, .dtype = if (i < 2) fuse.DT_DIR else fuse.DT_REG }; // Everything fits: five records, off = index + 1. var big: [1024]u8 = undefined; const used = packDirents(&entries, 0, &big); var pos: usize = 0; var idx: usize = 0; while (pos < used) : (idx += 1) { const d = std.mem.bytesToValue(fuse.Dirent, big[pos..][0..@sizeOf(fuse.Dirent)]); try testing.expectEqual(@as(u64, 100 + idx), d.ino); try testing.expectEqual(@as(u64, idx + 1), d.off); try testing.expectEqualStrings(names[idx], big[pos + @sizeOf(fuse.Dirent) ..][0..d.namelen]); pos += (@sizeOf(fuse.Dirent) + d.namelen + 7) & ~@as(usize, 7); } try testing.expectEqual(names.len, idx); // A buffer that fits exactly two records ("." = 32, ".." = 32) stops there… var small: [64]u8 = undefined; const first_used = packDirents(&entries, 0, &small); try testing.expectEqual(@as(usize, 64), first_used); const last = std.mem.bytesToValue(fuse.Dirent, small[32..][0..@sizeOf(fuse.Dirent)]); try testing.expectEqual(@as(u64, 2), last.off); // …and resuming at the last `off` yields "one" next. const second_used = packDirents(&entries, last.off, &small); const next = std.mem.bytesToValue(fuse.Dirent, small[0..@sizeOf(fuse.Dirent)]); try testing.expectEqualStrings("one", small[@sizeOf(fuse.Dirent)..][0..next.namelen]); try testing.expectEqual(@as(u64, 3), next.off); try testing.expect(second_used > 0); // Past the end: nothing (EOF for the kernel). try testing.expectEqual(@as(usize, 0), packDirents(&entries, names.len, &big)); try testing.expectEqual(@as(usize, 0), packDirents(&entries, 1000, &big)); } test "attr mapping saturates hostile lengths instead of overflowing" { const a = attrFromStat(testStat("f", 0o600, std.math.maxInt(u64), 4), 4, 0, 0); try testing.expectEqual(@as(u64, std.math.maxInt(i64)), a.size); try testing.expectEqual(@as(u64, std.math.maxInt(u64) / 512 + 1), a.blocks); const b = attrFromStat(testStat("f", 0o600, 1024, 4), 4, 0, 0); try testing.expectEqual(@as(u64, 2), b.blocks); try testing.expectEqual(@as(u64, 1024), b.size); } test "dirent names the kernel would reject are dropped from listings" { try testing.expect(validDirentName("a")); try testing.expect(validDirentName("x" ** 1024)); try testing.expect(!validDirentName("")); try testing.expect(!validDirentName("a/b")); try testing.expect(!validDirentName("a\x00b")); try testing.expect(!validDirentName(".")); try testing.expect(!validDirentName("..")); try testing.expect(!validDirentName("x" ** 1025)); var buf: [4096]u8 = undefined; var n: usize = 0; for ([_][]const u8{ ".", "..", "", "a/b", "keep", "x" ** 1025, "also" }) |name| { n += (try cloud9.Stat.encode(testStat(name, 0o644, 1, 0x30), buf[n..])).len; } var list: DirList = .{}; defer list.deinit(testing.allocator); try parseDirRecords(testing.allocator, buf[0..n], &list); try testing.expectEqual(@as(usize, 2), list.entries.items.len); try testing.expectEqualStrings("keep", list.entries.items[0].name); 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.store(7, .seq_cst); 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.store(8, .seq_cst); 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.store(9, .seq_cst); 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.store(10, .seq_cst); 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 }); list.deinit(testing.allocator); } test "mntgen node id layout: index in the top bits, local ids below" { // Index 0 is reserved for the synthetic root (node 1); mounts start at 1. try testing.expectEqual(fuse.root_id, mountNode(0, 1)); try testing.expectEqual(@as(u64, 1) << 32 | 1, mountNode(1, 1)); try testing.expectEqual(@as(u64, 7) << 32 | 12345, mountNode(7, 12345)); try testing.expectEqual(@as(u32, 0), mountIndex(fuse.root_id)); try testing.expectEqual(@as(u32, 1), mountIndex(mountNode(1, 1))); try testing.expectEqual(@as(u32, 7), mountIndex(mountNode(7, 12345))); try testing.expectEqual(@as(u32, 4095), mountIndex(4095 << 32 | 2)); // Every local id stays under the mask; the index never bleeds below it. for ([_]u64{ 1, 2, 0xFFFF_FFFF }) |local| { const node = mountNode(12, local); try testing.expectEqual(local, node & mount_node_mask); try testing.expectEqual(@as(u32, 12), mountIndex(node)); } try testing.expect(mount_node_mask == (1 << 32) - 1); } test "mntgen synthetic subdirectory node ids route apart from mounts" { for (1..max_synth_dirs + 1) |i| { const node = mountNode(synth_index, i); try testing.expectEqual(synth_index, mountIndex(node)); try testing.expectEqual(i, node & mount_node_mask); } // The reserved index can never be a mount ordinal. try testing.expect(synth_index >= max_mounts); } test "mntgen synthetic-root inos are deterministic and name-derived" { const a = nameIno("alpha"); try testing.expectEqual(a, nameIno("alpha")); try testing.expect(a != nameIno("beta")); try testing.expect(a != 0); try testing.expect(nameIno("") != nameIno("x")); } test "mntgen: ino_xor keeps two servers' identical qid.paths distinct" { // Two ramfs instances hand out the same qid.path; without the mount // index mixed in, `find` would see one file twice as a hardlink. var b1: Bridge = .{ .gpa = testing.allocator, .fuse_fd = -1, .nine = undefined, .opts = .{ .uid = 0, .gid = 0 }, .root_id = mountNode(1, 1), .ino_xor = @as(u64, 1 + 1) << 48 }; var b2: Bridge = .{ .gpa = testing.allocator, .fuse_fd = -1, .nine = undefined, .opts = .{ .uid = 0, .gid = 0 }, .root_id = mountNode(2, 1), .ino_xor = @as(u64, 2 + 1) << 48 }; const qid = cloud9.Qid{ .type = 0, .version = 0, .path = 0x11 }; const ino1 = b1.inoOf(1, qid); const ino2 = b2.inoOf(1, qid); try testing.expect(ino1 != ino2); // Within one mount the qid.path still decides (same file two ways = one ino). try testing.expectEqual(ino1, b1.inoOf(2, qid)); // And the single-connection mode is unchanged (xor 0). var b0: Bridge = .{ .gpa = testing.allocator, .fuse_fd = -1, .nine = undefined, .opts = .{ .uid = 0, .gid = 0 } }; try testing.expectEqual(inoFromPath(qid.path), b0.inoOf(1, qid)); }