summaryrefslogtreecommitdiff
path: root/9ns/src
diff options
context:
space:
mode:
Diffstat (limited to '9ns/src')
-rw-r--r--9ns/src/bridge.zig750
-rw-r--r--9ns/src/main.zig116
-rw-r--r--9ns/src/nine.zig55
3 files changed, 896 insertions, 25 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));
+}
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, &reg_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/<name>` with `--name` or a
// name derived from the transport.
var name_buf: [512]u8 = undefined;
@@ -516,6 +596,26 @@ test "parseArgs" {
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);
const help = [_][:0]const u8{ "9ns", "--help" };
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)));