diff options
| -rw-r--r-- | 9ns/build.zig | 6 | ||||
| -rw-r--r-- | 9ns/docs/DESIGN.md | 187 | ||||
| -rw-r--r-- | 9ns/src/bridge.zig | 750 | ||||
| -rw-r--r-- | 9ns/src/main.zig | 116 | ||||
| -rw-r--r-- | 9ns/src/nine.zig | 55 | ||||
| -rwxr-xr-x | 9ns/test/mntgen.sh | 222 | ||||
| -rw-r--r-- | 9proc/src/linux/probe.zig | 50 | ||||
| -rw-r--r-- | docs/design.md | 86 | ||||
| -rw-r--r-- | src/post.zig | 1024 | ||||
| -rw-r--r-- | src/root.zig | 4 | ||||
| -rw-r--r-- | src/serve.zig | 175 |
11 files changed, 2627 insertions, 48 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, ®_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 ] diff --git a/9proc/src/linux/probe.zig b/9proc/src/linux/probe.zig index 4db277d..c2f1026 100644 --- a/9proc/src/linux/probe.zig +++ b/9proc/src/linux/probe.zig @@ -67,6 +67,10 @@ pub const Error = error{ TooManyProviders, PathTooLong, BadAddress, + /// The unix path holds a live server; nothing is deleted or taken. + AlreadyListening, + /// The unix path holds a foreign non-socket entry; never deleted. + Occupied, /// A syscall failed; `last_errno` says which error. Syscall, }; @@ -128,6 +132,8 @@ pub fn Probe(comptime Srv: type) type { wake_fd: i32 = -1, unix_path: [108]u8 = undefined, unix_len: usize = 0, + /// The bound entry's inode; `stop` unlinks only its own socket. + unix_ino: u64 = 0, thread: ?std.Thread = null, thread_tid: std.atomic.Value(u32) = .init(0), nclients: std.atomic.Value(u32) = .init(0), @@ -233,7 +239,11 @@ pub fn Probe(comptime Srv: type) type { p.listen_fd = -1; } if (p.unix_len > 0) { - _ = linux.unlink(@ptrCast(&p.unix_path)); + // Only our own entry: a late stop must never unlink a + // name another server has since claimed. + if (unixIno(@ptrCast(&p.unix_path))) |ino| { + if (p.unix_ino != 0 and ino == p.unix_ino) _ = linux.unlink(@ptrCast(&p.unix_path)); + } p.unix_len = 0; } if (p.wake_fd >= 0) { @@ -522,14 +532,48 @@ pub fn Probe(comptime Srv: type) type { try p.check(rc); const lfd: i32 = @intCast(rc); errdefer _ = linux.close(lfd); - // No libc, so no "is it still listening" probe: unlink a stale socket and bind. - _ = linux.unlink(@ptrCast(&sa.path)); + // Nothing foreign is deleted: a non-socket entry at the path + // is refused (Occupied), a live server is refused + // (AlreadyListening), and only a socket that refuses a + // connect — a corpse — is unlinked. (A hand racing the swap + // between probe and unlink is the documented residual of this + // cheap protocol; cloud9.post claims names atomically when + // that matters.) + const st = unixStat(@ptrCast(&sa.path)) catch return error.Occupied; + if (st) |s| { + if (s.mode & linux.S.IFMT != linux.S.IFSOCK) return error.Occupied; + if (cloud9.post.probe(@ptrCast(&sa.path)) != .stale) return error.AlreadyListening; + _ = linux.unlink(@ptrCast(&sa.path)); + } try p.check(linux.bind(lfd, @ptrCast(&sa), @sizeOf(linux.sockaddr.un))); try p.check(linux.listen(lfd, 128)); p.listen_fd = lfd; p.own_listener = true; p.unix_path = sa.path; p.unix_len = path.len; + // Our entry's inode, so `stop` never unlinks a name another + // server has since claimed. + p.unix_ino = if (unixStat(@ptrCast(&p.unix_path)) catch null) |s| s.ino else 0; + } + + /// statx(2) of one path: null when it does not exist, and an + /// error when the kernel cannot say — the caller refuses rather + /// than guesses. + fn unixStat(path: [*:0]const u8) !?linux.Statx { + var stx: linux.Statx = undefined; + const rc = linux.statx(linux.AT.FDCWD, path, 0, .{ .TYPE = true, .INO = true }, &stx); + const s: isize = @bitCast(rc); + if (s == 0) return stx; + const noent: isize = @intCast(@intFromEnum(linux.E.NOENT)); + if (s == -noent) return null; + return error.StatFailed; + } + + /// The socket at `path`, by inode; null when it is gone or + /// unreadable. + fn unixIno(path: [*:0]const u8) ?u64 { + const stx = unixStat(path) catch return null; + return if (stx) |s| s.ino else null; } fn listenTcp(p: *Self, text: []const u8) Error!void { diff --git a/docs/design.md b/docs/design.md index f34a3ea..e68530f 100644 --- a/docs/design.md +++ b/docs/design.md @@ -170,8 +170,90 @@ still yields through `serve`, closes the socket and frees the slot. `stop()` cancels the group (`Io.Group.cancel` interrupts the blocking accepts and reads), joins every task and closes the listeners; a `serve` that waits on another thread must stop waiting once `stop()` has begun. -Unix socket paths are the application's: the runner neither unlinks, -chmods nor removes them. +Unix paths bound by `listen()` are the application's: the runner neither +unlinks, chmods nor removes them. `listenPosted(env, name, backlog)` is the +exception with a rule of its own: it posts through the registry (below) — +one posted name per runner, a second `listenPosted` is `error.AlreadyPosted` +(its socket would be orphaned in the registry, nothing left to unpost it) — +and `stop()` unposts, but only while the entry is still the runner's own +(`posted_ino`, below): a name re-posted by another server after this +runner's socket file was lost survives the stop. + +# Post registry + +`post` is the `/srv` translation: a server posts itself under a name in one +per-user registry directory and clients list the names and dial them, the +shape of Plan 9's `devsrv.c` and plan9port's `post9pservice`. The registry +is `$XDG_RUNTIME_DIR/9p/` (per-user tmpfs, the right lifetime); a name's +socket lives at `$XDG_RUNTIME_DIR/9p/<name>`. There is no union daemon and +no per-server directory: mounting and namespace policy are the client's +(9ns `--mntgen` lists the registry for its synthetic root and dials on +walk). + +The library never reads the process environment: the environment comes in +as the raw block `main` receives (`post.Env`, scanned by `post.getenv`, +the same threading 9ns uses). `registryDir`/`registryPath` build paths +into caller buffers, zero-terminated, and refuse when `XDG_RUNTIME_DIR` is +unset — there is no `/tmp` fallback. A name is legal when it would be a +legal file name for the engine (`legalName` in fs.zig: non-empty, not "." +or "..", no '/' or NUL) *and* fits the Unix socket path budget +(`transport.sun_path_len`); path traversal through a posted name is the +attack the caps exist for, and `max_name_len` derives from that budget +against a conforming 32-byte `/run/user/<uid>` prefix, with the real +prefix checked again per call. + +`post.post(io, env, name, backlog, path_buf)` posts: the registry directory +is created 0o750 (already present is fine) and the claim runs. The claiming +socket is bound and **listening at a private temp path first** — +`$XDG_RUNTIME_DIR/.post.sock.<pid>.<serial>`, beside the registry, never +inside it, so no listing ever sees it — and only then is the name taken, +entirely by atomic `renameat2` calls under an advisory `flock` on +`$XDG_RUNTIME_DIR/.post.lock` (the kernel drops the lock if the poster +dies; the dotfile persists, empty): a free name is claimed with +`RENAME_NOREPLACE`, which is the arbiter between two posts racing on one +name — exactly one ends up bound, the loser re-runs into +`error.AlreadyPosted`; a stale entry (a socket that refuses a connect) is +grabbed with `RENAME_EXCHANGE` against a private dummy file +(`.post.tmp.*`, also beside the registry) and verified by inode *and* a +fresh probe before the dead socket is removed, so a live server that came +up in between is swapped back untouched. **`post` never unlinks the +registry path**: what a claim displaces is parked under private temp names +and deleted only when provably its own dummy or the verified dead entry — +anything a foreign hand put in their place is left alone, intact. An +entry that is not a socket is `error.NotSocket` and is never touched at +all, learned the hard way in zmx. The lone exception is the legacy +fallback on filesystems without the `renameat2` flags: there the stale +entry is removed by an inode-checked in-place unlink (the replace window +that costs is confined to such filesystems). + +The claimed entry's inode comes back in `Posted{server, path, inode}`, and +`post.unpost(io, path, inode)` unlinks only that entry: a path whose +content has been replaced (the socket file lost, the name re-posted by +another server) is left strictly alone, so one server's late stop can +never unpost another's live name. Idempotent, errors swallowed. On the +server side, `serve.Runner.listenPosted` posts and `stop()` unposts. + +Probing is a nonblocking raw-syscall connect (`post.probe`: +`.none`/`.stale`/`.live`; uncertainty counts live), because `std.Io`'s +Unix connect does not promise ECONNREFUSED. The kernel answers a connect +to a non-socket entry with the same ECONNREFUSED as to a dead server's +socket, so `.stale` covers both — `post` tells them apart by stat before +anything is removed; callers that must know use `Io.Dir.statFile`. + +`post.posted(io, env, out)` lists the raw entries into a caller buffer as +`len:u8 name` staged records (the engine's readdir staging pattern; no +allocation, no connection made — a missing registry lists as empty) and +returns a `Names` iterator; staleness is the caller's concern, and a +buffer too small for the listing is `error.NoSpace`, never a silent drop. +`post.dial(io, env, name)` connects and answers `error.NotPosted` (no +entry) and `error.Stale` (entry present, connection refused) distinctly. +`post.Watch` is a Linux inotify watcher on the registry directory for +hosts that cache its listing (`init`, `add`, `next`, `deinit`; hosted +Linux only — everything else in the module is platform-neutral, and +`next` yields `added` and `removed` per name and, with an empty name, the +two signals a caching consumer must not miss: `overflow` — the kernel +dropped events, rescan — and `gone` — the watch itself ended (the +directory was removed); re-`add` and rescan. # Multiplexer (9web) diff --git a/src/post.zig b/src/post.zig new file mode 100644 index 0000000..47a7d3d --- /dev/null +++ b/src/post.zig @@ -0,0 +1,1024 @@ +//! The `/srv` translation: a server posts itself under a name in one +//! per-user registry directory, and clients list the names and dial them +//! (`post9pservice` in plan9port, `devsrv.c` in Plan 9). The registry is +//! `$XDG_RUNTIME_DIR/9p/` — per-user tmpfs, the right lifetime — and a +//! name's socket lives at `$XDG_RUNTIME_DIR/9p/<name>`. Mounting and +//! namespace policy stay with the client (9ns `--mntgen`); this module +//! only builds paths, posts, lists, dials and watches. +//! +//! House style as everywhere in cloud9: no assumed allocator, buffers +//! are the caller's, the listing is staged records, and the environment +//! is passed in as the raw block `main` receives (`post.Env`) — the +//! library never reads the process environment on its own. The Linux +//! syscalls (probe, dial, `Watch`) are raw and need no libc; the rest +//! goes through `std.Io`. See docs/design.md, "Post registry". +const std = @import("std"); +const builtin = @import("builtin"); +const Io = std.Io; +const linux = std.os.linux; +const transport = @import("transport.zig"); + +/// The raw environment block the C startup hands to `main`, the same +/// shape 9ns threads to `ns.getenv`. `null` terminates. +pub const Env = [*:null]const ?[*:0]const u8; + +/// The kernel's Unix socket path budget (`sun_path`). +pub const sun_path_len = transport.sun_path_len; + +/// The longest name a post may carry: the socket path +/// `$XDG_RUNTIME_DIR/9p/<name>` must stay inside `sun_path_len`, and a +/// conforming `$XDG_RUNTIME_DIR` (`/run/user/<uid>`) is budgeted at 32 +/// bytes, leaving the rest of the budget for `/9p/`, the name and the +/// terminating zero. `registryPath` checks the real prefix again. +pub const max_name_len = sun_path_len - "/9p/".len - 1 - 32; + +comptime { + if (max_name_len < 8) @compileError("9P post name budget is too small to be useful"); +} + +/// One entry of a raw environment block. +pub fn getenv(env: Env, name: []const u8) ?[]const u8 { + var i: usize = 0; + while (env[i]) |entry| : (i += 1) { + const kv = std.mem.span(entry); + const eq = std.mem.indexOfScalar(u8, kv, '=') orelse continue; + if (std.mem.eql(u8, kv[0..eq], name)) return kv[eq + 1 ..]; + } + return null; +} +/// Is `name` legal for posting? It mirrors the engine's `legalName` +// (empty, "." and ".." refused, no '/' or NUL) with the post-specific +// budget: path traversal through a posted name is THE attack, and the +// name must also fit a Unix socket path. +pub fn legalName(name: []const u8) bool { + if (name.len == 0 or name.len > max_name_len) return false; + if (std.mem.eql(u8, name, ".") or std.mem.eql(u8, name, "..")) return false; + return std.mem.indexOfAny(u8, name, "/\x00") == null; +} + +pub const PathError = error{ + /// `$XDG_RUNTIME_DIR` is unset. There is no `/tmp` fallback. + NotRuntimeDir, + /// The name cannot be posted (see `legalName`). + IllegalName, + /// The caller's buffer is too small for the path and its zero. + NoSpace, + /// The full socket path would not fit `sun_path_len`. + NameTooLong, +}; + +const xdg_runtime_dir = "XDG_RUNTIME_DIR"; + +/// Writes `$XDG_RUNTIME_DIR/9p` into `out` and returns it zero-terminated. +pub fn registryDir(env: Env, out: []u8) PathError![:0]u8 { + const root = getenv(env, xdg_runtime_dir) orelse return error.NotRuntimeDir; + if (root.len == 0) return error.NotRuntimeDir; + const total = root.len + "/9p".len; + if (total + 1 > out.len) return error.NoSpace; + @memcpy(out[0..root.len], root); + @memcpy(out[root.len..total], "/9p"); + out[total] = 0; + return out[0..total :0]; +} + +/// Writes `$XDG_RUNTIME_DIR/9p/<name>` into `out` and returns it +/// zero-terminated, ready for `probe`, `dial` or `Io.net.UnixAddress`. +pub fn registryPath(env: Env, name: []const u8, out: []u8) PathError![:0]u8 { + if (!legalName(name)) return error.IllegalName; + const root = getenv(env, xdg_runtime_dir) orelse return error.NotRuntimeDir; + if (root.len == 0) return error.NotRuntimeDir; + const total = root.len + "/9p/".len + name.len; + if (total >= sun_path_len) return error.NameTooLong; + if (total + 1 > out.len) return error.NoSpace; + @memcpy(out[0..root.len], root); + @memcpy(out[root.len..][0.."/9p/".len], "/9p/"); + @memcpy(out[total - name.len ..][0..name.len], name); + out[total] = 0; + return out[0..total :0]; +} + +/// The raw registry entries, staged by `posted` as `len:u8 name` +/// records back to back in the caller's buffer (the engine's staged- +/// record pattern for readdir). Nothing is dialed; staleness is the +/// caller's concern (`probe`). +pub const Names = struct { + bytes: []const u8, + i: usize = 0, + + pub fn next(n: *Names) ?[]const u8 { + if (n.i + 1 > n.bytes.len) return null; + const len = n.bytes[n.i]; + if (n.i + 1 + len > n.bytes.len) return null; + const name = n.bytes[n.i + 1 ..][0..len]; + n.i += 1 + @as(usize, len); + return name; + } +}; + +pub const PostedError = PathError || error{NoSpace} || Io.Dir.OpenError || Io.Dir.Reader.Error; + +/// Lists the registry into `out` and returns an iterator over the names. +/// Nothing is dialed; a missing registry lists as empty (no server has +/// posted yet). The iteration buffer is this stack frame's, so only the +/// staged names outlive the call. +pub fn posted(io: Io, env: Env, out: []u8) PostedError!Names { + var dir_buf: [std.fs.max_path_bytes]u8 = undefined; + const dir_path = try registryDir(env, &dir_buf); + const dir = Io.Dir.openDirAbsolute(io, dir_path, .{ .iterate = true }) catch |err| switch (err) { + error.FileNotFound, error.NotDir => return .{ .bytes = out[0..0] }, + else => return err, + }; + defer Io.Dir.close(dir, io); + var read_buf: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined; + var reader = Io.Dir.Reader.init(dir, &read_buf); + var n: usize = 0; + while (try reader.next(io)) |entry| { + if (n + 1 + entry.name.len > out.len) return error.NoSpace; + out[n] = @intCast(entry.name.len); + @memcpy(out[n + 1 ..][0..entry.name.len], entry.name); + n += 1 + entry.name.len; + } + return .{ .bytes = out[0..n] }; +} + +/// What `probe` found at a registry path. Nothing is modified: a stale +/// entry is reported, not removed (`post` owns replacement). +pub const Probe = enum { none, stale, live }; + +/// Probes a registry entry without touching it: connect to the socket +/// path without blocking. Refused means the path holds nothing to talk +/// to — a dead server's socket *or an entry that is not a socket at +/// all*, which the kernel answers with the same ECONNREFUSED; `post` +/// tells them apart by stat and never deletes the latter. No entry at +/// all is `none`; connected — or busy, or any unexpected problem — is +/// `live`, because uncertainty must be owned by the server, never +/// resolved by deleting what may be someone's socket. +pub fn probe(path: [:0]const u8) Probe { + const fd = socketNonblocking() catch return .live; + defer _ = linux.close(fd); + var addr: linux.sockaddr.un = .{ .path = @splat(0) }; + @memcpy(addr.path[0..path.len], path); + const rc: isize = @bitCast(linux.connect(fd, @ptrCast(&addr), @sizeOf(linux.sockaddr.un))); + return switch (rawErrno(rc)) { + .SUCCESS, .AGAIN, .INPROGRESS, .PERM, .ACCES => .live, + .NOENT, .NOTDIR => .none, + .CONNREFUSED => .stale, + else => .live, + }; +} + +fn socketNonblocking() !i32 { + const rc = linux.socket(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.NONBLOCK | linux.SOCK.CLOEXEC, 0); + if (rawErrno(@bitCast(rc)) != .SUCCESS) return error.Socket; + return @intCast(rc); +} + +/// The errno of a raw Linux syscall return, which encodes it as the +/// negated value (0 on success). The std.Io Unix connect does not +/// promise ECONNREFUSED, so `probe` and `dial` speak to the kernel +/// directly. +fn rawErrno(rc: isize) linux.E { + const n: usize = if (rc < 0 and -rc < 4096) @intCast(-rc) else 0; + return @enumFromInt(n); +} + +pub const DialError = error{ + NotRuntimeDir, + IllegalName, + NoSpace, + NameTooLong, + /// No registry entry under this name. + NotPosted, + /// The entry exists but the connection is refused: it is stale. + Stale, + SystemResources, + ProcessFdQuotaExceeded, + SystemFdQuotaExceeded, + AccessDenied, + PermissionDenied, + SymLinkLoop, + NotDir, + Io, + Unexpected, +}; + +/// Dials the server posted under `name` and returns its stream. The +/// descriptor is blocking and close-on-exec; it is not the registry's — +/// closing it changes nothing in the registry. On Linux the connect is +/// a raw syscall (so ECONNREFUSED is distinguishable, which `std.Io`'s +/// Unix connect does not promise). +pub fn dial(io: Io, env: Env, name: []const u8) DialError!Io.net.Stream { + var path_buf: [sun_path_len]u8 = undefined; + const path = try registryPath(env, name, &path_buf); + const fd = connectBlocking(path) catch |err| switch (err) { + error.Socket => return error.SystemResources, + error.Noent => return error.NotPosted, + error.Refused => return error.Stale, + error.AccessDenied => return error.AccessDenied, + error.Loop => return error.SymLinkLoop, + error.NotDir => return error.NotDir, + }; + _ = io; + return .{ .socket = .{ .handle = fd, .address = .{ .ip4 = .loopback(0) } } }; +} + +const ConnectError = error{ Socket, Noent, Refused, AccessDenied, Loop, NotDir }; + +fn connectBlocking(path: [:0]const u8) ConnectError!i32 { + const rc = linux.socket(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0); + if (rawErrno(@bitCast(rc)) != .SUCCESS) return error.Socket; + const fd: i32 = @intCast(rc); + errdefer _ = linux.close(fd); + var addr: linux.sockaddr.un = .{ .path = @splat(0) }; + @memcpy(addr.path[0..path.len], path); + const crc: isize = @bitCast(linux.connect(fd, @ptrCast(&addr), @sizeOf(linux.sockaddr.un))); + switch (rawErrno(crc)) { + .SUCCESS => return fd, + .NOENT, .NOTDIR => return error.Noent, + .CONNREFUSED => return error.Refused, + .ACCES, .PERM => return error.AccessDenied, + .LOOP => return error.Loop, + else => return error.Socket, + } +} + +/// What `post` bound, and what the caller owes `unpost`: the registry +/// path and the inode of the bound entry, so a late `unpost` can prove +/// the entry is still its own before unlinking it. +pub const Posted = struct { + server: Io.net.Server, + path: [:0]const u8, + inode: u64, +}; + +pub const PostError = error{ + NotRuntimeDir, + IllegalName, + NoSpace, + NameTooLong, + /// A live server owns the name. + AlreadyPosted, + /// The registry entry exists but is not a socket. It is never + /// deleted; whoever put it there must remove it themselves. + NotSocket, +} || Io.net.UnixAddress.ListenError || Io.Dir.CreateDirPathError || Io.Dir.StatFileError || + Io.Dir.DeleteFileError || Io.UnexpectedError || Io.Cancelable; + +/// How many rounds of the claim protocol `post` runs before giving the +/// name up as contested. +const post_attempts = 8; + +/// Distinct temp socket names for posts sharing one pid (threads). +var temp_serial: std.atomic.Value(u64) = .init(0); + +/// Posts a listening socket under `name`: the registry directory is +/// created (0o750; already present is fine), the name's path — written +/// into `path_buf` — is bound and returned for `unpost`. +/// +/// The protocol: a live server owning the name is `AlreadyPosted`; an +/// entry that is not a socket is `NotSocket` and is never deleted. The +/// claiming socket is bound and *listening* at a private temp path +/// first (so any probe of it answers live — no window in which the +/// claim could be mistaken for a stale entry), and the name is taken +/// under the stale-removal lock by atomic `renameat2` calls only +/// (`claimName`): of two posts racing on one stale name exactly one +/// ends up owning the entry, the loser re-runs into `AlreadyPosted`, +/// and the registry path is never unlinked by a claim — what a claim +/// displaces is parked under a private temp name and left alone +/// unless provably ours or provably the dead entry. +pub fn post(io: Io, env: Env, name: []const u8, backlog: u31, path_buf: []u8) PostError!Posted { + const path = try registryPath(env, name, path_buf); + var dir_buf: [std.fs.max_path_bytes]u8 = undefined; + const dir_path = try registryDir(env, &dir_buf); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, dir_path, .fromMode(0o750)); + const parent = dir_path[0 .. dir_path.len - "/9p".len]; + const serial = temp_serial.fetchAdd(1, .monotonic); + var tmp_buf: [std.fs.max_path_bytes]u8 = undefined; + const tmp = std.fmt.bufPrintZ(&tmp_buf, "{s}/.post.sock.{d}.{d}", .{ parent, linux.getpid(), serial }) catch return error.NoSpace; + var attempts: usize = 0; + var saw_not_socket = false; + claim: while (attempts < post_attempts) : (attempts += 1) { + // A live listener at a private temp path, invisible to registry + // listings, before the name is even looked at. + const addr = try Io.net.UnixAddress.init(tmp); + var server = addr.listen(io, .{ .kernel_backlog = backlog }) catch |err| switch (err) { + error.AddressInUse => { + // A crashed former self under a reused pid, or stale + // garbage: the dead temp is ours to clear. + _ = linux.unlink(tmp.ptr); + continue :claim; + }, + else => { + _ = linux.unlink(tmp.ptr); + return err; + }, + }; + switch (claimName(io, dir_path, path, tmp)) { + .won => { + // Ours, and live: nothing in the protocol displaces a + // live entry. Record the inode `unpost` checks — and + // only from a socket: whatever a foreign hand may have + // put in the entry's place between the claim and here + // is not ours to name. + const st = Io.Dir.statFile(.cwd(), io, path, .{}) catch |err| switch (err) { + error.FileNotFound => { + // Unlinked behind our back: not ours now. (The + // temp is already spent; whatever a foreign hand + // parked there was deliberately left alone.) + server.deinit(io); + continue :claim; + }, + else => { + server.deinit(io); + return err; + }, + }; + if (st.kind != .unix_domain_socket) { + server.deinit(io); + continue :claim; + } + return .{ .server = server, .path = path, .inode = st.inode }; + }, + .contended => { + _ = linux.unlink(tmp.ptr); + server.deinit(io); + return error.AlreadyPosted; + }, + .not_socket => { + saw_not_socket = true; + _ = linux.unlink(tmp.ptr); + server.deinit(io); + // Another post's stale-swap parks a dummy here for a + // moment; give the dust a beat to settle before the + // entry is declared bogus. + io.sleep(.fromMilliseconds(1), .awake) catch {}; + continue :claim; + }, + .retry => { + _ = linux.unlink(tmp.ptr); + server.deinit(io); + continue :claim; + }, + } + } + return if (saw_not_socket) error.NotSocket else error.AlreadyPosted; +} + +/// The outcome of one claim round. +const Claim = enum { + /// The live listener now sits at the registry path. + won, + /// A live server owns the name. + contended, + /// The registry entry is not a socket (never deleted). + not_socket, + /// The landscape raced; look again. + retry, +}; + +/// Takes the registry path for the listener bound at `sock_tmp`, under +/// the stale-removal lock: +/// +/// * No entry: an atomic `renameat2(RENAME_NOREPLACE)` claims it — two +/// posts race to exactly one winner. +/// * A live entry: `contended`. +/// * A non-socket entry: `not_socket` — it stays strictly alone. +/// * A stale entry: it is grabbed with `renameat2(RENAME_EXCHANGE)` +/// against a private dummy and verified by inode *and* a fresh probe +/// (a live server that came up since the probe is swapped back, +/// untouched), and only then is the live listener exchanged into the +/// entry's place. The registry path is never unlinked here — the +/// dummy and the dead entry die under private temp names, provably +/// by inode, and anything a foreign hand parked in their place is +/// left alone. +fn claimName(io: Io, registry_dir: [:0]const u8, path: [:0]const u8, sock_tmp: [:0]const u8) Claim { + const lock_fd = lockStaleRemoval(registry_dir); + defer { + if (lock_fd >= 0) _ = linux.close(lock_fd); + } + const parent = registry_dir[0 .. registry_dir.len - "/9p".len]; + + const st = Io.Dir.statFile(.cwd(), io, path, .{}) catch |err| switch (err) { + error.FileNotFound => { + // A fresh name: the rename is the arbiter. + const rc: isize = @bitCast(linux.renameat2(linux.AT.FDCWD, sock_tmp.ptr, linux.AT.FDCWD, path.ptr, .{ .NOREPLACE = true })); + return switch (rawErrno(rc)) { + .SUCCESS => .won, + .EXIST => .retry, + .INVAL, .NOSYS, .PERM, .OPNOTSUPP => legacyClaim(path, sock_tmp), + else => .retry, + }; + }, + else => return .retry, + }; + if (st.kind != .unix_domain_socket) return .not_socket; + switch (probe(path)) { + .live => return .contended, + .none, .stale => {}, + } + + // A dead entry. Grab it for a look: the dummy never enters the + // registry directory, so no listing ever sees it. + const dummy_serial = temp_serial.fetchAdd(1, .monotonic); + var dummy_buf: [std.fs.max_path_bytes]u8 = undefined; + const dummy = std.fmt.bufPrintZ(&dummy_buf, "{s}/.post.tmp.{d}.{d}", .{ parent, linux.getpid(), dummy_serial }) catch return .retry; + var dummy_file = Io.Dir.createFileAbsolute(io, dummy, .{ .exclusive = true, .truncate = false }) catch return .retry; + dummy_file.close(io); + const dummy_st = Io.Dir.statFile(.cwd(), io, dummy, .{}) catch { + _ = linux.unlink(dummy.ptr); + return .retry; + }; + + const rc: isize = @bitCast(linux.renameat2(linux.AT.FDCWD, path.ptr, linux.AT.FDCWD, dummy.ptr, .{ .EXCHANGE = true })); + switch (rawErrno(rc)) { + .NOENT => { + // The entry unposted meanwhile; look again. + unlinkIfOurs(io, dummy, dummy_st.inode); + return .retry; + }, + .INVAL, .NOSYS, .PERM, .OPNOTSUPP => { + // No RENAME_EXCHANGE on this filesystem: the grab cannot + // be made safe here; the documented fallback is an + // inode-checked in-place unlink (the only place a claim + // may delete a raced non-socket, on filesystems without + // the primitive). + unlinkIfOurs(io, dummy, dummy_st.inode); + const now = Io.Dir.statFile(.cwd(), io, path, .{}) catch return .retry; + if (now.inode != st.inode) return .retry; + Io.Dir.deleteFileAbsolute(io, path) catch return .retry; + return switch (claimFresh(path, sock_tmp)) { + .won => .won, + else => .retry, + }; + }, + .SUCCESS => {}, + else => { + unlinkIfOurs(io, dummy, dummy_st.inode); + return .retry; + }, + } + // The dummy sits at the entry's place; the grabbed entry is at + // `dummy`. Verify it is still the dead one we probed. + const grabbed = Io.Dir.statFile(.cwd(), io, dummy, .{}) catch { + _ = swapBack(path, dummy); + unlinkIfOurs(io, dummy, dummy_st.inode); + return .retry; + }; + if (grabbed.inode != st.inode or probe(dummy) == .live) { + // A live server took the name between the probe and the grab: + // restore it, untouched. + const back = swapBack(path, dummy); + if (rawErrno(back) == .SUCCESS) unlinkIfOurs(io, dummy, dummy_st.inode); + return if (grabbed.inode != st.inode) .retry else .contended; + } + // Verified dead: exchange the live listener into the entry's + // place. What the exchange parks at `sock_tmp` is our dummy — + // unless a foreign hand replaced it in the meantime, in which + // case it is left alone, intact. + const claim_rc: isize = @bitCast(linux.renameat2(linux.AT.FDCWD, sock_tmp.ptr, linux.AT.FDCWD, path.ptr, .{ .EXCHANGE = true })); + if (rawErrno(claim_rc) != .SUCCESS) { + _ = swapBack(path, dummy); + unlinkIfOurs(io, dummy, dummy_st.inode); + return .retry; + } + unlinkIfOurs(io, sock_tmp, dummy_st.inode); + // The dead entry dies under the private temp name, provably by + // inode; anything a foreign hand parked there survives. + unlinkIfOurs(io, dummy, grabbed.inode); + return .won; +} + +/// The NOREPLACE rename for a provably free path. +fn claimFresh(path: [:0]const u8, sock_tmp: [:0]const u8) Claim { + const rc: isize = @bitCast(linux.renameat2(linux.AT.FDCWD, sock_tmp.ptr, linux.AT.FDCWD, path.ptr, .{ .NOREPLACE = true })); + return switch (rawErrno(rc)) { + .SUCCESS => .won, + else => .retry, + }; +} + +/// The claim on a filesystem without `renameat2` flags at all: a plain +/// rename over a provably empty path (the replace window is the price +/// of such a filesystem). +fn legacyClaim(path: [:0]const u8, sock_tmp: [:0]const u8) Claim { + if (probe(path) != .none) return .retry; + const rc: isize = @bitCast(linux.rename(sock_tmp.ptr, path.ptr)); + return if (rawErrno(rc) == .SUCCESS) .won else .retry; +} + +/// Unlinks `p` only while it still holds inode `ino`: a file some +/// other hand has put in a private temp's place is left alone (a +/// leaked dotfile beats deleting what may be someone's file). +fn unlinkIfOurs(io: Io, p: [:0]const u8, ino: u64) void { + const st = Io.Dir.statFile(.cwd(), io, p, .{}) catch return; + if (st.inode != ino) return; + Io.Dir.deleteFileAbsolute(io, p) catch {}; +} + +/// Puts a swapped-out entry back: `path` and `tmp` re-exchange, so the +/// displaced entry returns to the registry and our dummy comes back to +/// the temp name. The result is ignored: the removal is serialized by +/// `lockStaleRemoval`, so the entry can only fail to return when the +/// path was cleared by hand in the meantime. +fn swapBack(path: [:0]const u8, tmp: [:0]const u8) isize { + return @bitCast(linux.renameat2(linux.AT.FDCWD, path.ptr, linux.AT.FDCWD, tmp.ptr, .{ .EXCHANGE = true })); +} + +/// `flock(LOCK_EX)` on the registry-side `.post.lock`, serializing the +/// stale-removal window between posts (including threads: the lock is +/// taken per open file description). The file is an empty dotfile; the +/// kernel drops the lock when the holder dies. A lock that cannot be +/// taken does not block posting — the window just loses its +/// serialization. +fn lockStaleRemoval(registry_dir: [:0]const u8) i32 { + var buf: [std.fs.max_path_bytes]u8 = undefined; + const lock_path = std.fmt.bufPrintZ(&buf, "{s}/.post.lock", .{registry_dir[0 .. registry_dir.len - "/9p".len]}) catch return -1; + const fd: isize = @bitCast(linux.open(lock_path.ptr, .{ .ACCMODE = .RDWR, .CREAT = true, .CLOEXEC = true }, 0o600)); + if (rawErrno(fd) != .SUCCESS) return -1; + const lfd: i32 = @intCast(fd); + while (rawErrno(@bitCast(linux.flock(lfd, lock_ex))) == .INTR) {} + return lfd; +} + +/// `LOCK_EX` for the raw `flock` call (the constant is not in std's +/// linux namespace). +const lock_ex: i32 = 2; + +/// Unposts: unlinks the registry path — but only the caller's own +/// entry. `inode` is the one `post` bound (in `Posted`); a path whose +/// entry has been replaced (the socket lost and the name re-posted by +/// another server) is left strictly alone, so one server's late stop +/// can never unpost another's live name. Idempotent; all errors are +/// swallowed — an unpost must never be the reason a server fails to +/// shut down. +pub fn unpost(io: Io, path: [:0]const u8, inode: u64) void { + const st = Io.Dir.statFile(.cwd(), io, path, .{}) catch return; + if (st.inode != inode) return; + Io.Dir.deleteFileAbsolute(io, path) catch {}; +} + +/// An inotify watcher on the registry directory, for hosts (9ns's +/// mntgen) that cache its listing: `init`, `add` the directory, `next` +/// yields names as they are added and removed. Hosted Linux only; the +/// rest of `post` is platform-neutral. +pub const Watch = if (builtin.os.tag == .linux) struct { + const Self = @This(); + + fd: i32, + /// Event staging, so `next` never truncates a record. + stage: [4096]u8 = undefined, + head: usize = 0, + tail: usize = 0, + + /// The kind of change. `overflow` (empty name) reports that the + /// kernel dropped events — the queue overflowed; a caching consumer + /// must rescan. `gone` (empty name) reports that the watch itself + /// ended (the directory was removed, or the kernel dropped the + /// watch): re-`add` it and rescan; until then `next` is null. + pub const Kind = enum { added, removed, overflow, gone }; + pub const Event = struct { name: []const u8, kind: Kind }; + + pub const Error = error{ SystemResources, ProcessFdQuotaExceeded, Io, Unexpected }; + + pub fn init() Error!Self { + const rc = linux.inotify_init1(linux.IN.CLOEXEC | linux.IN.NONBLOCK); + return switch (rawErrno(@bitCast(rc))) { + .SUCCESS => .{ .fd = @intCast(rc) }, + .NOMEM => error.SystemResources, + .MFILE, .NFILE => error.ProcessFdQuotaExceeded, + else => error.Io, + }; + } + + /// Watches the directory for names appearing and disappearing. + pub fn add(w: *Self, dir_path: [:0]const u8) Error!void { + const mask = linux.IN.CREATE | linux.IN.DELETE | linux.IN.MOVED_TO | linux.IN.MOVED_FROM; + const rc = linux.inotify_add_watch(w.fd, dir_path.ptr, mask); + switch (rawErrno(@bitCast(rc))) { + .SUCCESS => {}, + .NOMEM => return error.SystemResources, + else => return error.Io, + } + } + /// Yields the next registry change, or null when nothing is + /// pending. `Event.name` points into the watch's staging and is + /// valid until the following `next`. Events for anything but a name + /// appearing or disappearing (directories, metadata) are skipped — + /// except the two a caching consumer must not miss: a queue + /// overflow (`Kind.overflow`) and the watch ending + /// (`Kind.gone`). + pub fn next(w: *Self) Error!?Event { + while (true) { + if (w.head + @sizeOf(linux.inotify_event) <= w.tail) { + const ev: *const linux.inotify_event = @ptrCast(@alignCast(w.stage[w.head..].ptr)); + const total = @sizeOf(linux.inotify_event) + ev.len; + if (w.head + total <= w.tail) { + const kind: ?Kind = if (ev.mask & linux.IN.Q_OVERFLOW != 0) + // Events were dropped: the consumer's cache may + // be arbitrarily wrong and must rescan. + .overflow + else if (ev.mask & linux.IN.IGNORED != 0) + // The watch itself ended (the directory was + // removed); nothing further will be reported. + .gone + else if (ev.mask & linux.IN.ISDIR != 0) + null + else if (ev.mask & (linux.IN.CREATE | linux.IN.MOVED_TO) != 0) + .added + else if (ev.mask & (linux.IN.DELETE | linux.IN.MOVED_FROM) != 0) + .removed + else + null; + const name = if (ev.getName()) |n| n else ""; + w.head += total; + if (kind) |k| return .{ .name = name, .kind = k }; + continue; + } + } + // Not enough staged: keep the tail, refill. + if (w.head > 0) { + const left = w.tail - w.head; + std.mem.copyForwards(u8, w.stage[0..left], w.stage[w.head..w.tail]); + w.head = 0; + w.tail = left; + } + const rc: isize = @bitCast(linux.read(w.fd, @as([*]u8, &w.stage) + w.tail, w.stage.len - w.tail)); + switch (rawErrno(rc)) { + .SUCCESS => w.tail += @intCast(rc), + .AGAIN => return null, // nothing pending + .INTR => continue, + else => return error.Io, + } + } + } + + pub fn deinit(w: *Self) void { + _ = linux.close(w.fd); + } +} else @compileError("post.Watch requires Linux inotify"); + +// ---- tests ---- + +const testing = std.testing; + +/// A scratch registry: a fake environment block whose XDG_RUNTIME_DIR is +/// a per-test temporary directory. +const Scratch = struct { + dir: testing.TmpDir, + path_buf: [std.fs.max_path_bytes]u8 = undefined, + env_buf: [std.fs.max_path_bytes]u8 = undefined, + env: [2]?[*:0]const u8 = undefined, + + fn start(s: *Scratch) !void { + const io = testing.io; + s.dir = testing.tmpDir(.{}); + errdefer s.dir.cleanup(); + const len = try s.dir.dir.realPath(io, &s.path_buf); + const value = try std.fmt.bufPrintZ(&s.env_buf, "XDG_RUNTIME_DIR={s}", .{s.path_buf[0..len]}); + s.env[0] = @ptrCast(value.ptr); + s.env[1] = null; + } + + fn end(s: *Scratch) void { + s.dir.cleanup(); + } + + fn envp(s: *Scratch) Env { + return @ptrCast(&s.env); + } +}; + +test "post getenv: the block is scanned, keys match whole" { + const env = [_:null]?[*:0]const u8{ "A=1", "XDG_RUNTIME_DIR=/run/user/1000", "X=", "XDG=no" }; + try testing.expectEqualStrings("/run/user/1000", getenv(@ptrCast(&env), "XDG_RUNTIME_DIR").?); + try testing.expectEqualStrings("no", getenv(@ptrCast(&env), "XDG").?); + try testing.expect(getenv(@ptrCast(&env), "NOPE") == null); + try testing.expect(getenv(@ptrCast(&env), "PAT") == null); + try testing.expect(getenv(@ptrCast(&env), "A") != null); +} + +test "post legalName: traversal, empties and the socket budget" { + try testing.expect(legalName("demo")); + try testing.expect(legalName("...")); + try testing.expect(legalName("a" ** max_name_len)); + // Empty and dot names are not postable. + try testing.expect(!legalName("")); + try testing.expect(!legalName(".")); + try testing.expect(!legalName("..")); + // Path traversal is THE attack: no separators, no NUL. + try testing.expect(!legalName("a/b")); + try testing.expect(!legalName("/")); + try testing.expect(!legalName("..\x00..")); + try testing.expect(!legalName("a\x00b")); + // The name must fit a Unix socket path. + try testing.expect(!legalName("a" ** (max_name_len + 1))); + try testing.expect(!legalName("a" ** 255)); +} + +test "post registryPath: unset XDG is refused, no /tmp fallback" { + const empty = [_:null]?[*:0]const u8{"PATH=/bin"}; + var buf: [sun_path_len]u8 = undefined; + try testing.expectError(error.NotRuntimeDir, registryPath(@ptrCast(&empty), "demo", &buf)); + try testing.expectError(error.NotRuntimeDir, registryDir(@ptrCast(&empty), &buf)); + const empty_value = [_:null]?[*:0]const u8{"XDG_RUNTIME_DIR="}; + try testing.expectError(error.NotRuntimeDir, registryPath(@ptrCast(&empty_value), "demo", &buf)); +} + +test "post registryPath: shape, budget and caller buffer" { + const env = [_:null]?[*:0]const u8{"XDG_RUNTIME_DIR=/run/user/1000"}; + const e: Env = @ptrCast(&env); + var buf: [sun_path_len]u8 = undefined; + const p = try registryPath(e, "demo", &buf); + try testing.expectEqualStrings("/run/user/1000/9p/demo", p); + try testing.expectEqual(@as(u8, 0), buf[p.len]); + var dbuf: [32]u8 = undefined; + try testing.expectEqualStrings("/run/user/1000/9p", try registryDir(e, &dbuf)); + // Illegal names never build a path at all. + try testing.expectError(error.IllegalName, registryPath(e, "a/b", &buf)); + try testing.expectError(error.IllegalName, registryPath(e, "..", &buf)); + // The 108-byte sockaddr budget, checked against the real prefix. + const long_env = [_:null]?[*:0]const u8{"XDG_RUNTIME_DIR=/tmp/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}; + const name = "n" ** max_name_len; // legal against the 32-byte XDG budget + try testing.expectError(error.NameTooLong, registryPath(@ptrCast(&long_env), name, &buf)); + // Too small a caller buffer is NoSpace, not a truncation. + var tiny: [4]u8 = undefined; + try testing.expectError(error.NoSpace, registryPath(e, "demo", &tiny)); +} + +test "post stale protocol: fresh, stale, live and non-socket entries" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var pbuf: [sun_path_len]u8 = undefined; + + // A fresh post binds and creates the 0o750 registry directory. + var first = try post(io, env, "demo", 4, &pbuf); + // While the listener is open the name is live and owned. + try testing.expect(probe(first.path) == .live); + try testing.expectError(error.AlreadyPosted, post(io, env, "demo", 4, &pbuf)); + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(env, &dbuf); + const st = try Io.Dir.statFile(.cwd(), io, reg, .{}); + try testing.expect(st.kind == .directory); + try testing.expect(st.permissions.toMode() & 0o750 == 0o750); + + // The listener dies without unlinking: the entry is stale, and a + // new post unlinks and replaces it. + first.server.deinit(io); + try testing.expect(probe(first.path) == .stale); + var revived = try post(io, env, "demo", 4, &pbuf); + try testing.expect(probe(revived.path) == .live); + unpost(io, revived.path, revived.inode); + revived.server.deinit(io); + + // A non-socket entry is refused, never deleted. + const plain = try registryPath(env, "plain", &pbuf); + { + var f = try Io.Dir.createFileAbsolute(io, plain, .{}); + f.close(io); + } + try testing.expectError(error.NotSocket, post(io, env, "plain", 4, &pbuf)); + const after = try Io.Dir.statFile(.cwd(), io, plain, .{}); + try testing.expect(after.kind == .file); + + // dial tells the three apart. + try testing.expectError(error.NotPosted, dial(io, env, "missing")); + try testing.expectError(error.Stale, dial(io, env, "plain")); + var live = try post(io, env, "dials", 4, &pbuf); + var stream = try dial(io, env, "dials"); + stream.close(io); + // The listener dies without unposting: the entry stays, stale. + live.server.deinit(io); + try testing.expectError(error.Stale, dial(io, env, "dials")); + unpost(io, live.path, live.inode); +} + +test "post posted: listing stages the raw names, missing lists empty" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var stage: [512]u8 = undefined; + var pbuf: [sun_path_len]u8 = undefined; + + // Before anything posts (or even creates the registry): empty. + var names = try posted(io, env, &stage); + try testing.expect(names.next() == null); + + // Two posts and a decoy regular file: the listing is raw. + var a = try post(io, env, "alpha", 4, &pbuf); + var b = try post(io, env, "beta", 4, &pbuf); + const decoy = try registryPath(env, "decoy", &pbuf); + { + var f = try Io.Dir.createFileAbsolute(io, decoy, .{}); + f.close(io); + } + names = try posted(io, env, &stage); + var seen: usize = 0; + var has_alpha = false; + var has_beta = false; + var has_decoy = false; + while (names.next()) |name| { + seen += 1; + has_alpha = has_alpha or std.mem.eql(u8, name, "alpha"); + has_beta = has_beta or std.mem.eql(u8, name, "beta"); + has_decoy = has_decoy or std.mem.eql(u8, name, "decoy"); + } + try testing.expectEqual(@as(usize, 3), seen); + try testing.expect(has_alpha and has_beta and has_decoy); + + unpost(io, a.path, a.inode); + a.server.deinit(io); + unpost(io, b.path, b.inode); + b.server.deinit(io); + try Io.Dir.deleteFileAbsolute(io, decoy); +} + +test "post watch: names appear and disappear" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(env, &dbuf); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, reg, .fromMode(0o750)); + + var w = try Watch.init(); + defer w.deinit(); + try w.add(reg); + try testing.expect((try w.next()) == null); + + var pbuf: [sun_path_len]u8 = undefined; + var a = try post(io, env, "alpha", 4, &pbuf); + const seen = try w.next(); + try testing.expect(seen != null); + try testing.expectEqualStrings("alpha", seen.?.name); + try testing.expectEqual(Watch.Kind.added, seen.?.kind); + unpost(io, a.path, a.inode); + a.server.deinit(io); + const gone = try w.next(); + try testing.expect(gone != null); + try testing.expectEqualStrings("alpha", gone.?.name); + try testing.expectEqual(Watch.Kind.removed, gone.?.kind); + try testing.expect((try w.next()) == null); +} + +test "post unpost is ownership-checked: a re-posted name survives a late unpost" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var pbuf: [sun_path_len]u8 = undefined; + + // A posts; the socket file is lost behind its back; B re-posts. + var a = try post(io, env, "svc", 4, &pbuf); + try Io.Dir.deleteFileAbsolute(io, a.path); + var b = try post(io, env, "svc", 4, &pbuf); + try testing.expect(probe(b.path) == .live); + + // A's late unpost (a delayed stop, say) must leave B's entry alone. + unpost(io, a.path, a.inode); + const st = try Io.Dir.statFile(.cwd(), io, b.path, .{}); + try testing.expect(st.kind == .unix_domain_socket); + try testing.expect(probe(b.path) == .live); + + // B's own unpost still works — and is idempotent. + unpost(io, b.path, b.inode); + unpost(io, b.path, b.inode); + try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, b.path, .{})); + a.server.deinit(io); + b.server.deinit(io); +} + +test "post names: a 255-byte entry stages and iterates (NAME_MAX)" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(env, &dbuf); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, reg, .fromMode(0o750)); + // A 255-byte name can never be a socket (108-byte budget), but an + // attacker can drop one in the registry; the raw listing must carry + // it without overflowing the staged `len:u8` record. + const fat = "z" ** 255; + var fat_buf: [std.fs.max_path_bytes]u8 = undefined; + const fat_path = try std.fmt.bufPrintZ(&fat_buf, "{s}/{s}", .{ reg, fat }); + { + var f = try Io.Dir.createFileAbsolute(io, fat_path, .{}); + f.close(io); + } + var stage: [1024]u8 = undefined; + var names = try posted(io, env, &stage); + var seen: usize = 0; + var have_fat = false; + while (names.next()) |name| { + seen += 1; + have_fat = have_fat or std.mem.eql(u8, name, fat); + } + try testing.expectEqual(@as(usize, 1), seen); + try testing.expect(have_fat); + try Io.Dir.deleteFileAbsolute(io, fat_path); +} + +test "post long name: max_name_len posts, probes, dials and unposts" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + // The budget assumes a conforming (short) XDG_RUNTIME_DIR; the + // testing tmpdir's path is longer, so the boundary needs a short + // scratch directory of its own. + const io = testing.io; + var xdg_buf: [64]u8 = undefined; + const xdg = try std.fmt.bufPrintZ(&xdg_buf, "/tmp/.p9t{d}", .{linux.getpid()}); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, xdg, .fromMode(0o700)); + var env_buf: [96]u8 = undefined; + const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{xdg}); + const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; + const envp: Env = @ptrCast(&env); + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(envp, &dbuf); + defer { + _ = linux.rmdir(reg.ptr); + _ = linux.rmdir(xdg.ptr); + } + var pbuf: [sun_path_len]u8 = undefined; + const name = "l" ** max_name_len; + try testing.expect(legalName(name)); + var p = try post(io, envp, name, 4, &pbuf); + try testing.expect(probe(p.path) == .live); + var stream = try dial(io, envp, name); + stream.close(io); + unpost(io, p.path, p.inode); + p.server.deinit(io); + // One past the cap is not even a path. + const too_long = "l" ** (max_name_len + 1); + try testing.expect(!legalName(too_long)); + try testing.expectError(error.IllegalName, registryPath(envp, too_long, &pbuf)); +} + +test "post watch: queue overflow surfaces, a deleted registry is gone" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(env, &dbuf); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, reg, .fromMode(0o750)); + var w = try Watch.init(); + defer w.deinit(); + try w.add(reg); + + // Overwhelm the kernel queue (default 16384 events): the overflow + // must be surfaced, not silently skipped, or a caching consumer + // stays stale forever. + var name_buf: [std.fs.max_path_bytes]u8 = undefined; + const spam = 20000; + for (0..spam) |i| { + const p = try std.fmt.bufPrintZ(&name_buf, "{s}/q{d}", .{ reg, i }); + var f = try Io.Dir.createFileAbsolute(io, p, .{}); + f.close(io); + } + var overflow = false; + var events: usize = 0; + while (try w.next()) |ev| { + events += 1; + if (ev.kind == .overflow) overflow = true; + } + try testing.expect(overflow); + + // The registry directory itself is replaced: the watch must say so. + for (0..spam) |i| { + const p = try std.fmt.bufPrint(&name_buf, "{s}/q{d}", .{ reg, i }); + try Io.Dir.deleteFileAbsolute(io, p); + } + while (try w.next()) |_| {} + _ = linux.rmdir(reg.ptr); // empty now + try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, reg, .{})); + const gone = try w.next(); + try testing.expect(gone != null); + try testing.expectEqual(Watch.Kind.gone, gone.?.kind); + try testing.expect((try w.next()) == null); + // Re-added, the (recreated) registry is watched again. + _ = try Io.Dir.createDirPathStatus(.cwd(), io, reg, .fromMode(0o750)); + try w.add(reg); + { + var f = try Io.Dir.createFileAbsolute(io, try std.fmt.bufPrintZ(&name_buf, "{s}/back", .{reg}), .{}); + f.close(io); + } + const back = try w.next(); + try testing.expect(back != null); + try testing.expectEqualStrings("back", back.?.name); + try Io.Dir.deleteFileAbsolute(io, try std.fmt.bufPrint(&name_buf, "{s}/back", .{reg})); +} diff --git a/src/root.zig b/src/root.zig index 586a41c..fb0bf5c 100644 --- a/src/root.zig +++ b/src/root.zig @@ -50,6 +50,10 @@ test { pub const http = @import("http.zig"); pub const transport = @import("transport.zig"); pub const serve = @import("serve.zig"); +pub const post = @import("post.zig"); +test { + _ = @import("post.zig"); +} pub const Quic = @import("quic.zig").Quic; test { _ = @import("session_test.zig"); diff --git a/src/serve.zig b/src/serve.zig index 255e2ff..19a8ca2 100644 --- a/src/serve.zig +++ b/src/serve.zig @@ -12,6 +12,7 @@ const std = @import("std"); const Io = std.Io; const fs = @import("fs.zig"); const transport = @import("transport.zig"); +const post = @import("post.zig"); /// Comptime bounds of one runner. pub const Limits = struct { @@ -72,6 +73,7 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits }; pub const ListenError = error{TooManyListeners} || Io.net.IpAddress.ListenError || Io.net.UnixAddress.ListenError || Io.net.UnixAddress.InitError || Io.ConcurrentError; + pub const ListenPostedError = post.PostError || error{TooManyListeners} || Io.ConcurrentError; io: Io, root: u64, @@ -81,6 +83,14 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits conns: [limits.connections]Conn, listeners: [limits.listeners]Io.net.Server, nlisteners: usize, + /// The registry socket of a `listenPosted`, zero-terminated; + /// `stop()` unlinks it (unpost on stop) — but only while it is + /// still this runner's entry (`posted_ino`). + posted_path: [transport.sun_path_len + 1]u8 = @splat(0), + posted_len: usize = 0, + /// The bound registry entry's inode, the ownership proof for + /// the unpost in `stop()`. + posted_ino: u64 = 0, /// Accept tasks and connection tasks; `stop()` cancels it. group: Io.Group, /// Guards `Conn.used`. @@ -292,6 +302,7 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits r.seed = o.seed; r.handler = o.handler; r.greet_timeout_ms = o.greet_timeout_ms; + r.posted_len = 0; r.nlisteners = 0; r.group = .init; r.slots = .init; @@ -307,15 +318,46 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits pub fn listen(r: *Self, address: transport.Address, backlog: u31) ListenError!Io.net.IpAddress { if (r.nlisteners == limits.listeners) return error.TooManyListeners; if (r.stopping.load(.acquire)) return error.TooManyListeners; + var server = try transport.listen(r.io, address, backlog); + errdefer server.deinit(r.io); + try r.startListener(server); + return server.socket.address; + } + + /// Posts the runner on the registry socket + /// `$XDG_RUNTIME_DIR/9p/<name>` and starts accepting on it: + /// `post.post` creates the 0o750 registry directory and runs the + /// stale protocol (a refused entry is replaced; a live server + /// owning the name is `AlreadyPosted`; a non-socket entry is + /// never deleted). One posted name per runner: a second + /// `listenPosted` is `AlreadyPosted` (its socket would otherwise + /// be orphaned in the registry — nothing would unpost it). + /// `stop()` unposts — the socket is unlinked when the runner + /// stops, as long as the entry is still the runner's own. + pub fn listenPosted(r: *Self, env: post.Env, name: []const u8, backlog: u31) ListenPostedError!void { + if (r.nlisteners == limits.listeners) return error.TooManyListeners; + if (r.stopping.load(.acquire)) return error.TooManyListeners; + if (r.posted_len != 0) return error.AlreadyPosted; + var p = try post.post(r.io, env, name, backlog, &r.posted_path); + errdefer { + post.unpost(r.io, p.path, p.inode); + p.server.deinit(r.io); + } + try r.startListener(p.server); + r.posted_len = p.path.len; + r.posted_ino = p.inode; + } + + /// Registers a bound listener and starts its accept task. + fn startListener(r: *Self, server: Io.net.Server) Io.ConcurrentError!void { const i = r.nlisteners; - r.listeners[i] = try transport.listen(r.io, address, backlog); + r.listeners[i] = server; errdefer r.listeners[i].deinit(r.io); r.nlisteners += 1; r.group.concurrent(r.io, acceptLoop, .{ r, i }) catch |err| { r.nlisteners -= 1; return err; }; - return r.listeners[i].socket.address; } /// Connections held right now. @@ -333,13 +375,21 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits for (&r.conns) |*c| if (c.live()) c.close(); } - /// Stops accepting, hangs every connection up, waits for their tasks - /// and closes the listeners. Idempotent; the runner is spent after. + /// Stops accepting, hangs every connection up, waits for their + /// tasks and closes the listeners. A posted listener is unposted + /// — its registry socket is unlinked, but only while the entry + /// is still the runner's own: a name that was re-posted by + /// another server (this one's socket file having been lost) + /// survives the stop. Idempotent; the runner is spent after. pub fn stop(r: *Self) void { if (r.stopping.swap(true, .acq_rel)) return; r.group.cancel(r.io); for (r.listeners[0..r.nlisteners]) |*l| l.deinit(r.io); r.nlisteners = 0; + if (r.posted_len != 0) { + post.unpost(r.io, r.posted_path[0..r.posted_len :0], r.posted_ino); + r.posted_len = 0; + } } fn acceptLoop(r: *Self, i: usize) void { @@ -846,3 +896,120 @@ test "serve: a backend answered from another thread under lock()" { rig.runner.stop(); try expectHangup(tc.one(.{ .stat = .{ .fid = 0 } })); } + +test "serve: listenPosted serves the registry name and stop() unposts" { + if (@import("builtin").os.tag != .linux) return error.SkipZigTest; + var rig: Rig = .{ .dir = undefined }; + const io = testing.io; + // A scratch registry: XDG_RUNTIME_DIR is the rig's own temp dir. + rig.dir = testing.tmpDir(.{}); + errdefer rig.dir.cleanup(); + var real_buf: [std.fs.max_path_bytes]u8 = undefined; + const len = try rig.dir.dir.realPath(io, &real_buf); + var env_buf: [std.fs.max_path_bytes]u8 = undefined; + const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{real_buf[0..len]}); + const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; + const envp: post.Env = @ptrCast(&env); + + rig.stub = .{}; + rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &rig.stub, .serve = serveNow } }); + try rig.runner.listenPosted(envp, "posted", 4); + var path_buf: [transport.sun_path_len]u8 = undefined; + const path = try post.registryPath(envp, "posted", &path_buf); + + // The name is posted, live, and serves a full 9P session. + try testing.expect(post.probe(path) == .live); + var tc: TestClient = .{ .io = undefined, .stream = undefined }; + try tc.open(io, .{ .unix = path }); + defer tc.close(); + try tc.handshake(); + try tc.readIndex(1); + + // A second post of the same name is refused while the runner lives. + var pbuf: [transport.sun_path_len]u8 = undefined; + try testing.expectError(error.AlreadyPosted, post.post(io, envp, "posted", 4, &pbuf)); + + // The listing sees it; stop() unposts and the entry disappears. + var stage: [512]u8 = undefined; + var names = try post.posted(io, envp, &stage); + var seen = false; + while (names.next()) |n| seen = seen or std.mem.eql(u8, n, "posted"); + try testing.expect(seen); + rig.runner.stop(); + try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, path, .{})); + names = try post.posted(io, envp, &stage); + try testing.expect(names.next() == null); + rig.dir.cleanup(); +} + +test "serve: a second listenPosted is refused; stop() never unposts another's name" { + if (@import("builtin").os.tag != .linux) return error.SkipZigTest; + var rig: Rig = .{ .dir = undefined }; + const io = testing.io; + rig.dir = testing.tmpDir(.{}); + errdefer rig.dir.cleanup(); + var real_buf: [std.fs.max_path_bytes]u8 = undefined; + const len = try rig.dir.dir.realPath(io, &real_buf); + var env_buf: [std.fs.max_path_bytes]u8 = undefined; + const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{real_buf[0..len]}); + const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; + const envp: post.Env = @ptrCast(&env); + + rig.stub = .{}; + rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &rig.stub, .serve = serveNow } }); + try rig.runner.listenPosted(envp, "twice", 4); + // One posted name per runner: a second would overwrite the first's + // path and orphan its socket in the registry. + try testing.expectError(error.AlreadyPosted, rig.runner.listenPosted(envp, "twice", 4)); + try testing.expectError(error.AlreadyPosted, rig.runner.listenPosted(envp, "other", 4)); + + var path_buf: [transport.sun_path_len]u8 = undefined; + const path = try post.registryPath(envp, "twice", &path_buf); + // The runner's socket file is lost behind its back (rm, crash + // cleanup), and another server takes the now-free name. + try Io.Dir.deleteFileAbsolute(io, path); + var thief = try post.post(io, envp, "twice", 4, &path_buf); + try testing.expect(post.probe(thief.path) == .live); + // stop() unposts only what it still owns: the thief survives. + rig.runner.stop(); + const st = try Io.Dir.statFile(.cwd(), io, thief.path, .{}); + try testing.expect(st.kind == .unix_domain_socket); + try testing.expect(post.probe(thief.path) == .live); + post.unpost(io, thief.path, thief.inode); + thief.server.deinit(io); + rig.dir.cleanup(); +} + +test "serve: listenPosted beside listen(): stop unposts only the registry name" { + if (@import("builtin").os.tag != .linux) return error.SkipZigTest; + var rig: Rig = .{ .dir = undefined }; + const io = testing.io; + rig.dir = testing.tmpDir(.{}); + errdefer rig.dir.cleanup(); + var real_buf: [std.fs.max_path_bytes]u8 = undefined; + const len = try rig.dir.dir.realPath(io, &real_buf); + var env_buf: [std.fs.max_path_bytes]u8 = undefined; + const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{real_buf[0..len]}); + const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; + const envp: post.Env = @ptrCast(&env); + + rig.stub = .{}; + rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &rig.stub, .serve = serveNow } }); + // A plain Unix listener (the application's path policy) beside the + // posted name: both listener slots fill. + var unix_buf: [std.fs.max_path_bytes]u8 = undefined; + const unix_path = try std.fmt.bufPrintZ(&unix_buf, "{s}/plain.sock", .{real_buf[0..len]}); + _ = try rig.runner.listen(.{ .unix = unix_path }, 4); + try rig.runner.listenPosted(envp, "mixed", 4); + + var path_buf: [transport.sun_path_len]u8 = undefined; + const posted_path = try post.registryPath(envp, "mixed", &path_buf); + try testing.expect(post.probe(posted_path) == .live); + rig.runner.stop(); + // The posted name is unposted; the plain path is the application's. + try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, posted_path, .{})); + const plain_st = try Io.Dir.statFile(.cwd(), io, unix_path, .{}); + try testing.expect(plain_st.kind == .unix_domain_socket); + try Io.Dir.deleteFileAbsolute(io, unix_path); + rig.dir.cleanup(); +} |
