summaryrefslogtreecommitdiff
path: root/9ns
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-21 14:13:43 -0300
committerGabriel Schneider <[email protected]>2026-09-21 14:13:43 -0300
commit3a23f6a29e47ace901bd4d82b9db4055fcc12bb9 (patch)
treeb82d6e7c3ebe108434ce00ca75db59cf037917e0 /9ns
parentf1b53c1533539aecbf16ad19fd9156deae091f92 (diff)
downloadcloud9-3a23f6a29e47ace901bd4d82b9db4055fcc12bb9.tar.gz
cloud9-3a23f6a29e47ace901bd4d82b9db4055fcc12bb9.zip
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.
Diffstat (limited to '9ns')
-rw-r--r--9ns/build.zig6
-rw-r--r--9ns/docs/DESIGN.md187
-rw-r--r--9ns/src/bridge.zig750
-rw-r--r--9ns/src/main.zig116
-rw-r--r--9ns/src/nine.zig55
-rwxr-xr-x9ns/test/mntgen.sh222
6 files changed, 1297 insertions, 39 deletions
diff --git a/9ns/build.zig b/9ns/build.zig
index 7f2155f..4e8d23b 100644
--- a/9ns/build.zig
+++ b/9ns/build.zig
@@ -52,7 +52,9 @@ pub fn add(b: *std.Build, ctx: Context) Artifacts {
test_step.dependOn(&b.addRunArtifact(b.addTest(.{ .root_module = ns_mod })).step);
// End-to-end suites: real namespaces, real FUSE, a real 9P server.
- const itest = b.step("9ns-itest", "Run 9ns/test/integration.sh (needs unprivileged user namespaces and /dev/fuse)");
+ // integration.sh: the single-connection transports; mntgen.sh: the
+ // posted-registry multi-server mode (registry servers + lazy dial).
+ const itest = b.step("9ns-itest", "Run 9ns/test/integration.sh and 9ns/test/mntgen.sh (needs unprivileged user namespaces and /dev/fuse)");
const adv = b.step("9ns-adv", "Run 9ns/test/adversarial.sh (hostile servers, namespaces, stress; several minutes)");
const demo = ctx.proc_demo orelse {
const fail = b.addFail("9ns-itest and 9ns-adv need the 9proc-demo server (build with -D9proc=true)");
@@ -60,7 +62,7 @@ pub fn add(b: *std.Build, ctx: Context) Artifacts {
adv.dependOn(&fail.step);
return .{ .exe = ns, .test_step = test_step, .itest_step = itest, .adv_step = adv };
};
- inline for (.{ .{ itest, "integration" }, .{ adv, "adversarial" } }) |pair| {
+ inline for (.{ .{ itest, "integration" }, .{ itest, "mntgen" }, .{ adv, "adversarial" } }) |pair| {
const run = b.addSystemCommand(&.{"bash"});
run.addFileArg(b.path("9ns/test/" ++ pair[1] ++ ".sh"));
run.addArtifactArg(ns);
diff --git a/9ns/docs/DESIGN.md b/9ns/docs/DESIGN.md
index f1589d0..1d66387 100644
--- a/9ns/docs/DESIGN.md
+++ b/9ns/docs/DESIGN.md
@@ -102,6 +102,101 @@ protocol we need directly against `/usr/include/linux/fuse.h`.
The FUSE fd is shared with the child only until exec (CLOEXEC); the parent's
copy keeps the connection alive.
+### mntgen: one mount, many servers (`9ns --mntgen`)
+
+`9ns --mntgen [--mount DIR] -- PROGRAM` (default mountpoint `/mnt/9p`) is a
+transport of its own, mutually exclusive with `--unix/--tcp/--fd/--spawn`;
+`--name` is rejected (there is no single server to name) while `--uname`,
+`--aname`, `--msize`, `--cache`, `--no-direct-io` and `--debug` apply to
+every per-server dial. The mount it builds is the `/srv` view of the
+posted-9P registry (`cloud9.post`, `$XDG_RUNTIME_DIR/9p`): a walk into a
+posted name reaches that server's whole 9P tree, and nothing is connected
+until something walks. XDG_RUNTIME_DIR unset is fatal before anything is
+forked (the registry is not guessable; no `/tmp` fallback).
+
+```
+ program 9ns parent
+ in new userns │
+ /mnt/9p ─FUSE─▶ kernel ─▶ │ dispatcher (main thread)
+ alpha/ beta/ │ ├─ synthetic root (node 1): lists the registry
+ ...each a server │ ├─ LOOKUP(alpha) ── dial+attach+stat ─▶ worker 1 ─ bridge ─ 9P session ─ server alpha
+ │ └─ LOOKUP(beta) ── dial+attach+stat ─▶ worker 2 ─ bridge ─ 9P session ─ server beta
+```
+
+Threads and node ids:
+
+* **Dispatcher** (the main thread) is the only reader of `/dev/fuse`. It
+ answers INIT/DESTROY, serves the **synthetic root** (node 1) itself, and
+ routes every other request by the node id's top bits — the **mount
+ index** — to the owning server's mount; a FUSE_INTERRUPT is routed by
+ scanning the mounts for the one currently serving the interrupted unique
+ and dropping the target into its interrupt pipe.
+* A **mount** is a dialed server: its own `nine.Session`, its own bridge
+ state (inode table, handles — the ordinary single-server translation,
+ unchanged) and a **worker thread**. The worker pops copied requests off a
+ queue and serves them one at a time, exactly the single-connection
+ contract; the dispatcher keeps reading `/dev/fuse` meanwhile, so one slow
+ server never blocks the other names. Replies go straight back on the FUSE
+ fd (one `writev` per reply; the kernel processes each write as one
+ message).
+* Node id layout: `nodeid = (mount_index << 32) | local`. Index 0 is the
+ synthetic root; per-mount local ids start at 1 (the server's 9P root) and
+ never exceed 2^32 (a bridge never reuses one). Mount indexes are
+ **ordinals and are never reused** (cap 4096 per process), so a stale
+ kernel-side inode of a dead server can never be conflated with a fresh
+ inode of its replacement. Reported `st_ino` mixes the index into the
+ qid.path (`(index+1) << 48` XOR), so two servers handing out the same
+ qid.path (two ramfs instances) still get distinct inode numbers.
+* **Lazy dial**: LOOKUP of an unmounted name checks the registry, dials,
+ attaches, stats the root and spawns the worker — all on the dispatcher
+ thread, in service of the walk that triggered it. No eager connection is
+ ever made: `ls` of the root reads the registry directory only (a plain
+ file dropped there is listed too — and yields EIO on the walk, never
+ deleted). A walk into a **stale** entry (socket present, connect refused)
+ answers EIO. The whole dial watches `stop_fd` (the session is built with
+ `nine.Session.connectWatched`, so the `Tversion` exchange is covered too):
+ when the program exits while a walk is parked in a dial, the dial fails
+ with `Stopped`, the pending LOOKUP answers EIO and 9ns follows the program
+ out. There is no dial timeout of our own (a slow server delays the walk,
+ like it would delay any 9P client), and a server that accepts but never
+ answers `Tversion` still parks the dispatcher until the program exits —
+ including the unkillable corner where the *blocked walk itself* is the
+ only thing keeping the program alive (the task sits in D state until the
+ filesystem answers; a same-user self-DoS, accepted with the pinned
+ "dispatcher dials, no concurrent dial" design).
+* **Death and re-dial**: when a worker's session dies mid-request, dispatch
+ has already answered that request EIO, the mount is marked dead, and
+ everything further routed to that subtree answers EIO (a FORGET is
+ dropped). Nothing reconnects eagerly. Because synthetic-root entries are
+ served with zero entry-validity, the next walk into the name LOOKUPs it
+ again; a dead mount is skipped and the name is dialed afresh — a new
+ mount under a new index, so kernel-held inodes of the corpse keep
+ answering EIO until forgotten. `ls` still lists the dead name (listing
+ connects to nothing). Death is discovered lazily: the first walk after a
+ silent death answers EIO (it marks the mount dead), the next walk re-dials.
+* **Interrupts** work per mount: the worker's session polls the mount's
+ interrupt pipe while a 9P reply is outstanding; the dispatcher writes the
+ interrupted request's unique into it and the usual `Tflush` dance
+ (see *Interrupts*) follows. `stop_fd` (the child's death) is watched by
+ every session, so no worker can stay blocked on a hung server past the
+ program's exit; teardown wakes every worker, joins them, and tears down
+ their sessions and bridge state.
+* The synthetic root is read-only (`dr-xr-xr-x`, like `/srv`): services are
+ posted and unposted by their servers (`cloud9.post`'s
+ `post`/`listenPosted`/`unpost`), not created and removed through files.
+ Capability probes the kernel makes before trusting a file — xattr ops and
+ `STATX` — answer `ENOSYS` (as they do on server subtrees through the
+ ordinary dispatch): `EPERM` would leak into userland as
+ "Operation not permitted" blamed on the mount root by `ls -l` and plain
+ `stat`. Modifying ops (create, mkdir, rename, ...) keep `EPERM` — the
+ root's read-only nature — and unknown ops are refused, never fatal.
+ Root readdir snapshots the registry per OPENDIR (a new OPENDIR sees new
+ posts); because nothing the kernel can cache is ever served stale (zero
+ timeouts, snapshot per opendir), there is nothing to invalidate and no
+ watcher is needed (`post.Watch` exists in the library for future caching).
+* Naming: `--mount` (default `/mnt/9p`) is the whole mount; `$NINE_MOUNT`
+ points at it as usual.
+
### Mountpoint policy
Default mountpoint: `/mnt/9p/<name>`, where the name is `--name NAME` or is
@@ -286,12 +381,33 @@ pub const Options = struct {
};
/// Runs until the FUSE fd reports ENODEV or `stop_fd` becomes readable.
pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, nine: *nine.Session, root_fid: u32, stop_fd: i32, opts: Options) !void;
+
+// mntgen (see "mntgen: one mount, many servers" under Process model):
+pub const MntgenOptions = struct {
+ io: std.Io, // dispatcher-thread only: post.posted / post.dial
+ env: post.Env, // XDG_RUNTIME_DIR names the registry
+ uname: []const u8, aname: []const u8 = "", msize: u32 = 131072,
+};
+/// One FUSE mount whose root lists the posted-9P registry; servers dialed
+/// lazily, one worker thread each; node ids carry the mount index in the
+/// top bits (`mount_shift = 32`, indexes are ordinals from 1, never reused,
+/// `max_mounts` = 4096; the registry snapshot buffer is 8 KiB, `max_root_dirs`
+/// = 64 concurrent OPENDIRs). Runs until ENODEV, DESTROY or `stop_fd`.
+pub fn serveMntgen(gpa: std.mem.Allocator, fuse_fd: i32, stop_fd: i32, mo: MntgenOptions, opts: Options) !void;
+pub fn mountNode(index: u32, local: u64) u64; // (index << 32) | local
+pub fn mountIndex(nodeid: u64) u32; // nodeid >> 32
+pub fn nameIno(name: []const u8) u64; // FNV-1a of a name: synthetic-root dirent inos
```
State:
* `inodes: AutoHashMap(u64 /*nodeid*/, Inode{ fid: u32, qid: Qid, nlookup: u64 })`.
- Node 1 is the root (`root_fid`, never forgotten).
+ Node 1 is the root (`root_fid`, never forgotten). In mntgen mode the
+ same machinery runs once per mount with node ids that already carry the
+ mount index; the root node id is the `Bridge.root_id` field (`fuse.root_id`
+ single-connection) and reported inode numbers are `qid.path ^ ino_xor`
+ (`ino_xor` 0 single-connection). `cur_unique` is atomic so the mntgen
+ dispatcher can scan it to route FUSE_INTERRUPTs.
* `by_qid: AutoHashMap(u64 /*qid.path*/, u64 /*nodeid*/)` so that repeated
lookups of the same file map to the same inode (the old fid is clunked and
the fresh one kept). Dedupe only merges when the qid type (dir bit) also
@@ -431,10 +547,16 @@ Transport (exactly one):
--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, max 16 MiB)
@@ -446,11 +568,14 @@ PROGRAM defaults to $SHELL (else /bin/sh). The mountpoint is exported as $NINE_M
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).
```
`--name` and `--mount` may both be given; `--mount` wins. Exit codes: child's
status; 125 for 9ns's own failures (usage including a bad `--name`, connect,
-mount); 126/127 as usual for exec failures.
+mount, and for `--mntgen` an unset XDG_RUNTIME_DIR); 126/127 as usual for
+exec failures.
+
### `../9proc/demo/main.zig` — demo 9P2000 server (binary `9proc-demo`)
@@ -491,12 +616,26 @@ Everything under a temp dir. Skips (exit 0 with a notice) when
is untouched).
6. Kill tests: 9ns exits when the child exits; server death during use
yields `EIO`, not a hang.
+7. `test/mntgen.sh` (also `9ns-itest`): `--mntgen` against a scratch
+ registry (`XDG_RUNTIME_DIR` = temp dir, never the real one): two servers
+ posted under two names (9proc-demo's socket created inside the registry
+ directory — a socket at `$XDG_RUNTIME_DIR/9p/<name>` is a posted name —
+ plus plan9port `ramfs`, an independent 9P2000 implementation, posting
+ with `NAMESPACE=$REG`), the synthetic root listing without dialing
+ (including a non-socket file, which is listed, yields EIO on the walk,
+ and is never removed), lazy dial and per-server routing in one program
+ run, a post appearing after mount, server death → EIO on the subtree
+ with the name still listed and a re-dial after re-post, parallel reads on
+ both mounts, `$NINE_MOUNT` = `/mnt/9p` by default, and the usage errors
+ (`--mntgen` + any transport, `--name`, `--mntgen=x`, unset
+ `XDG_RUNTIME_DIR`).
## Verification
-`zig build 9ns-test` (unit), `zig build 9ns-itest` (88 end-to-end
-checks against 9proc-demo over unix/tcp/socketpair and against plan9port's
-`ramfs`) and `zig build 9ns-adv` (adversarial suites: a scriptable
+`zig build 9ns-test` (unit), `zig build 9ns-itest` (integration.sh + the
+mntgen suite: end-to-end checks against 9proc-demo over unix/tcp/socketpair,
+plan9port's `ramfs`, and the posted-registry multi-server mode) and
+`zig build 9ns-adv` (adversarial suites: a scriptable
hostile 9P server with ~30 misbehaviour modes, interrupt forwarding against
its `never_flush`/`never` modes with 28 checks, FUSE semantics through the
bridge, process/namespace/signal edge cases with 51 checks, and stress). The
@@ -507,11 +646,35 @@ ReleaseSafe.
## Out of scope for v1 (documented, not hidden)
-* One 9P request in flight at a time: a 9P read that blocks (event files)
- stalls the whole mount while it is outstanding (but not past the child's
- exit). It can be interrupted: killing or Ctrl-C-ing the reader sends
- `FUSE_INTERRUPT`, which becomes `Tflush`; servers that honour it unblock
- immediately, servers that don't still block the mount until they answer.
+* One 9P request in flight at a time, per server in mntgen mode (across
+ servers they proceed in parallel, one worker each): a 9P read that blocks
+ (event files) stalls that server's subtree while it is outstanding (but
+ not past the child's exit). It can be interrupted: killing or Ctrl-C-ing
+ the reader sends `FUSE_INTERRUPT`, which becomes `Tflush`; servers that
+ honour it unblock immediately, servers that don't still block that
+ subtree until they answer.
+* mntgen: the dispatcher dials on the main thread, in service of the walk
+ that triggered it: a server that accepts the connection but never
+ answers `Tversion` parks the whole mount for as long as the walk's
+ program keeps running (as it would delay any 9P client). The dial
+ watches `stop_fd`, so the program exiting ends it (`Stopped`, EIO to the
+ pending LOOKUP); no concurrent dial, no dial timeout of our own. The
+ remaining corner is a same-user self-DoS: the program blocked *in that
+ very walk* sits in D state until the filesystem answers and cannot be
+ killed to fire `stop_fd` — only killing the mute server (EOF) ends it.
+* mntgen: mount indexes are ordinals and never reused (cap 4096 dials per
+ process); per-mount node ids cap at 2^32 lookups. `st_ino` mixes the
+ index with the qid.path via XOR — collisions remain theoretically
+ possible, just not the practical ones (identical qid.paths across two
+ servers). A dead mount keeps its slot — and its session socket and
+ interrupt-pipe descriptors — until process exit (early close would race
+ the dispatcher's INTERRUPT scan against fd reuse), so with a small
+ `ulimit -n` a re-dial storm exhausts descriptors before the index cap;
+ dials then fail cleanly with EIO.
+* mntgen: no posting/unposting through the mount (the synthetic root is
+ read-only; servers manage their registry entries through
+ `cloud9.post`), and no per-name mount options: one set of
+ `--uname/--aname/--msize/--cache` applies to every dial.
* No 9P2000.u/.L: no symlinks, ownership, or extended attributes.
* No PID namespace, no `/proc` remount. `--tcp` needs an IP literal.
* Cross-directory rename returns `EXDEV` (9P2000 cannot move files).
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)));
diff --git a/9ns/test/mntgen.sh b/9ns/test/mntgen.sh
new file mode 100755
index 0000000..ea58e79
--- /dev/null
+++ b/9ns/test/mntgen.sh
@@ -0,0 +1,222 @@
+#!/usr/bin/env bash
+# Integration tests for 9ns --mntgen: one FUSE mount whose synthetic root
+# lists the posted-9P registry ($XDG_RUNTIME_DIR/9p), servers dialed lazily
+# on the first walk into their name, one worker thread per server.
+# Usage: bash 9ns/test/mntgen.sh <9ns> <9proc-demo> (zig build 9ns-itest)
+# Exit 0 on success (or when the machine cannot run the tests), 1 on failure.
+set -u
+
+NS=$(realpath "${1:?path to 9ns}")
+PROC=$(realpath "${2:?path to 9proc-demo}")
+TMP=$(mktemp -d "${TMPDIR:-/tmp}/9ns-mntgen.XXXXXX")
+PIDS=()
+FAILED=0
+PASSED=0
+M=/mnt/9p
+RAMFS=/usr/lib/plan9/bin/ramfs
+cleanup() {
+ # ramfs delegates to 9pserve, and tests may spawn servers inside the
+ # namespace: catch anything holding a path under $TMP.
+ for p in "${PIDS[@]:-}"; do [ -n "$p" ] && kill "$p" 2>/dev/null; done
+ pkill -f "$TMP" 2>/dev/null
+ rm -rf "$TMP"
+}
+trap cleanup EXIT
+
+if ! unshare -Urm true 2>/dev/null; then
+ echo "SKIP: unprivileged user namespaces unavailable"; exit 0
+fi
+if [ ! -c /dev/fuse ]; then
+ echo "SKIP: /dev/fuse missing"; exit 0
+fi
+
+pass() { PASSED=$((PASSED + 1)); echo "ok - $1"; }
+fail() { FAILED=$((FAILED + 1)); echo "FAIL - $1"; shift; [ $# -gt 0 ] && printf ' %s\n' "$@"; }
+expect_eq() { # name expected actual
+ if [ "$2" = "$3" ]; then pass "$1"; else fail "$1" "expected: $(printf %q "$2")" "actual: $(printf %q "$3")"; fi
+}
+expect_contains() { # name needle haystack
+ case "$3" in *"$2"*) pass "$1" ;; *) fail "$1" "missing: $(printf %q "$2")" "in: $(printf %q "$3")" ;; esac
+}
+
+wait_socket() { # path
+ for _ in $(seq 1 100); do [ -S "$1" ] && return 0; sleep 0.05; done
+ return 1
+}
+
+# A scratch registry: the /srv translation is per-user tmpfs keyed by
+# XDG_RUNTIME_DIR; tests must never touch the real /run/user/<uid>/9p.
+# mktemp paths are short enough for the 108-byte socket name budget.
+export XDG_RUNTIME_DIR="$TMP"
+mkdir -p "$TMP/9p"
+REG="$TMP/9p"
+
+# run_in "<shell script>" — inside a namespace with the mntgen mount on $M.
+run_in() { timeout 60 "$NS" --mntgen --mount "$M" -- sh -c "$1" 2>"$TMP/stderr"; }
+
+# ============================================================================
+echo "# two servers posted under two names"
+# 9proc-demo posts its socket directly in the registry directory: a socket
+# at $XDG_RUNTIME_DIR/9p/<name> IS a posted name (the /srv model).
+"$PROC" --unix "$REG/alpha" &
+ALPHA=$!
+PIDS+=($ALPHA)
+wait_socket "$REG/alpha" || { echo "9proc-demo did not post $REG/alpha"; exit 1; }
+# plan9port ramfs: an independent 9P2000 implementation; NAMESPACE points its
+# post9pservice at the registry (-S names the service; the default would
+# post under "ramfs").
+HAVE_RAMFS=no
+if [ -x "$RAMFS" ]; then
+ NAMESPACE="$REG" "$RAMFS" -S beta &
+ PIDS+=($!)
+ wait_socket "$REG/beta" && HAVE_RAMFS=yes
+fi
+
+echo "# synthetic root: posted names, no connection made for listing"
+OUT=$(run_in "ls $M")
+expect_contains "root lists alpha" "alpha" "$OUT"
+[ "$HAVE_RAMFS" = yes ] && expect_contains "root lists beta" "beta" "$OUT"
+# A plain file in the registry is listed but can never be dialed: this also
+# proves listing connects to nothing (there is nothing to connect to).
+touch "$REG/junk"
+expect_contains "root lists a non-socket entry" "junk" "$(run_in "ls $M")"
+
+echo "# lazy dial and per-server routing in one program run"
+OUT=$(run_in "ls $M; ls $M/alpha; cat $M/alpha/build/zig_version")
+expect_eq "alpha subtree served after lazy dial" "$(zig version)" "$(run_in "cat $M/alpha/build/zig_version")"
+expect_contains "root and subtree in one ls run" "build" "$OUT"
+expect_contains "root and subtree in one ls run (zig version)" "$(zig version)" "$OUT"
+if [ "$HAVE_RAMFS" = yes ]; then
+ expect_eq "ramfs file round trip through its own mount" "hello-from-beta" \
+ "$(run_in "echo hello-from-beta > $M/beta/f && cat $M/beta/f")"
+ # Both servers in one run: the dispatcher serves them through two workers.
+ expect_eq "two subtrees in one run" "alpha ok beta ok" \
+ "$(run_in "[ -f $M/alpha/build/zig_version ] && echo alpha ok; [ -f $M/beta/f ] && echo beta ok" | tr '\n' ' ' | sed 's/ $//')"
+ # Parallel readers on both mounts at once.
+ expect_eq "parallel reads on both mounts" "ok" \
+ "$(run_in 'for i in 1 2 3 4 5 6 7 8; do cat '"$M"'/alpha/build/zig_version >/dev/null & cat '"$M"'/beta/f >/dev/null & done; wait; echo ok')"
+fi
+
+echo "# a post after the mount is visible (snapshot per opendir)"
+"$PROC" --unix "$REG/gamma" &
+PIDS+=($!)
+wait_socket "$REG/gamma"
+expect_contains "gamma appears without remounting" "gamma" "$(run_in "ls $M")"
+expect_eq "gamma dials on walk" "$(zig version)" "$(run_in "cat $M/gamma/build/zig_version")"
+
+echo "# walking a non-socket entry fails; the entry is never removed"
+expect_eq "walk into junk yields EIO, exit status" "1" "$(run_in "cat $M/junk 2>/dev/null; echo \$?")"
+expect_contains "junk still listed" "junk" "$(run_in "ls $M")"
+[ -S "$REG/junk" ] && fail "junk was turned into a socket" || pass "junk untouched"
+
+echo "# server death: EIO on the subtree, name still listed, re-dial on re-post"
+# SIGKILL, not SIGTERM: a graceful 9proc-demo unposts (the /srv model — a
+# clean exit removes the name), while the acceptance case is a server that
+# dies without cleaning up: the stale socket stays and must be listed,
+# yield EIO on walks, and be replaced by the next server that posts.
+# (the kill itself happens inside the child, at a deterministic point)
+# One 9ns process stays mounted through the whole cycle. The child shares
+# the host PID namespace, so the program itself kills the server and
+# re-posts a fresh one under the same name at deterministic points.
+DEATH_OUT=$(run_in '
+ cat '"$M"'/alpha/build/zig_version >/dev/null && echo dial=ok
+ kill -9 '"$ALPHA"' 2>/dev/null
+ sleep 0.4
+ MSG=$(cat '"$M"'/alpha/build/zig_version 2>&1 >/dev/null); echo "dead_status=$? dead_msg=$MSG"
+ ls '"$M"' | grep -q "^alpha$" && echo listed=yes
+ '"$PROC"' --unix '"$REG"'/alpha >/dev/null 2>&1 &
+ NEW=$!
+ for i in $(seq 1 50); do cat '"$M"'/alpha/build/zig_version >/dev/null 2>&1 && break; sleep 0.1; done
+ cat '"$M"'/alpha/build/zig_version >/dev/null && echo redial=ok
+ kill -9 "$NEW" 2>/dev/null; wait "$NEW" 2>/dev/null
+')
+expect_contains "mount dialed the server first" "dial=ok" "$DEATH_OUT"
+expect_contains "dead server yields EIO, not a hang" "dead_status=1" "$DEATH_OUT"
+expect_contains "dead server errno is EIO" "Input/output error" "$DEATH_OUT"
+expect_contains "dead name still listed (stale socket, listing dials nothing)" "listed=yes" "$DEATH_OUT"
+expect_contains "re-dial on next walk after re-post" "redial=ok" "$DEATH_OUT"
+# The re-posted server was killed without cleanup inside the child, so a
+# fresh process walks a stale entry: EIO, again.
+expect_eq "stale name answers EIO in a fresh process" "1" "$(run_in "cat $M/alpha/build/zig_version 2>/dev/null; echo \$?")"
+# Leave a live alpha for the remaining sections.
+"$PROC" --unix "$REG/alpha" &
+ALPHA=$!
+PIDS+=($ALPHA)
+wait_socket "$REG/alpha"
+expect_eq "re-dial picks up the fresh post" "$(zig version)" "$(run_in "cat $M/alpha/build/zig_version")"
+
+echo "# mount defaults and environment"
+expect_eq "default mountpoint is /mnt/9p" "$M" "$(timeout 60 "$NS" --mntgen -- sh -c 'echo $NINE_MOUNT')"
+expect_contains "default mount is served" "alpha" "$(timeout 60 "$NS" --mntgen -- sh -c 'ls $NINE_MOUNT')"
+expect_eq "mount is fuse" "yes" "$(run_in "grep -q \"^9ns $M fuse\" /proc/mounts && echo yes")"
+
+echo "# usage errors"
+expect_eq "--mntgen with --unix is a usage error" "125" "$(timeout 10 "$NS" --mntgen --unix /tmp/x -- true 2>/dev/null; echo $?)"
+expect_eq "--mntgen with --tcp is a usage error" "125" "$(timeout 10 "$NS" --mntgen --tcp 127.0.0.1:564 -- true 2>/dev/null; echo $?)"
+expect_eq "--mntgen with --fd is a usage error" "125" "$(timeout 10 "$NS" --mntgen --fd 3 -- true 2>/dev/null; echo $?)"
+expect_eq "--mntgen with --spawn is a usage error" "125" "$(timeout 10 "$NS" --spawn "true" --mntgen -- true 2>/dev/null; echo $?)"
+expect_eq "--mntgen with --name is a usage error" "125" "$(timeout 10 "$NS" --mntgen --name foo -- true 2>/dev/null; echo $?)"
+expect_eq "--mntgen=x is a usage error" "125" "$(timeout 10 "$NS" --mntgen=x -- true 2>/dev/null; echo $?)"
+expect_eq "missing XDG_RUNTIME_DIR is fatal" "125" "$(env -u XDG_RUNTIME_DIR timeout 10 "$NS" --mntgen -- true 2>/dev/null; echo $?)"
+expect_eq "exit status propagates" "7" "$(run_in 'exit 7'; echo $?)"
+
+echo "# xattr probes on the synthetic root read as unsupported, not EPERM"
+# llistxattr/lgetxattr (ACL and capability probes from ls -l and stat) used to
+# get EPERM from the root's read-only fallback, and coreutils blamed the
+# mount root itself: "ls: /mnt/9p: Operation not permitted". They must read
+# ENOSYS (like every other unimplemented op), which the kernel turns into a
+# silent "no xattrs here".
+XATTR_OUT=$(run_in "stat $M > /dev/null 2>&1; echo stat_rc=\$?; ls -ld $M > /dev/null 2>&1; echo lsld_rc=\$?")
+expect_contains "stat of the mount root is clean" "stat_rc=0" "$XATTR_OUT"
+expect_contains "ls -ld of the mount root is clean" "lsld_rc=0" "$XATTR_OUT"
+# junk (a plain registry file) makes full ls -l fail with EIO on that entry
+# (expected); what must never appear is EPERM blamed on the root.
+BAD=$(run_in "ls -l $M 2>&1 | grep -i 'not permitted' | head -1")
+expect_eq "no EPERM leak from ls -l" "" "$BAD"
+
+echo "# --debug and --no-direct-io reach the mntgen dispatcher"
+# runMntgen used to drop both flags from the options it handed serveMntgen;
+# --debug produced no trace at all.
+DEBUG_ERR=$(timeout 60 "$NS" --mntgen --debug -- sh -c "ls $M/alpha >/dev/null" 2>&1 >/dev/null)
+expect_contains "--debug traces the dispatcher" "lookup 'alpha'" "$DEBUG_ERR"
+
+echo "# a dial parked on a mute server unwedges when the program exits"
+# A server that accepts the connection but never answers Tversion parks the
+# dispatcher in the dial (pinned: no dial timeout, no concurrent dial). The
+# dial must watch stop_fd through the whole handshake: when the program's
+# main flow exits while a background walk is parked there, 9ns must follow it
+# out instead of wedging forever (pre-fix it survived SIGTERM).
+if command -v python3 > /dev/null; then
+ python3 - "$REG/mute" <<'PYEOF' &
+import socket, sys, os
+path = sys.argv[1]
+try:
+ os.unlink(path)
+except FileNotFoundError:
+ pass
+s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
+s.bind(path)
+s.listen(8)
+while True:
+ conn, _ = s.accept() # accept, then never say a word
+PYEOF
+ MUTEPID=$!
+ PIDS+=($MUTEPID)
+ sleep 0.3
+ HANG_START=$(date +%s)
+ timeout 20 "$NS" --mntgen -- bash -c "(stat $M/mute >/dev/null 2>&1) & sleep 1" >/dev/null 2>&1
+ HANG_RC=$?
+ HANG_SECONDS=$(( $(date +%s) - HANG_START ))
+ if [ "$HANG_RC" -eq 124 ] || [ "$HANG_SECONDS" -ge 15 ]; then
+ fail "9ns unwedges after the program exits a parked dial" "rc=$HANG_RC after ${HANG_SECONDS}s (wedged)"
+ else
+ pass "9ns unwedges after the program exits a parked dial (rc=$HANG_RC after ${HANG_SECONDS}s)"
+ fi
+ rm -f "$REG/mute"
+else
+ echo "SKIP: mute-server dial-hang check needs python3"
+fi
+
+echo
+echo "passed=$PASSED failed=$FAILED"
+[ "$FAILED" -eq 0 ]