From 3a23f6a29e47ace901bd4d82b9db4055fcc12bb9 Mon Sep 17 00:00:00 2001 From: Gabriel Schneider Date: Mon, 21 Sep 2026 14:13:43 -0300 Subject: post registry + 9ns --mntgen: the /srv translation cloud9.post: servers post their socket under a name in $XDG_RUNTIME_DIR/9p (post/unpost, posted, dial, Watch) and serve.Runner.listenPosted posts a server by name, unposting on stop. Names are budget-checked against the 108-byte socket path; a claim binds+listens at a private temp path and takes the name with atomic renames under flock (RENAME_NOREPLACE for free names, RENAME_EXCHANGE grab-verify-commit for stale ones): the registry path is never unlinked by a claim, live names refuse with AlreadyPosted, foreign files with NotSocket, and unpost removes only the caller's inode-matched entry. Watch surfaces inotify overflow and a replaced registry dir. 9ns --mntgen [--mount DIR] -- PROGRAM: one FUSE mount at /mnt/9p whose synthetic root lists the posted registry (no connection made); a walk into an unmounted name dials it and runs the existing bridge dispatch in a per-server worker thread, routed by mount index in the node id's top bits (ordinals never reused, cap 4096); a dead server answers EIO on its subtree and is re-dialed on the next walk. The dial watches stop_fd through Tversion (connectWatched). All existing 9ns forms are unchanged. 9proc's unix listener no longer blind-unlinks its path: a foreign non-socket is refused (Occupied), a live server is refused (AlreadyListening), only a refused socket is cleared, and stop() unlinks only the listener's own inode-matched socket. Hardened by adversarial review (GLM 5.3 x2 + DeepSeek V4.1 Flash, all high-thinking): double-bind races on one name (0 in 180k rounds), foreign-file TOCTOU deletions (0 in 4M flips), a 255-byte-name listing panic, inotify queue overflow silently dropped, listenPosted silently overwriting, dial-time Tversion hangs wedging the dispatcher, --debug silently ignored in mntgen, and xattr/statx probes answering EPERM on the synthetic root (broke `ls -l /mnt/9p`). Tests: root 80/80, 9ns 47/47, 9proc 60/60, integration 88/88 + mntgen 37/37, adversarial 213/0, freestanding riscv32 gate green. --- 9ns/src/bridge.zig | 750 +++++++++++++++++++++++++++++++++++++++++++++++++++-- 9ns/src/main.zig | 116 ++++++++- 9ns/src/nine.zig | 55 ++++ 3 files changed, 896 insertions(+), 25 deletions(-) (limited to '9ns/src') 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)); +} diff --git a/9ns/src/main.zig b/9ns/src/main.zig index 076aa42..386552f 100644 --- a/9ns/src/main.zig +++ b/9ns/src/main.zig @@ -7,6 +7,7 @@ const std = @import("std"); const linux = std.os.linux; +const cloud9 = @import("cloud9"); const ns = @import("ns.zig"); const nine = @import("nine.zig"); const bridge = @import("bridge.zig"); @@ -20,10 +21,16 @@ const usage_text = \\ --tcp IP:PORT TCP (IPv4/IPv6 literal) \\ --fd N already-connected inherited descriptor \\ --spawn CMD run CMD (via /bin/sh -c) with a socketpair on its stdin/stdout + \\ --mntgen mount the posted-9P registry ($XDG_RUNTIME_DIR/9p): one + \\ mount whose root lists the posted names; walking into a + \\ name dials that server (mutually exclusive with the rest) \\Options: \\ --name NAME mount name: the tree appears at /mnt/9p/NAME (one path - \\ component; default derived from the transport, see below) - \\ --mount PATH mountpoint inside the new namespace (overrides --name) + \\ component; default derived from the transport, see below; + \\ not with --mntgen) + \\ --mount PATH mountpoint inside the new namespace (overrides --name; + \\ with --mntgen the mount is the registry view itself, + \\ default /mnt/9p) \\ --uname NAME 9P user name (default $USER, else "none") \\ --aname NAME 9P tree to attach (default "") \\ --msize BYTES maximum 9P message size to request (default 131072) @@ -35,7 +42,7 @@ const usage_text = \\Default name: --unix PATH -> basename of PATH without .sock/.9p/.socket; \\--tcp IP:PORT -> tcp-IP-PORT (':' becomes '-'); --spawn CMD -> basename of its \\first word; --fd N -> fdN; 9p when nothing usable comes out of that. - \\ + \\--mntgen: no per-server name; the registry mount goes to --mount (default /mnt/9p). ; /// Where `--name NAME` mounts: `mount_root/NAME`. @@ -61,10 +68,12 @@ fn printStdout(text: []const u8) void { } } } - const Config = struct { address: ?nine.Address = null, spawn_cmd: ?[]const u8 = null, + /// `--mntgen`: the mount lists the posted-9P registry and dials servers + /// lazily (see `runMntgen`); mutually exclusive with the transports. + mntgen: bool = false, /// `--mount`: wins over `name` when set. mount: ?[]const u8 = null, /// `--name`: null means "derive from the transport" (see `defaultName`). @@ -119,15 +128,15 @@ fn parseArgs(arena: std.mem.Allocator, args: []const [:0]const u8) !ParseResult name = arg[0..eq]; inline_value = arg[eq + 1 ..]; } - const Opt = enum { unix, tcp, fd, spawn, name, mount, uname, aname, msize, cache, @"no-direct-io", debug, help, version, unknown }; + const Opt = enum { unix, tcp, fd, spawn, mntgen, name, mount, uname, aname, msize, cache, @"no-direct-io", debug, help, version, unknown }; const opt = std.meta.stringToEnum(Opt, name[2..]) orelse .unknown; switch (opt) { - .@"no-direct-io", .debug, .help, .version => if (inline_value != null) return usageError("{s} takes no value", .{name}), + .mntgen, .@"no-direct-io", .debug, .help, .version => if (inline_value != null) return usageError("{s} takes no value", .{name}), .unknown => return usageError("unknown option {s}", .{name}), else => {}, } const value: []const u8 = switch (opt) { - .@"no-direct-io", .debug, .help, .version, .unknown => "", + .mntgen, .@"no-direct-io", .debug, .help, .version, .unknown => "", else => inline_value orelse blk: { i += 1; if (i >= args.len) return usageError("{s} needs a value", .{name}); @@ -155,6 +164,10 @@ fn parseArgs(arena: std.mem.Allocator, args: []const [:0]const u8) !ParseResult cfg.spawn_cmd = value; transports += 1; }, + .mntgen => { + cfg.mntgen = true; + transports += 1; + }, .name => { if (!validName(value)) return usageError("--name wants a single path component (not empty, no '/', not . or ..), got '{s}'", .{value}); cfg.name = value; @@ -181,7 +194,8 @@ fn parseArgs(arena: std.mem.Allocator, args: []const [:0]const u8) !ParseResult .unknown => unreachable, } } - if (transports == 0) return usageError("one transport is required (--unix, --tcp, --fd or --spawn)", .{}); + if (cfg.mntgen and cfg.name != null) return usageError("--name is not meaningful with --mntgen (the registry mount is --mount, default {s})", .{mount_root}); + if (transports == 0) return usageError("one transport is required (--unix, --tcp, --fd, --spawn or --mntgen)", .{}); if (transports > 1) return usageError("exactly one transport is allowed", .{}); if (program_start) |start| { const prog = try arena.alloc([]const u8, args.len - start); @@ -338,6 +352,71 @@ fn describeAddress(a: nine.Address, buf: []u8) []const u8 { .fd => |fd| std.fmt.bufPrint(buf, "fd {d}", .{fd}) catch "fd", }; } +/// `--mntgen`: one FUSE mount whose root lists the posted-9P registry +/// (`$XDG_RUNTIME_DIR/9p`, the /srv translation of cloud9.post). Servers +/// are dialed lazily when the program walks into their name; see +/// `bridge.serveMntgen` for the process model. Fails before anything is +/// forked when XDG_RUNTIME_DIR is unset or /dev/fuse is unusable. +fn runMntgen(init: std.process.Init, envp: [*:null]const ?[*:0]const u8, cfg: Config, uname: []const u8) !u8 { + const gpa = init.gpa; + const mount_arg = cfg.mount orelse mount_root; + const mountpoint = ns.resolveMountpoint(gpa, mount_arg) catch |err| { + std.debug.print("9ns: --mount {s}: {t}\n", .{ mount_arg, err }); + return own_failure; + }; + defer gpa.free(mountpoint); + + // The registry must be nameable before anything is forked; there is no + // fallback directory (post.registryDir errors rather than guess /tmp). + var reg_buf: [128]u8 = undefined; + _ = cloud9.post.registryDir(envp, ®_buf) catch |err| { + std.debug.print("9ns: --mntgen: {t} (the posted-9P registry is $XDG_RUNTIME_DIR/9p)\n", .{err}); + return own_failure; + }; + + if (!probeFuseDevice()) return own_failure; + + // Writes to a dead server socket must not kill us. + ignoreSignal(.PIPE); + + var child_pid: i32 = 0; + const stop_fd = ns.installSignals(&child_pid) catch return own_failure; + + const uid = linux.getuid(); + const gid = linux.getgid(); + const child = ns.spawn(gpa, .{ + .argv = cfg.program, + .envp = envp, + .mountpoint = mountpoint, + .uid = uid, + .gid = gid, + .max_read = bridge.max_write, + }) catch return own_failure; + + bridge.serveMntgen(gpa, child.fuse_fd, stop_fd, .{ + .io = init.io, + .env = envp, + .uname = uname, + .aname = cfg.aname, + .msize = cfg.msize, + }, .{ + .uid = uid, + .gid = gid, + .attr_timeout_ns = cfg.cache_ns, + .direct_io = cfg.direct_io, + .debug = cfg.debug, + }) catch |err| { + std.debug.print("9ns: fuse: {t}\n", .{err}); + }; + + // Closing the device aborts the FUSE connection: anything still using + // the mount gets ENOTCONN instead of hanging on an unserved request. + _ = linux.close(child.fuse_fd); + + const status = ns.reapIfExited(child.pid) orelse ns.waitChild(child.pid) catch own_failure; + _ = ns.reportExecFailure(child); + return status; +} pub fn main(init: std.process.Init) !u8 { const gpa = init.gpa; @@ -361,6 +440,7 @@ pub fn main(init: std.process.Init) !u8 { cfg.program = try arena.dupe([]const u8, &.{shell}); } const uname = cfg.uname orelse ns.getenv(envp, "USER") orelse "none"; + if (cfg.mntgen) return runMntgen(init, envp, cfg, uname); // `--mount PATH` wins; otherwise `/mnt/9p/` with `--name` or a // name derived from the transport. var name_buf: [512]u8 = undefined; @@ -515,6 +595,26 @@ test "parseArgs" { const okmsize = [_][:0]const u8{ "9ns", "--fd", "3", "--msize", "16777216" }; try std.testing.expectEqual(@as(u32, 16777216), (try parseArgs(arena, &okmsize)).run.msize); } + { + // --mntgen is a transport: exclusive with the others, no value, + // --name rejected, options still apply to the per-server dials. + const ok = [_][:0]const u8{ "9ns", "--mntgen", "--mount", "/m", "--msize=8192", "--", "sh" }; + const r = try parseArgs(arena, &ok); + defer arena.free(r.run.program); + try std.testing.expect(r.run.mntgen); + try std.testing.expect(r.run.address == null); + try std.testing.expectEqualStrings("/m", r.run.mount.?); + try std.testing.expectEqual(@as(u32, 8192), r.run.msize); + try std.testing.expectEqual(@as(usize, 1), r.run.program.len); + const withunix = [_][:0]const u8{ "9ns", "--mntgen", "--unix", "/s", "--", "sh" }; + try std.testing.expectEqual(@as(u8, 125), (try parseArgs(arena, &withunix)).exit); + const withspawn = [_][:0]const u8{ "9ns", "--spawn", "x", "--mntgen", "--", "sh" }; + try std.testing.expectEqual(@as(u8, 125), (try parseArgs(arena, &withspawn)).exit); + const withname = [_][:0]const u8{ "9ns", "--mntgen", "--name", "foo", "--", "sh" }; + try std.testing.expectEqual(@as(u8, 125), (try parseArgs(arena, &withname)).exit); + const withvalue = [_][:0]const u8{ "9ns", "--mntgen=x", "--", "sh" }; + try std.testing.expectEqual(@as(u8, 125), (try parseArgs(arena, &withvalue)).exit); + } { const ver = [_][:0]const u8{ "9ns", "--version" }; try std.testing.expectEqualStrings(version_string ++ "\n", (try parseArgs(arena, &ver)).info); diff --git a/9ns/src/nine.zig b/9ns/src/nine.zig index c89a343..cf9c49a 100644 --- a/9ns/src/nine.zig +++ b/9ns/src/nine.zig @@ -75,6 +75,17 @@ pub const Session = struct { /// Connect to `address`, then negotiate the protocol version. /// `msize` is the maximum message size to ask for (0 = the buffers' size). pub fn connect(gpa: std.mem.Allocator, address: Address, msize: u32) !Session { + return connectWatched(gpa, address, msize, -1); + } + + /// `connect`, with `stop_fd` watched for the whole handshake (the version + /// rpc included). A server that accepts the connection and then never + /// answers the Tversion would otherwise pin the caller in a blocking read + /// with no way out: 9ns's mntgen dispatcher dials on the strength of the + /// program's walk, so it must come back when that program is gone. The + /// field stays set on the returned session, so the attach and stat that + /// follow a dial keep watching it too; -1 disables the watch. + pub fn connectWatched(gpa: std.mem.Allocator, address: Address, msize: u32, stop_fd: i32) !Session { const want: u32 = if (msize == 0) 8192 else @max(msize, 24); const fd = try openTransport(address); errdefer if (address != .fd) { @@ -93,6 +104,7 @@ pub const Session = struct { .in_buf = in_buf, .out_buf = out_buf, .msize = want, + .stop_fd = stop_fd, }; const r = try s.rpc(.{ .version = .{ .msize = want } }); if (!std.mem.eql(u8, r.version.version, "9P2000")) return error.Protocol; @@ -1037,6 +1049,49 @@ test "rpc wait loop: a read interrupted after some data is a short read" { try testing.expectEqual(@as(usize, 0), s.client.pending()); } +test "connectWatched: a silent server cannot pin the handshake past stop_fd" { + // A server that accepts and then never answers: the version handshake has + // nothing to read. With a readable stop_fd the connect must come back with + // error.Stopped instead of blocking in readSocket (the fd is blocking), and + // the caller's descriptor must survive: `Address.fd` is not ours to close. + var sv: [2]i32 = undefined; + try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &sv))); + defer _ = linux.close(sv[0]); + defer _ = linux.close(sv[1]); + var p: [2]i32 = undefined; + try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.pipe2(&p, .{ .CLOEXEC = true, .NONBLOCK = true }))); + defer _ = linux.close(p[0]); + defer _ = linux.close(p[1]); + try testing.expectEqual(@as(usize, 1), linux.write(p[1], "x", 1)); + + const Probe = struct { + const Self = @This(); + done: std.atomic.Value(bool) = .init(false), + stopped: std.atomic.Value(bool) = .init(false), + fd_open: std.atomic.Value(bool) = .init(false), + + fn run(w: *Self, client: i32, stop: i32) void { + if (Session.connectWatched(testing.allocator, .{ .fd = client }, 8192, stop)) |session| { + var s = session; + s.deinit(); + } else |e| w.stopped.store(e == error.Stopped, .release); + w.fd_open.store(linux.errno(linux.fcntl(client, linux.F.GETFD, 0)) == .SUCCESS, .release); + w.done.store(true, .release); + } + }; + var w: Probe = .{}; + const th = try std.Thread.spawn(.{}, Probe.run, .{ &w, sv[0], p[0] }); + var waited_ms: usize = 0; + while (!w.done.load(.acquire) and waited_ms < 3000) : (waited_ms += 10) { + const ts: linux.timespec = .{ .sec = 0, .nsec = 10 * std.time.ns_per_ms }; + _ = linux.nanosleep(&ts, null); + } + try testing.expect(w.done.load(.acquire)); + try testing.expect(w.stopped.load(.acquire)); + try testing.expect(w.fd_open.load(.acquire)); + th.join(); +} + test "session against an in-process cloud9.Server" { var fds: [2]i32 = undefined; try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &fds))); -- cgit v1.3