summaryrefslogtreecommitdiff
path: root/9ns/src/bridge.zig
diff options
context:
space:
mode:
Diffstat (limited to '9ns/src/bridge.zig')
-rw-r--r--9ns/src/bridge.zig750
1 files changed, 733 insertions, 17 deletions
diff --git a/9ns/src/bridge.zig b/9ns/src/bridge.zig
index 089c3ca..76327c3 100644
--- a/9ns/src/bridge.zig
+++ b/9ns/src/bridge.zig
@@ -12,6 +12,7 @@
//! 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;
@@ -118,13 +119,13 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_
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, fuse.root_id);
+ 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, fuse.root_id, .{ .fid = root_fid, .qid = root_qid, .nlookup = 1, .parent = fuse.root_id });
+ 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 },
@@ -172,6 +173,14 @@ const Bridge = struct {
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.
@@ -181,9 +190,10 @@ const Bridge = struct {
/// 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,
+ /// `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,
@@ -249,8 +259,8 @@ const Bridge = struct {
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) {
+ 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;
@@ -288,9 +298,9 @@ const Bridge = struct {
b.reply(h.unique, &.{}) catch {};
return false;
}
- b.cur_unique = h.unique;
+ b.cur_unique.store(h.unique, .seq_cst);
b.interrupted = false;
- defer b.cur_unique = null;
+ 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.?;
@@ -449,7 +459,7 @@ const Bridge = struct {
ino.nlookup += 1;
ino.qid = qid;
nodeid = existing;
- if (existing == fuse.root_id) {
+ if (existing == b.root_id) {
b.clunkQuiet(newfid);
} else {
// Keep the fresh fid (it is bound to the current file at this
@@ -493,7 +503,7 @@ const Bridge = struct {
}
fn forget(b: *Bridge, nodeid: u64, n: u64) HandlerError!void {
- if (nodeid == fuse.root_id) return;
+ if (nodeid == b.root_id) return;
const ino = b.inodes.getPtr(nodeid) orelse return;
if (ino.nlookup > n) {
ino.nlookup -= n;
@@ -794,9 +804,8 @@ const Bridge = struct {
/// 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 {
- _ = b;
_ = nodeid;
- return inoFromPath(qid.path);
+ return inoFromPath(qid.path) ^ b.ino_xor;
}
fn inoOfNode(b: *const Bridge, nodeid: u64) u64 {
@@ -817,6 +826,671 @@ const Bridge = struct {
}
};
+// ===========================================================================
+// 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;
+
+/// 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 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),
+
+ 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;
+ }
+ }
+ 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;
+ },
+ else => {},
+ }
+ if (h.nodeid == fuse.root_id) {
+ try mg.handleRoot(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 {};
+ }
+ }
+
+ /// 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)});
+ }
+ const m = mg.dialMount(name) catch |e| switch (e) {
+ error.NotPosted => {
+ mg.trace(" lookup '{s}': nothing posted under that name", .{name});
+ return mg.replyError(u, .NOENT);
+ },
+ 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;
+ }
+
+ // -- 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 {
+ if (mg.next_index >= max_mounts) return error.TooManyMounts;
+ const index: u32 = mg.next_index;
+
+ const stream = try post.dial(mg.mo.io, mg.mo.env, name);
+ 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, name),
+ .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.
@@ -1106,7 +1780,7 @@ test "interrupt source: INTERRUPT in flight → Tflush, other INTERRUPTs ignored
// 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;
+ b.cur_unique.store(7, .seq_cst);
pi.unique = 7;
try pi.inject(6);
try testing.expectError(error.Interrupted, b.read(1, 0, &buf));
@@ -1116,7 +1790,7 @@ test "interrupt source: INTERRUPT in flight → Tflush, other INTERRUPTs ignored
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;
+ 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());
@@ -1127,7 +1801,7 @@ test "interrupt source: INTERRUPT in flight → Tflush, other INTERRUPTs ignored
@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;
+ 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;
@@ -1143,7 +1817,7 @@ test "interrupt source: INTERRUPT in flight → Tflush, other INTERRUPTs ignored
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;
+ 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);
@@ -1155,3 +1829,45 @@ test "DirList frees its names" {
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-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));
+}