diff options
Diffstat (limited to '9ns')
| -rw-r--r-- | 9ns/README.md | 61 | ||||
| -rw-r--r-- | 9ns/docs/DESIGN.md | 170 | ||||
| -rw-r--r-- | 9ns/src/bridge.zig | 212 | ||||
| -rw-r--r-- | 9ns/src/fuse.zig | 56 | ||||
| -rw-r--r-- | 9ns/src/main.zig | 133 | ||||
| -rw-r--r-- | 9ns/src/nine.zig | 438 | ||||
| -rw-r--r-- | 9ns/src/ns.zig | 83 | ||||
| -rwxr-xr-x | 9ns/test/adv_bridge_hostile.py | 11 | ||||
| -rwxr-xr-x | 9ns/test/adv_bridge_hostile.sh | 6 | ||||
| -rwxr-xr-x | 9ns/test/adv_bridge_interrupt.sh | 136 | ||||
| -rwxr-xr-x | 9ns/test/adv_bridge_semantics.sh | 7 | ||||
| -rwxr-xr-x | 9ns/test/adv_bridge_stress.sh | 2 | ||||
| -rwxr-xr-x | 9ns/test/adversarial.sh | 2 | ||||
| -rwxr-xr-x | 9ns/test/integration.sh | 27 |
14 files changed, 1198 insertions, 146 deletions
diff --git a/9ns/README.md b/9ns/README.md index 04edffb..7d6babe 100644 --- a/9ns/README.md +++ b/9ns/README.md @@ -4,15 +4,20 @@ Mount a 9P2000 file tree into a fresh mount namespace and run a program in it, as a plain user, without touching the host's mount table. ```sh -9ns --unix /run/user/1000/acme -- fish # a shell that sees the tree at /mnt/9p -9ns --tcp 127.0.0.1:564 -- claude # an agent that sees it too -9ns --spawn '9proc-demo --stdio' -- bash # start the server yourself, talk over a socketpair +9ns --unix /run/user/1000/acme -- fish # a shell that sees the tree at /mnt/9p/acme +9ns --tcp 127.0.0.1:564 -- claude # an agent that sees it at /mnt/9p/tcp-127.0.0.1-564 +9ns --spawn '9proc-demo --stdio' -- bash # start the server yourself; /mnt/9p/9proc-demo +9ns --unix /tmp/9debug.sock --name dbg -- bash # pick the name: /mnt/9p/dbg ``` Inside, the tree is ordinary files: `ls`, `cat`, `echo x > ctl`, editors, -`find`, `rsync`, whatever. `$NINE_MOUNT` tells programs where it is -(default `/mnt/9p`). When the program exits, 9ns exits with its status -and the namespace, mount and connection disappear. +`find`, `rsync`, whatever. Every mount lives under `/mnt/9p/<name>`, where +the name comes from `--name` or, by default, from the transport (the socket's +basename, `tcp-IP-PORT`, the spawned command's basename, `fdN`); `--mount +PATH` puts it anywhere else. `$NINE_MOUNT` tells programs where it is. When +the program exits, 9ns exits with its status and the namespace, mount and +connection disappear. Nesting works: `9ns --name a -- 9ns --name b -- fish` +gives a shell that sees both `/mnt/9p/a` and `/mnt/9p/b`. ## How it works @@ -24,7 +29,7 @@ of the kernel FUSE protocol needed, straight from `linux/fuse.h`. ``` program (fish/bash/claude) 9ns (parent) 9P server in a new user+mount namespace │ - /mnt/9p ── FUSE ──▶ kernel ────▶│ fuse.zig ─▶ bridge.zig ─▶ nine.zig ──▶ unix / tcp / socketpair + /mnt/9p/<name> ─FUSE─▶ kernel ─▶│ fuse.zig ─▶ bridge.zig ─▶ nine.zig ──▶ unix / tcp / socketpair │ (framing) (translation) (cloud9 Client) ``` @@ -43,10 +48,15 @@ Files are opened with `FOPEN_DIRECT_IO`, so synthetic files that report length travels inside the 9P open mode (`OTRUNC`) rather than as a separate truncate. Repeated lookups of the same qid map to the same inode. -If `/mnt/9p` does not exist and cannot be created (the normal case), 9ns -mounts a tmpfs over `/mnt` *inside the namespace only* and bind-mounts every -existing entry of `/mnt` back into it, so nothing is hidden. Pass `--mount DIR` -to use any other directory. +The mountpoint `/mnt/9p/<name>` normally does not exist and cannot be +created by a plain user. 9ns then walks down the path to the deepest existing +directory (`/mnt` on a host without `/mnt/9p`, `/mnt/9p` if the host has one), +mounts a tmpfs over it *inside the namespace only*, bind-mounts every existing +entry of that directory back into the tmpfs so nothing is hidden, and creates +the missing components inside. Inside a 9ns namespace `/mnt/9p` is already +that writable tmpfs directory, so a nested 9ns just adds its own name next to +the outer mount (the same name mounts over it). Pass `--mount DIR` to use any +other directory. ## Building and testing @@ -87,7 +97,9 @@ Transport (exactly one): --spawn CMD run CMD via /bin/sh -c with a socketpair on its stdin/stdout Options: - --mount PATH mountpoint inside the new namespace (default /mnt/9p) + --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) --uname NAME 9P user name (default $USER) --aname NAME 9P tree to attach (default "") --msize BYTES maximum 9P message size to request (default 131072, max 16 MiB) @@ -101,6 +113,17 @@ PROGRAM defaults to `$SHELL`. Exit status is the program's (`128+signal` if it was killed); 125 means 9ns itself failed (usage, connect, namespace, mount); 126/127 are exec failures as usual. +The default name comes from the transport: `--unix PATH` → the basename of +PATH without a trailing `.sock`, `.9p` or `.socket` (`/tmp/9debug.sock` → +`9debug`); `--tcp IP:PORT` → `tcp-IP-PORT` with every `:` turned into `-` +(`[::1]:564` → `tcp---1-564`); `--spawn CMD` → the basename of CMD's first +word (`/x/9proc-demo --stdio` → `9proc-demo`); `--fd N` → `fdN`. A name must +be a single path component (not empty, no `/`, not `.` or `..`); an invalid +`--name` is a usage error, an unusable derived name falls back to `9p`. With +`--mount PATH` the name is simply PATH's last component. The program gets +`NINE_MOUNT=<mountpoint>` (replacing any inherited value, so a nested 9ns +overwrites it). + ## 9proc-demo: a demo 9P server `9proc-demo` is a single-binary 9P2000 server whose file tree is the binary @@ -131,14 +154,18 @@ zig-out/bin/9ns --unix /tmp/intro.sock -- sh -c ' echo "fib 20" > $NINE_MOUNT/runtime/ctl; cat $NINE_MOUNT/runtime/ctl' ``` -The same command with `claude -p "explore /mnt/9p ..."` as the program gives an -agent a live, file-shaped view into a running process; that is the intended -use. +The same command with `claude -p "explore /mnt/9p/intro ..."` as the program +gives an agent a live, file-shaped view into a running process; that is the +intended use. ## Limitations -* One 9P request is in flight at a time; a server read that blocks (event - files) stalls the mount until it returns. +* One 9P request is in flight at a time, so a server read that blocks (event + files, consoles) stalls the mount while it is outstanding. The blocked + request can be interrupted, though: killing or Ctrl-C-ing the reader makes + the kernel send `FUSE_INTERRUPT`, which 9ns forwards as 9P `Tflush`; + servers that honour it release the reader with `EINTR` at once, servers + that ignore it keep the mount stalled until they answer. * Base 9P2000 only: no symlinks, ownership, xattrs or locks. Every file is reported as owned by the invoking user. Cross-directory rename is `EXDEV`. * No PID namespace and no `/proc` remount. `--tcp` takes IP literals only diff --git a/9ns/docs/DESIGN.md b/9ns/docs/DESIGN.md index 7f943a8..f1589d0 100644 --- a/9ns/docs/DESIGN.md +++ b/9ns/docs/DESIGN.md @@ -15,7 +15,7 @@ So 9ns is a tiny FUSE server that speaks 9P2000 to the real server: ``` program (fish/bash/claude) 9ns (parent) 9P server in new user+mount namespace │ (9proc-demo, - /mnt/9p ─── FUSE ───▶ kernel ──▶│ fuse.zig ──▶ bridge.zig ──▶ nine.zig ──▶ ramfs, ...) + /mnt/9p/<name> ─FUSE─▶ kernel ─▶│ fuse.zig ──▶ bridge.zig ──▶ nine.zig ──▶ ramfs, ...) │ (framing) (translation) (cloud9 Client) ``` @@ -79,7 +79,8 @@ protocol we need directly against `/usr/include/linux/fuse.h`. 8. `statx` of the mountpoint: this forces one GETATTR, which the parent serves. Without it the kernel keeps the root inode's initial uid 0 (unmapped in the namespace) and every create in the root gets `EACCES`. - 9. Sets `NINE_MOUNT=<mountpoint>` in the environment. + 9. Sets `NINE_MOUNT=<mountpoint>` in the environment (replacing any + inherited value; a nested 9ns overwrites it). 10. `execve` of PROGRAM with PATH search (implemented by hand; no libc). Exec failures are reported through the `CLOEXEC` status socket (errno + message); the parent prints them after the serve loop ends. @@ -90,7 +91,10 @@ protocol we need directly against `/usr/include/linux/fuse.h`. signalled). The self-pipe is also watched by the 9P session while a reply is outstanding (`Session.stop_fd` → `error.Stopped`), so a server that never answers cannot keep 9ns alive after the child is gone; a 3 s - watchdog armed from the SIGCHLD handler is the last resort. + watchdog armed from the SIGCHLD handler is the last resort. The FUSE fd + is watched during that wait too (`Session.interrupt`): a `FUSE_INTERRUPT` + for the request being served becomes a `Tflush` (see *Interrupts* under + `src/bridge.zig`). 4. Signals in the parent: `SIGINT`/`SIGQUIT` ignored (the child owns the tty and gets them itself); `SIGTERM`/`SIGHUP` forwarded to the child; `SIGPIPE` ignored; `SIGCHLD` → self-pipe. @@ -100,17 +104,43 @@ copy keeps the connection alive. ### Mountpoint policy -Default mountpoint: `/mnt/9p`. A relative `--mount` is resolved against cwd. +Default mountpoint: `/mnt/9p/<name>`, where the name is `--name NAME` or is +derived from the transport (`main.zig`, `defaultName`): -* If the path is a directory: use it. -* Else try `mkdir`. If that fails with `EACCES`/`EPERM`/`EROFS` (the normal - case for `/mnt/9p` as a plain user), **shadow the parent directory**: - open an fd to the parent, mount a `tmpfs` over it, then recreate every - existing entry inside the tmpfs: directories → `mkdir` + bind mount from - `/proc/self/fd/<fd>/<name>`; symlinks → `readlinkat` + `symlink`; anything - else → empty regular file + bind mount. Then `mkdir` the target inside. - Refuse (with a clear message) if the parent has more than 4096 entries or - is `/`. This only affects the new namespace. +| transport | default name | +|---|---| +| `--unix PATH` | basename of PATH with one trailing `.sock`/`.9p`/`.socket` removed (`/tmp/9debug.sock` → `9debug`) | +| `--tcp IP:PORT` | `tcp-IP-PORT` with `:` → `-` (brackets are already gone: `[::1]:564` → `tcp---1-564`) | +| `--spawn CMD` | basename of the first whitespace-separated word of CMD | +| `--fd N` | `fdN` | + +A name is one path component: non-empty, no `/`, no NUL, not `.`/`..`. An +invalid `--name` is a usage error (125); a derived name that is not valid +(empty basename, `..`) falls back to `9p`. `--mount PATH` overrides all of +this (the name is then PATH's last component and is not used for anything); +a relative `--mount` is resolved against cwd; `/` is rejected. + +`ns.ensureMountpoint` then makes the path a directory inside the new +namespace: + +* If the path is a directory: use it (a nested 9ns with the same name + therefore mounts over the outer mount at that path). +* Else walk up to the deepest existing ancestor (which must be a directory; + a dangling symlink anywhere is an error) and `mkdir` the missing + components under it one by one. If the first of those fails with + `EACCES`/`EPERM`/`EROFS` (the normal case for `/mnt/9p/<name>` as a plain + user), **shadow that ancestor**: open an fd to it, mount a `tmpfs` over + it, then recreate every existing entry inside the tmpfs: directories → + `mkdir` + bind mount from `/proc/self/fd/<fd>/<name>`; symlinks → + `readlinkat` + `symlink`; anything else → empty regular file + bind mount. + Then create the missing components inside. Refuse (with a clear message) + if the ancestor is `/` (also through `/proc/self/root`), is `/proc` or + below it, or has more than 4096 entries. This only affects the new + namespace. So on a host without `/mnt/9p` the shadow goes over `/mnt` and + `9p/<name>` is created inside; with a root-owned `/mnt/9p` it goes over + `/mnt/9p`; inside a 9ns namespace `/mnt/9p` is a directory of the outer + tmpfs owned by our uid, so a nested 9ns just creates `<name>` next to the + outer mount and both are visible. * Else fail with the errno and a hint to pass `--mount` an existing dir. ## Module contracts @@ -151,6 +181,10 @@ I/O helpers (blocking fd, no allocation beyond the caller's buffer): pub const Request = struct { header: InHeader, body: []const u8 }; /// One kernel request. Returns null on ENODEV (unmounted). Retries EINTR/EAGAIN/ENOENT. pub fn readRequest(fd: i32, buf: []u8) !?Request; +/// Same, one read(2) only: EINTR/EAGAIN/ENOENT → error.Retry (for reads that follow a poll). +pub fn readRequestOnce(fd: i32, buf: []u8) !?Request; +/// O_NONBLOCK on the device, so a request withdrawn between poll and read cannot block us. +pub fn setNonblocking(fd: i32) !void; /// Success reply: header + concatenated payload slices, single writev. pub fn reply(fd: i32, unique: u64, payloads: []const []const u8) !void; /// Error reply: negative errno. @@ -172,11 +206,20 @@ single-threaded). Fids are allocated from a free list. ```zig pub const Address = union(enum) { unix: []const u8, tcp: struct { host: []const u8, port: u16 }, fd: i32 }; +/// Owner's hook into the reply wait (the bridge's FUSE fd); see "Interrupts" under bridge.zig. +pub const Interrupt = struct { + ctx: *anyopaque, + watch: *const fn (ctx) i32, // fd to poll alongside the socket, or -1 + onReadable: *const fn (ctx) Session.Error!bool, // consume it; true = flush the request in flight + armed: *const fn (ctx) bool, // a flush was asked for the operation in progress +}; pub const Session = struct { - pub const Error = error{ Nine, Protocol, Io, Closed, TooLarge, OutOfMemory }; + pub const Error = error{ Nine, Protocol, Io, Closed, Stopped, Interrupted, TooLarge, OutOfMemory }; /// After error.Nine, `ename` holds the server's Rerror text (copied, bounded). ename: [256]u8, ename_len: usize, msize: u32, + stop_fd: i32 = -1, // readable → pending rpc fails with error.Stopped + interrupt: ?Interrupt = null, pub fn connect(gpa: std.mem.Allocator, address: Address, msize: u32) !Session; // socket+connect, version pub fn deinit(s: *Session) void; @@ -209,12 +252,20 @@ negotiated msize is `result.version.msize`; if the server answered `"unknown"`, fail with `error.Protocol`. `rpc`: submit, write all of `client.output()` (calling `wrote`), then loop: -`take()`; if null, `read` from the fd into a temp buffer and `push` (push -returns how much fit; the frame is at most msize so it always fits after a -`take`). If the fd returns 0 → `error.Closed`. If the client dies → -`error.Protocol`. A `.fail` result copies the ename and returns `error.Nine`. +`take()`; if null, poll the fd together with `stop_fd` and the interrupt +source's descriptor; when the fd is readable, `read` into a temp buffer and +`push` (push returns how much fit; the frame is at most msize so it always +fits after a `take`). If the fd returns 0 → `error.Closed`. If the client +dies → `error.Protocol`. `stop_fd` readable → `error.Stopped`. When the +interrupt source asks for it, submit `.flush = .{ .oldtag = tag }` and keep +waiting: the original reply → returned normally (the Rflush that follows is +skipped by a later call); the Rflush first → `error.Interrupted`. A `.fail` +result copies the ename and returns `error.Nine`. `read`/`write` chunk loops +stop between chunks once the source reports `armed`, and turn an +`Interrupted` chunk into a short count when earlier chunks moved data. Rerror text → errno mapping (case-insensitive substring, in this order): +`"interrupt"` → `EINTR` (a server answering a flushed request with an error); `"not exist"`, `"not found"`, `"no such"` → `ENOENT`; `"exists"` → `EEXIST`; `"not empty"` → `ENOTEMPTY`; `"not a dir"` → `ENOTDIR`; `"is a dir"` → `EISDIR`; `"permission"`, `"denied"` → `EACCES`; @@ -276,7 +327,7 @@ Op mapping (9P2000 has no symlinks, links, xattrs, locks, mknod): | STATFS | constant `Kstatfs{ bsize = 4096, namelen = 255, frsize = 4096 }` | | ACCESS | `ENOSYS` (kernel stops asking; the server enforces permissions on open) | | READLINK, SYMLINK, LINK, MKNOD, *XATTR, *LK, IOCTL, POLL, BMAP, FALLOCATE, LSEEK, COPY_FILE_RANGE, TMPFILE, STATX | `ENOSYS` | -| INTERRUPT | ignored (reply nothing) | +| INTERRUPT | read while a 9P reply is outstanding: for the request in flight → `Tflush` (see *Interrupts*); otherwise ignored (reply nothing) | | DESTROY | return from `serve` | Attr mapping from `cloud9.Stat`: `mode = (S_IFDIR if DMDIR else S_IFREG) | @@ -301,6 +352,51 @@ rather than poisoning the whole READDIR reply; `length` near 2^64 is clamped to `i64` max; the errno of a failing 9P call is latched before any cleanup clunk overwrites the session's ename. +#### Interrupts (FUSE_INTERRUPT → Tflush) + +One request at a time, but not deaf. `serve` puts the FUSE fd in `O_NONBLOCK` +mode (every read follows a poll) and installs a `nine.Interrupt` source that +`Session.rpc` polls together with the socket and `stop_fd` whenever a reply +is outstanding, including the initial root stat: + +* `onReadable` reads the request the kernel has ready into a second + 8-aligned buffer of `request_buf_len` bytes (`spare_buf`). A + `FUSE_INTERRUPT` whose `InterruptIn.unique` names the request being served + (`cur_unique`) sets `interrupted` and returns true: `rpc` sends + `Tflush(oldtag)` and keeps waiting until either the original reply arrives + (the interrupt raced it: the result is returned as if nothing happened and + the Rflush that follows is swallowed by a later call) or the Rflush does + (`error.Interrupted` → `EINTR`; a server that instead answers the flushed + request with an Rerror containing "interrupt", as Pardes does, lands on the + same errno through the ename table). An INTERRUPT for any other unique is + consumed and dropped (the kernel expects no reply). Any other request + (FORGET, RELEASE, INIT during the root stat, a second process's LOOKUP) is + stashed in a one-slot queue that `serve` dispatches, after swapping the two + buffers, before it polls again; while the slot is full `watch` returns -1, + so a second one cannot arrive. +* The INTERRUPT applies to `cur_unique` only and is consumed when read: the + clunks that unwind a half-done lookup/create/mkdir after an `Interrupted` + walk or stat are ordinary rpcs and are not re-interrupted by it. An + interrupted walk clunks its new fid (the Rflush alone does not say whether + the server bound it); an Rerror'd walk only frees it locally. +* Multi-step operations fail with `EINTR` at whichever rpc was flushed and + release the fids they had allocated (`--debug` prints `fids=N` per + request; the interrupt suite checks it and the server-side count). The + chunk loops (`Session.read/write`, `loadDir`) also stop between chunks + while `armed`, so a reply that won the race cannot lead into another + blocking chunk: a read or write that already moved data returns the partial + count like read(2), one that moved nothing and a directory listing return + `EINTR`. +* The kernel sends one INTERRUPT per request, for a fatal signal (SIGKILL + included) as much as for a caught one, and then waits for the reply; our + `EINTR` is what finally lets the killed task die. + +Limits: a server that ignores Tflush still blocks the mount until it answers +(the hostile `never` mode; `SIGTERM` to 9ns ends the session as before), and +while the stash is full the FUSE fd is not read, so an INTERRUPT that arrives +after another process's request was parked is seen only once the blocked +request completes (a multi-slot stash would lift that). + ### `src/ns.zig` — namespace and process plumbing ```zig @@ -316,7 +412,7 @@ pub const Child = struct { pid: i32 }; /// fork; the child sets up the namespace, mounts, and execs. Returns once exec succeeded /// (status pipe closed) or fails with the child's error (message on stderr). pub fn spawn(gpa: std.mem.Allocator, s: Spawn) !Child; -pub fn ensureMountpoint(path: [:0]const u8) !void; // the shadowing logic, testable alone +pub fn ensureMountpoint(gpa: Allocator, path: [:0]const u8) !void; // walk down, mkdir -p, shadow; testable alone pub fn resolveMountpoint(gpa, path: []const u8) ![:0]u8; // absolute, no trailing slash pub fn findInPath(gpa, envp, name) ![:0]u8; ``` @@ -336,7 +432,9 @@ Transport (exactly one): --fd N already-connected inherited descriptor --spawn CMD run CMD (via /bin/sh -c) with a socketpair on its stdin/stdout Options: - --mount PATH mountpoint inside the new namespace (default /mnt/9p) + --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) --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) @@ -345,9 +443,13 @@ Options: --debug trace FUSE and 9P operations on stderr --help, --version PROGRAM defaults to $SHELL (else /bin/sh). The mountpoint is exported as $NINE_MOUNT. +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. ``` -Exit codes: child's status; 125 for 9ns's own failures (usage, connect, +`--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. ### `../9proc/demo/main.zig` — demo 9P2000 server (binary `9proc-demo`) @@ -373,23 +475,30 @@ Everything under a temp dir. Skips (exit 0 with a notice) when `mv` across dirs fails with `EXDEV`-ish message, `rm`, `rmdir`, 1 MiB random file round trip compared with `sha256sum`, `dd` with odd block sizes, many small files, `find`, exit-status propagation (`exit 7` → 7), - `$NINE_MOUNT` set, nested `9ns` inside `9ns`. + `$NINE_MOUNT` set, nested `9ns` inside `9ns`. The suites pin the + mountpoint with `--mount /mnt/9p`; the naming section then checks the + default `/mnt/9p/<name>` (unix socket basename, `--name`, `--name=`, + invalid names → 125, nested runs with two names both visible under + `/mnt/9p`, the same name twice mounting over, `--spawn` and `--tcp` + derived names, `--mount` beating `--name`). 2. `--spawn "<9proc-demo> --stdio"` variant. 3. `--tcp 127.0.0.1:<port>` variant. 4. If `/usr/lib/plan9/bin/ramfs` exists: `NAMESPACE=$tmp ramfs -s ramfs` creates `$tmp/ramfs`; run the scratch battery against it. -5. `--mount` with an existing dir, with a relative path, and the default - `/mnt/9p` (exercises parent shadowing; verify `/mnt`'s other entries are - still visible inside). +5. `--mount` with an existing dir, with a relative path, `/mnt/9p` and the + default `/mnt/9p/<name>` (both exercise the shadowing of `/mnt`; verify + `/mnt`'s other entries are still visible inside and the host mount table + is untouched). 6. Kill tests: 9ns exits when the child exits; server death during use yields `EIO`, not a hang. ## Verification -`zig build 9ns-test` (unit), `zig build 9ns-itest` (74 end-to-end +`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 -hostile 9P server with ~30 misbehaviour modes, FUSE semantics through the +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 suites that attack the 9proc server itself (a hostile raw-9P client with 181 checks, the core, the Linux layer) moved with it to @@ -399,7 +508,10 @@ 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 until it returns (but not past the child's exit). + 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. * 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 3a072ec..089c3ca 100644 --- a/9ns/src/bridge.zig +++ b/9ns/src/bridge.zig @@ -4,6 +4,12 @@ //! Everything here is single-threaded and one request at a time. State is three //! tables: inodes (nodeid → fid/qid, deduplicated by qid.path), open handles //! (fh → fid plus a cached directory listing), and the reverse qid map. +//! +//! One request at a time does not mean deaf: while a 9P reply is outstanding +//! the session polls the FUSE descriptor too (`nine.Interrupt`). A +//! FUSE_INTERRUPT for the request being served becomes a Tflush, and if the +//! server honours it the request fails with EINTR; anything else the kernel +//! sends meanwhile is parked in a one-slot stash and served next. const std = @import("std"); const cloud9 = @import("cloud9"); const fuse = @import("fuse.zig"); @@ -89,12 +95,22 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_ defer b.deinit(); b.req_buf = try gpa.alignedAlloc(u8, .@"8", request_buf_len); + b.spare_buf = try gpa.alignedAlloc(u8, .@"8", request_buf_len); b.data_buf = try gpa.alloc(u8, max_write); + // Every read of the FUSE fd follows a poll; non-blocking makes sure a + // request the kernel withdrew in between cannot park us in read(2) while + // a 9P reply is due. + fuse.setNonblocking(fuse_fd) catch return error.FuseIo; + // Abandon any pending 9P reply once the child is gone (stop_fd readable), // including the initial root stat below: a silent server must not pin us. session.stop_fd = stop_fd; defer session.stop_fd = -1; + // And watch the FUSE fd meanwhile: INTERRUPTs become Tflush, other + // requests (INIT arrives during the root stat) wait in the stash. + session.interrupt = b.interruptSource(); + defer session.interrupt = null; // Node 1 is the root; its qid comes from a stat so lookups resolving back to // it (e.g. via a walk) dedupe onto node 1. @@ -115,6 +131,17 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_ .{ .fd = stop_fd, .events = linux.POLL.IN, .revents = 0 }, }; while (true) { + // A request that arrived while a 9P reply was outstanding goes first. + // It lives in the spare buffer; swap so that the spare is free again + // for anything that arrives while this one is being served. + if (b.stash) |req| { + b.stash = null; + std.mem.swap([]align(8) u8, &b.req_buf, &b.spare_buf); + if (!try b.dispatch(req)) return; + continue; + } + if (b.fuse_gone) return; + if (b.fuse_fail) |e| return e; pfds[0].revents = 0; pfds[1].revents = 0; const rc = linux.poll(&pfds, pfds.len, -1); @@ -128,7 +155,8 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, session: *nine.Session, root_ return; } if (pfds[0].revents == 0) continue; - const req = (fuse.readRequest(fuse_fd, b.req_buf) catch |e| switch (e) { + const req = (fuse.readRequestOnce(fuse_fd, b.req_buf) catch |e| switch (e) { + error.Retry => continue, error.Protocol => return error.FuseProtocol, else => return error.FuseIo, }) orelse { @@ -145,7 +173,23 @@ const Bridge = struct { nine: *nine.Session, opts: Options, req_buf: []align(8) u8 = &.{}, + /// Second request buffer: what the interrupt poll reads into. Holds the + /// stashed request until `serve` swaps it in. + spare_buf: []align(8) u8 = &.{}, data_buf: []u8 = &.{}, + /// A non-INTERRUPT request read while a 9P reply was outstanding (its body + /// 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, + /// 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, + /// The FUSE fd reported ENODEV / a failure while the session was waiting. + fuse_gone: bool = false, + fuse_fail: ?error{ FuseIo, FuseProtocol } = null, inodes: std.AutoHashMapUnmanaged(u64, Inode) = .empty, by_qid: std.AutoHashMapUnmanaged(u64, u64) = .empty, handles: std.AutoHashMapUnmanaged(u64, Handle) = .empty, @@ -165,9 +209,65 @@ const Bridge = struct { b.inodes.deinit(b.gpa); b.by_qid.deinit(b.gpa); if (b.req_buf.len != 0) b.gpa.free(b.req_buf); + if (b.spare_buf.len != 0) b.gpa.free(b.spare_buf); if (b.data_buf.len != 0) b.gpa.free(b.data_buf); } + // -- interrupt source (polled by nine.Session while a reply is outstanding) ----- + + fn interruptSource(b: *Bridge) nine.Interrupt { + return .{ .ctx = b, .watch = interruptWatch, .onReadable = interruptReadable, .armed = interruptArmed }; + } + + /// Poll the FUSE fd only while the stash has room: with it full a second + /// request would have nowhere to go. + fn interruptWatch(ctx: *anyopaque) i32 { + const b: *Bridge = @ptrCast(@alignCast(ctx)); + return if (b.stash == null and !b.fuse_gone and b.fuse_fail == null) b.fuse_fd else -1; + } + + /// Reads the request the kernel has ready. An INTERRUPT for the request in + /// flight asks the session to flush it; one for any other request is + /// dropped (the kernel expects no reply); anything else is stashed. + fn interruptReadable(ctx: *anyopaque) nine.Session.Error!bool { + const b: *Bridge = @ptrCast(@alignCast(ctx)); + const req = (fuse.readRequestOnce(b.fuse_fd, b.spare_buf) catch |e| switch (e) { + error.Retry => return false, + error.Protocol => { + b.fuse_fail = error.FuseProtocol; + return error.Stopped; + }, + else => { + b.fuse_fail = error.FuseIo; + return error.Stopped; + }, + }) orelse { + b.trace("fuse fd reports ENODEV while a 9P reply is outstanding", .{}); + b.fuse_gone = true; + return error.Stopped; + }; + 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) { + b.trace("<- interrupt for unique={d} (in flight): sending Tflush", .{in.unique}); + b.interrupted = true; + return true; + } + b.trace("<- interrupt for unique={d} (not in flight; ignored)", .{in.unique}); + return false; + } + b.trace("<- {s} unique={d} nodeid={d} stashed while a 9P reply is outstanding", .{ opName(h.op()), h.unique, h.nodeid }); + b.stash = req; + return false; + } + + fn interruptArmed(ctx: *anyopaque) bool { + const b: *Bridge = @ptrCast(@alignCast(ctx)); + return b.interrupted; + } + fn trace(b: *const Bridge, comptime fmt: []const u8, args: anytype) void { if (b.opts.debug) std.debug.print("9ns: " ++ fmt ++ "\n", args); } @@ -188,7 +288,12 @@ const Bridge = struct { b.reply(h.unique, &.{}) catch {}; return false; } + b.cur_unique = h.unique; + b.interrupted = false; + defer b.cur_unique = null; 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.?; const code: linux.E = switch (e) { error.Nine => b.last_err, error.BadRequest => .INVAL, @@ -201,6 +306,7 @@ const Bridge = struct { error.TooLarge => .NAMETOOLONG, error.BadDir => .IO, error.Closed, error.Protocol, error.Io, error.Stopped => .IO, + error.Interrupted => .INTR, error.FuseIo => return error.FuseIo, }; if (wants_reply) try b.replyError(h.unique, code); @@ -565,15 +671,15 @@ const Bridge = struct { while (true) { // A server that ignores the offset would otherwise feed us forever. if (offset >= max_dir_bytes) return error.BadDir; + // A flushed read whose reply still won the race: the listing is + // incomplete either way, so stop here rather than read on. + if (b.interrupted) return error.Interrupted; const n = try b.read(fid, offset, b.data_buf); if (n == 0) break; try parseDirRecords(b.gpa, b.data_buf[0..n], &list); offset += n; } - // Entries carrying the root's own qid.path get the root's ino (1), as GETATTR would report it. - for (list.entries.items[2..]) |*e| if (e.ino == b.root_path) { - e.ino = fuse.root_id; - }; + for (list.entries.items[2..]) |*e| e.ino = inoFromPath(e.ino); return list; } @@ -586,21 +692,29 @@ const Bridge = struct { } fn walkName(b: *Bridge, fid: u32, name: []const u8) nine.Session.Error!u32 { - const newfid = b.nine.allocFid(); - _ = b.nine.walk(fid, newfid, &.{name}) catch |e| { - b.nine.freeFid(newfid); - return b.nineErr("walk", fid, e); - }; + const newfid = try b.walkTo(fid, &.{name}); b.trace(" 9p walk fid={d} newfid={d} name={s} -> ok", .{ fid, newfid, name }); return newfid; } fn clone(b: *Bridge, fid: u32) nine.Session.Error!u32 { - const newfid = b.nine.clone(fid) catch |e| return b.nineErr("clone", fid, e); + const newfid = try b.walkTo(fid, &.{}); b.trace(" 9p walk fid={d} newfid={d} (clone) -> ok", .{ fid, newfid }); return newfid; } + /// allocFid + walk. On Rerror the new fid was never bound; after an + /// interruption the server may or may not have bound it (the Rflush + /// tells us only that no reply is coming), so it is clunked to be sure. + fn walkTo(b: *Bridge, fid: u32, names: []const []const u8) nine.Session.Error!u32 { + const newfid = b.nine.allocFid(); + _ = b.nine.walk(fid, newfid, names) catch |e| { + if (e == error.Interrupted) b.clunkQuiet(newfid) else b.nine.freeFid(newfid); + return b.nineErr("walk", fid, e); + }; + return newfid; + } + fn open9(b: *Bridge, fid: u32, mode: u8) nine.Session.Error!nine.Session.Open { const o = b.nine.open(fid, mode) catch |e| return b.nineErr("open", fid, e); b.trace(" 9p open fid={d} mode={d} -> iounit={d}", .{ fid, mode, o.iounit }); @@ -674,12 +788,18 @@ const Bridge = struct { // -- attrs --------------------------------------------------------------------------- + /// The inode number reported to the kernel is the 9P qid.path, for the root + /// too: FUSE only needs the root's *nodeid* to be 1, and a server may hand + /// qid.path 1 to some other file (Pardes gives it to /self), which would + /// 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 { - return if (nodeid == fuse.root_id or qid.path == b.root_path) fuse.root_id else qid.path; + _ = b; + _ = nodeid; + return inoFromPath(qid.path); } fn inoOfNode(b: *const Bridge, nodeid: u64) u64 { - if (nodeid == fuse.root_id) return fuse.root_id; const ino = b.inodes.get(nodeid) orelse return nodeid; return b.inoOf(nodeid, ino.qid); } @@ -700,6 +820,11 @@ const Bridge = struct { // -- pure helpers (unit-tested) ------------------------------------------------------------ /// Attr from a 9P Stat: DMDIR → S_IFDIR else S_IFREG, low 9 permission bits kept. +/// qid.path → inode number; 0 becomes a sentinel (inode 0 reads as "invalid" to many tools). +pub fn inoFromPath(path: u64) u64 { + return if (path == 0) std.math.maxInt(u64) - 1 else path; +} + pub fn attrFromStat(st: cloud9.Stat, ino: u64, uid: u32, gid: u32) fuse.Attr { const ftype: u32 = if (st.mode & cloud9.dmdir != 0) fuse.S_IFDIR else fuse.S_IFREG; return .{ @@ -964,6 +1089,67 @@ test "dirent names the kernel would reject are dropped from listings" { try testing.expectEqualStrings("also", list.entries.items[1].name); } +test "interrupt source: INTERRUPT in flight → Tflush, other INTERRUPTs ignored, requests stashed" { + var pi = try nine.PipeInterrupt.init(0); // only its fake FUSE fd and inject() are used + defer pi.deinit(); + var fs: nine.FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush, .on_hang_inject = &pi }; + var pair = try fs.start(); + defer pair.close(); + const s = &pair.session; + var b: Bridge = .{ .gpa = testing.allocator, .fuse_fd = pi.read_end, .nine = s, .opts = .{ .uid = 0, .gid = 0 } }; + defer b.deinit(); + b.spare_buf = try testing.allocator.alignedAlloc(u8, .@"8", request_buf_len); + s.interrupt = b.interruptSource(); + defer s.interrupt = null; + try testing.expectEqual(pi.read_end, Bridge.interruptWatch(&b)); + + // 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; + pi.unique = 7; + try pi.inject(6); + try testing.expectError(error.Interrupted, b.read(1, 0, &buf)); + try testing.expect(b.interrupted); + try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst)); + try testing.expectEqual(fs.hung_tag.load(.seq_cst), fs.flush_oldtag.load(.seq_cst)); + 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; + try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40])); + try testing.expectEqual(@as(usize, 0), s.client.pending()); + + // A FORGET arriving during a wait is stashed, and the fd is then not watched. + var wire: [48]u8 = undefined; + const hdr = fuse.InHeader{ .len = 48, .opcode = @intFromEnum(fuse.Opcode.forget), .unique = 99, .nodeid = 5, .uid = 0, .gid = 0, .pid = 0, .total_extlen = 0, .padding = 0 }; + @memcpy(wire[0..40], std.mem.asBytes(&hdr)); + @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; + 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; + try testing.expectEqual(fuse.Opcode.forget, stashed.header.op()); + try testing.expectEqual(@as(u64, 5), stashed.header.nodeid); + try testing.expectEqual(@as(u64, 1), (try fuse.body(fuse.ForgetIn, stashed)).nlookup); + try testing.expectEqual(@as(i32, -1), Bridge.interruptWatch(&b)); + // With the stash full an INTERRUPT is not even looked at. + try pi.inject(9); + try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40])); + try testing.expect(!b.interrupted); + b.stash = null; + 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; + 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); + try testing.expect(b.stash == null); +} + test "DirList frees its names" { var list: DirList = .{}; try list.entries.append(testing.allocator, .{ .name = try testing.allocator.dupe(u8, "x"), .ino = 1, .dtype = fuse.DT_REG }); diff --git a/9ns/src/fuse.zig b/9ns/src/fuse.zig index 682b3d5..216e616 100644 --- a/9ns/src/fuse.zig +++ b/9ns/src/fuse.zig @@ -314,7 +314,7 @@ comptime { // Request / reply helpers // --------------------------------------------------------------------------- -pub const Error = error{ Protocol, Io, TooManyPayloads }; +pub const Error = error{ Protocol, Io, TooManyPayloads, Retry }; pub const Request = struct { header: InHeader, @@ -327,22 +327,42 @@ pub const Request = struct { /// `max_write + 4096` bytes and 8-byte aligned so `body()` can view it. pub fn readRequest(fd: i32, buf: []u8) Error!?Request { while (true) { - const rc = linux.read(fd, buf.ptr, buf.len); - switch (linux.errno(rc)) { - .SUCCESS => { - const n: usize = rc; - if (n < @sizeOf(InHeader)) return error.Protocol; - const header = std.mem.bytesToValue(InHeader, buf[0..@sizeOf(InHeader)]); - if (header.len != n) return error.Protocol; - return .{ .header = header, .body = buf[@sizeOf(InHeader)..n] }; - }, - .INTR, .AGAIN, .NOENT => continue, - .NODEV => return null, - else => return error.Io, - } + return readRequestOnce(fd, buf) catch |e| switch (e) { + error.Retry => continue, + else => return e, + }; + } +} + +/// One `read(2)` attempt: like `readRequest` but EINTR/EAGAIN/ENOENT surface as +/// `error.Retry` instead of being retried, so a caller that only reads after +/// `poll` (or on a non-blocking fd) never blocks in here. +pub fn readRequestOnce(fd: i32, buf: []u8) Error!?Request { + const rc = linux.read(fd, buf.ptr, buf.len); + switch (linux.errno(rc)) { + .SUCCESS => { + const n: usize = rc; + if (n < @sizeOf(InHeader)) return error.Protocol; + const header = std.mem.bytesToValue(InHeader, buf[0..@sizeOf(InHeader)]); + if (header.len != n) return error.Protocol; + return .{ .header = header, .body = buf[@sizeOf(InHeader)..n] }; + }, + .INTR, .AGAIN, .NOENT => return error.Retry, + .NODEV => return null, + else => return error.Io, } } +/// Sets O_NONBLOCK on `fd` so that a read after `poll` cannot block when the +/// kernel withdrew the request in between (a killed waiter, for instance). +pub fn setNonblocking(fd: i32) Error!void { + const cur = linux.fcntl(fd, linux.F.GETFL, 0); + if (linux.errno(cur) != .SUCCESS) return error.Io; + const nonblock: u32 = @bitCast(linux.O{ .NONBLOCK = true }); + const rc = linux.fcntl(fd, linux.F.SETFL, @as(usize, cur) | nonblock); + if (linux.errno(rc) != .SUCCESS) return error.Io; +} + /// Maximum number of payload slices a single `reply` can carry. pub const max_payloads = 7; @@ -650,4 +670,12 @@ test "readRequest parses one request from a pipe and rejects bad lengths" { std.mem.bytesAsValue(InHeader, bad[0..40]).len = 40; try testing.expectEqual(@as(usize, 48), linux.write(fds[1], &bad, bad.len)); try testing.expectError(error.Protocol, readRequest(fds[0], &buf)); + + // Non-blocking and empty: readRequestOnce reports Retry instead of waiting. + try setNonblocking(fds[0]); + try testing.expectError(error.Retry, readRequestOnce(fds[0], &buf)); + try testing.expectEqual(@as(usize, 48), linux.write(fds[1], &wire, wire.len)); + const again = (try readRequestOnce(fds[0], &buf)) orelse return error.Io; + try testing.expectEqual(@as(u64, 3), again.header.unique); + try testing.expectError(error.Retry, readRequestOnce(fds[0], &buf)); } diff --git a/9ns/src/main.zig b/9ns/src/main.zig index 26bd699..076aa42 100644 --- a/9ns/src/main.zig +++ b/9ns/src/main.zig @@ -21,7 +21,9 @@ const usage_text = \\ --fd N already-connected inherited descriptor \\ --spawn CMD run CMD (via /bin/sh -c) with a socketpair on its stdin/stdout \\Options: - \\ --mount PATH mountpoint inside the new namespace (default /mnt/9p) + \\ --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) \\ --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) @@ -30,9 +32,17 @@ const usage_text = \\ --debug trace FUSE and 9P operations on stderr \\ --help, --version \\PROGRAM defaults to $SHELL (else /bin/sh). The mountpoint is exported as $NINE_MOUNT. + \\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. \\ ; +/// Where `--name NAME` mounts: `mount_root/NAME`. +const mount_root = "/mnt/9p"; +/// Name used when nothing usable can be derived from the transport. +const fallback_name = "9p"; + const own_failure: u8 = 125; /// Largest 9P message size we agree to request: the session allocates two /// buffers of this size up front, before the server negotiates it down. @@ -55,7 +65,10 @@ fn printStdout(text: []const u8) void { const Config = struct { address: ?nine.Address = null, spawn_cmd: ?[]const u8 = null, - mount: []const u8 = "/mnt/9p", + /// `--mount`: wins over `name` when set. + mount: ?[]const u8 = null, + /// `--name`: null means "derive from the transport" (see `defaultName`). + name: ?[]const u8 = null, uname: ?[]const u8 = null, aname: []const u8 = "", msize: u32 = 131072, @@ -106,7 +119,7 @@ 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, mount, uname, aname, msize, cache, @"no-direct-io", debug, help, version, unknown }; + const Opt = enum { unix, tcp, fd, spawn, 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}), @@ -142,6 +155,10 @@ fn parseArgs(arena: std.mem.Allocator, args: []const [:0]const u8) !ParseResult cfg.spawn_cmd = value; 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; + }, .mount => { if (value.len == 0) return usageError("--mount wants a path", .{}); cfg.mount = value; @@ -183,6 +200,52 @@ fn parseTcp(spec: []const u8) ?nine.Address { return .{ .tcp = .{ .host = host, .port = port } }; } +/// A mount name is one path component: non-empty, no '/', no NUL, not `.` +/// or `..`. +fn validName(name: []const u8) bool { + if (name.len == 0) return false; + if (std.mem.eql(u8, name, ".") or std.mem.eql(u8, name, "..")) return false; + for (name) |c| if (c == '/' or c == 0) return false; + return true; +} + +/// The mount name derived from the transport when `--name` is absent: +/// `--unix PATH` → basename of PATH without a trailing `.sock`/`.9p`/ +/// `.socket`; `--tcp IP:PORT` → `tcp-IP-PORT` with every ':' turned into +/// '-' (so an IPv6 literal stays one component); `--spawn CMD` → basename +/// of CMD's first word; `--fd N` → `fdN`. Anything that does not come out +/// as a valid name (empty basename, `..`, ...) becomes `9p`. The result is +/// written into `buf` (at most `buf.len` bytes; longer inputs fall back). +fn defaultName(buf: []u8, cfg: Config) []const u8 { + const raw: []const u8 = blk: { + if (cfg.spawn_cmd) |cmd| { + var words = std.mem.tokenizeAny(u8, cmd, " \t\r\n"); + break :blk std.fs.path.basename(words.next() orelse ""); + } + switch (cfg.address orelse return fallback_name) { + .unix => |path| { + const base = std.fs.path.basename(path); + inline for (.{ ".sock", ".socket", ".9p" }) |ext| { + if (base.len > ext.len and std.mem.endsWith(u8, base, ext)) break :blk base[0 .. base.len - ext.len]; + } + break :blk base; + }, + .tcp => |t| { + const text = std.fmt.bufPrint(buf, "tcp-{s}-{d}", .{ t.host, t.port }) catch return fallback_name; + std.mem.replaceScalar(u8, text, ':', '-'); + return if (validName(text)) text else fallback_name; + }, + .fd => |fd| { + const text = std.fmt.bufPrint(buf, "fd{d}", .{fd}) catch return fallback_name; + return text; + }, + } + }; + if (!validName(raw) or raw.len > buf.len) return fallback_name; + @memcpy(buf[0..raw.len], raw); + return buf[0..raw.len]; +} + /// `--spawn`: run CMD under /bin/sh with one end of a socketpair as its /// stdin/stdout; the other end is the 9P transport. const Server = struct { pid: i32, fd: i32 }; @@ -298,8 +361,15 @@ 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"; - const mountpoint = ns.resolveMountpoint(gpa, cfg.mount) catch |err| { - std.debug.print("9ns: --mount {s}: {t}\n", .{ cfg.mount, err }); + // `--mount PATH` wins; otherwise `/mnt/9p/<name>` with `--name` or a + // name derived from the transport. + var name_buf: [512]u8 = undefined; + const mount_arg: []const u8 = cfg.mount orelse blk: { + const name = cfg.name orelse defaultName(&name_buf, cfg); + break :blk try std.fmt.allocPrint(arena, mount_root ++ "/{s}", .{name}); + }; + 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); @@ -405,7 +475,21 @@ test "parseArgs" { const r = try parseArgs(arena, &args); try std.testing.expectEqual(@as(i32, 3), r.run.address.?.fd); try std.testing.expectEqual(@as(usize, 0), r.run.program.len); - try std.testing.expectEqualStrings("/mnt/9p", r.run.mount); + try std.testing.expect(r.run.mount == null); + try std.testing.expect(r.run.name == null); + } + { + const named = [_][:0]const u8{ "9ns", "--fd", "3", "--name", "bar", "--mount=/x" }; + const r = try parseArgs(arena, &named); + try std.testing.expectEqualStrings("bar", r.run.name.?); + try std.testing.expectEqualStrings("/x", r.run.mount.?); + const eq = [_][:0]const u8{ "9ns", "--fd", "3", "--name=baz" }; + try std.testing.expectEqualStrings("baz", (try parseArgs(arena, &eq)).run.name.?); + // Invalid names: a path, empty, . and .. + for ([_][:0]const u8{ "a/b", "", ".", "..", "/" }) |bad| { + const args = [_][:0]const u8{ "9ns", "--fd", "3", "--name", bad }; + try std.testing.expectEqual(@as(u8, 125), (try parseArgs(arena, &args)).exit); + } } { // Two transports, no transport, unknown option, missing value: all 125. @@ -439,6 +523,43 @@ test "parseArgs" { } } +test "defaultName" { + var buf: [512]u8 = undefined; + const Case = struct { cfg: Config, want: []const u8 }; + const cases = [_]Case{ + .{ .cfg = .{ .address = .{ .unix = "/tmp/9debug.sock" } }, .want = "9debug" }, + .{ .cfg = .{ .address = .{ .unix = "/run/user/1000/acme" } }, .want = "acme" }, + .{ .cfg = .{ .address = .{ .unix = "ramfs.9p" } }, .want = "ramfs" }, + .{ .cfg = .{ .address = .{ .unix = "/x/y.socket" } }, .want = "y" }, + .{ .cfg = .{ .address = .{ .unix = "/x/.sock" } }, .want = ".sock" }, // the whole name, not empty + .{ .cfg = .{ .address = .{ .unix = "/x/y/" } }, .want = "y" }, + .{ .cfg = .{ .address = .{ .unix = "/" } }, .want = "9p" }, + .{ .cfg = .{ .address = .{ .unix = "/x/.." } }, .want = "9p" }, + .{ .cfg = .{ .address = .{ .tcp = .{ .host = "127.0.0.1", .port = 564 } } }, .want = "tcp-127.0.0.1-564" }, + .{ .cfg = .{ .address = .{ .tcp = .{ .host = "::1", .port = 9999 } } }, .want = "tcp---1-9999" }, + .{ .cfg = .{ .address = .{ .fd = 3 } }, .want = "fd3" }, + .{ .cfg = .{ .spawn_cmd = "/x/9proc-demo --stdio" }, .want = "9proc-demo" }, + .{ .cfg = .{ .spawn_cmd = " ramfs\t-s" }, .want = "ramfs" }, + .{ .cfg = .{ .spawn_cmd = " " }, .want = "9p" }, + .{ .cfg = .{}, .want = "9p" }, + }; + for (cases) |c| try std.testing.expectEqualStrings(c.want, defaultName(&buf, c.cfg)); +} + +test "validName" { + try std.testing.expect(validName("a")); + try std.testing.expect(validName("tcp-127.0.0.1-564")); + try std.testing.expect(validName("...")); + try std.testing.expect(!validName("")); + try std.testing.expect(!validName(".")); + try std.testing.expect(!validName("..")); + try std.testing.expect(!validName("a/b")); + try std.testing.expect(!validName("a\x00b")); +} + test { _ = ns; + _ = @import("nine.zig"); + _ = @import("bridge.zig"); + _ = @import("fuse.zig"); } diff --git a/9ns/src/nine.zig b/9ns/src/nine.zig index 70633e6..c89a343 100644 --- a/9ns/src/nine.zig +++ b/9ns/src/nine.zig @@ -29,8 +29,23 @@ pub const dontcare = cloud9.Stat{ .muid = "", }; +/// How the owner of a session (the FUSE bridge) gets a say while an rpc waits +/// for its reply. `watch` names a descriptor to poll alongside the socket, or +/// -1 to poll nothing extra right now; when it becomes readable `onReadable` +/// consumes whatever is there and returns true if the request in flight should +/// be cancelled with a Tflush. `armed` reports whether such a cancellation was +/// requested earlier for the operation in progress; the chunked read/write +/// loops stop between chunks when it is set (a reply that raced the flush still +/// leaves the caller wanting out). +pub const Interrupt = struct { + ctx: *anyopaque, + watch: *const fn (ctx: *anyopaque) i32, + onReadable: *const fn (ctx: *anyopaque) Session.Error!bool, + armed: *const fn (ctx: *anyopaque) bool, +}; + pub const Session = struct { - pub const Error = error{ Nine, Protocol, Io, Closed, Stopped, TooLarge, OutOfMemory }; + pub const Error = error{ Nine, Protocol, Io, Closed, Stopped, Interrupted, TooLarge, OutOfMemory }; pub const Walk = struct { nwqid: u16, wqid: [cloud9.max_welem]cloud9.Qid }; pub const Open = struct { qid: cloud9.Qid, iounit: u32 }; @@ -53,6 +68,9 @@ pub const Session = struct { /// readable (the bridge's "child exited" pipe) the pending rpc fails with /// `error.Stopped` instead of blocking on a server that never answers. stop_fd: i32 = -1, + /// Optional interrupt source (the bridge's FUSE descriptor) consulted while + /// a reply is outstanding; see `Interrupt`. + interrupt: ?Interrupt = null, /// Connect to `address`, then negotiate the protocol version. /// `msize` is the maximum message size to ask for (0 = the buffers' size). @@ -117,9 +135,17 @@ pub const Session = struct { } /// Generic RPC. Result slices borrow the input buffer until the next call. + /// + /// While the reply is outstanding the socket is polled together with + /// `stop_fd` (→ `error.Stopped`) and the interrupt source's descriptor. When + /// the latter asks for a cancellation a Tflush for the request's tag goes out + /// and the wait continues until either the original reply arrives (the flush + /// lost the race; the result is returned as if nothing happened and the + /// Rflush is swallowed by a later call) or the Rflush does (→ + /// `error.Interrupted`; the server has dropped the request). pub fn rpc(s: *Session, req: cloud9.Client.Request) Error!cloud9.Client.Result { s.ename_len = 0; - _ = s.client.submit(req) catch |e| switch (e) { + const tag = s.client.submit(req) catch |e| switch (e) { error.NoTags, error.Handshake, error.Dead => return error.Protocol, error.NoSpace, error.TooLarge => return error.TooLarge, error.BadRequest => { @@ -128,29 +154,52 @@ pub const Session = struct { }, }; try s.flush(); + var flush_tag: ?u16 = null; var tmp: [64 * 1024]u8 = undefined; while (true) { - if (s.client.take()) |done| { - switch (done.result) { - .fail => |ename| { - s.setEname(ename); - return error.Nine; - }, - else => return done.result, + while (s.client.take()) |done| { + if (done.tag == tag) { + switch (done.result) { + .fail => |ename| { + s.setEname(ename); + return error.Nine; + }, + else => return done.result, + } } + if (flush_tag != null and done.tag == flush_tag.?) return error.Interrupted; + // An Rflush for a flush whose original reply won the race in an + // earlier call: the client has released both tags; nothing to do. + if (done.op == .flush) continue; + return error.Protocol; } if (s.client.dead) return error.Protocol; // After take() returned null the previous frame is gone, so the free // space is at least what the pending frame still needs. const room = s.client.in.len - s.client.in_len; if (room == 0) return error.Protocol; - const n = try readSome(s.fd, s.stop_fd, tmp[0..@min(room, tmp.len)]); - if (n == 0) return error.Closed; - const pushed = s.client.push(tmp[0..n]); - if (pushed != n) return error.Protocol; + switch (try s.wait()) { + .socket => { + const n = try readSocket(s.fd, tmp[0..@min(room, tmp.len)]); + if (n == 0) return error.Closed; + const pushed = s.client.push(tmp[0..n]); + if (pushed != n) return error.Protocol; + }, + .cancel => if (flush_tag == null) { + flush_tag = s.client.submit(.{ .flush = .{ .oldtag = tag } }) catch return error.Protocol; + try s.flush(); + }, + } } } + /// True when the interrupt source has asked for the operation in progress to + /// stop; consulted between the chunks of a read or write. + pub fn interruptArmed(s: *const Session) bool { + const i = s.interrupt orelse return false; + return i.armed(i.ctx); + } + /// Walk `names` from `fid` to `newfid`. A partial walk leaves `newfid` unbound /// (9P semantics) and reports `error.Nine` with ename "file does not exist". pub fn walk(s: *Session, fid: u32, newfid: u32, names: []const []const u8) Error!Walk { @@ -245,6 +294,50 @@ pub const Session = struct { return chunkSize(s.client.maxWrite(), s.iounits.get(fid) orelse 0); } + const Ready = enum { socket, cancel }; + + /// Blocks until the socket is readable (`.socket`), the interrupt source + /// wants the request in flight cancelled (`.cancel`), or `stop_fd` fires + /// (`error.Stopped`). Anything the interrupt source consumes without asking + /// for a cancellation simply resumes the wait. + fn wait(s: *Session) Error!Ready { + while (true) { + var pfds: [3]linux.pollfd = undefined; + var n: usize = 0; + pfds[n] = .{ .fd = s.fd, .events = linux.POLL.IN, .revents = 0 }; + n += 1; + const stop_at: ?usize = if (s.stop_fd >= 0) n else null; + if (stop_at != null) { + pfds[n] = .{ .fd = s.stop_fd, .events = linux.POLL.IN, .revents = 0 }; + n += 1; + } + const ifd: i32 = if (s.interrupt) |i| i.watch(i.ctx) else -1; + const int_at: ?usize = if (ifd >= 0) n else null; + if (int_at != null) { + pfds[n] = .{ .fd = ifd, .events = linux.POLL.IN, .revents = 0 }; + n += 1; + } + if (n == 1) return .socket; + const prc = linux.poll(&pfds, @intCast(n), -1); + switch (linux.errno(prc)) { + .SUCCESS => {}, + .INTR, .AGAIN => continue, + else => return error.Io, + } + // A reply that is already there wins over everything else. + if (pfds[0].revents != 0) return .socket; + if (stop_at) |i| { + if (pfds[i].revents != 0) return error.Stopped; + } + if (int_at) |i| { + if (pfds[i].revents != 0) { + const src = s.interrupt.?; + if (try src.onReadable(src.ctx)) return .cancel; + } + } + } + } + /// Writes everything in the client's output buffer to the socket. fn flush(s: *Session) Error!void { while (s.client.output().len != 0) { @@ -280,8 +373,13 @@ fn readWith( if (max_chunk == 0) return error.Protocol; var done: usize = 0; while (done < buf.len) { + // Like read(2): an interruption after some data arrived is a short read. + if (done != 0 and s.interruptArmed()) break; const want: u32 = @intCast(@min(buf.len - done, max_chunk)); - const r = try rpcFn(s, .{ .read = .{ .fid = fid, .offset = offset + done, .count = want } }); + const r = rpcFn(s, .{ .read = .{ .fid = fid, .offset = offset + done, .count = want } }) catch |e| { + if (e == error.Interrupted and done != 0) break; + return e; + }; const data = r.read; @memcpy(buf[done..][0..data.len], data); done += data.len; @@ -302,29 +400,20 @@ fn writeWith( if (max_chunk == 0) return error.Protocol; var done: usize = 0; while (done < data.len) { + if (done != 0 and s.interruptArmed()) break; const want: usize = @min(data.len - done, max_chunk); - const r = try rpcFn(s, .{ .write = .{ .fid = fid, .offset = offset + done, .data = data[done..][0..want] } }); + const r = rpcFn(s, .{ .write = .{ .fid = fid, .offset = offset + done, .data = data[done..][0..want] } }) catch |e| { + if (e == error.Interrupted and done != 0) break; + return e; + }; done += r.write; if (r.write < want) break; } return done; } -fn readSome(fd: i32, stop_fd: i32, buf: []u8) Session.Error!usize { +fn readSocket(fd: i32, buf: []u8) Session.Error!usize { while (true) { - if (stop_fd >= 0) { - var pfds = [_]linux.pollfd{ - .{ .fd = fd, .events = linux.POLL.IN, .revents = 0 }, - .{ .fd = stop_fd, .events = linux.POLL.IN, .revents = 0 }, - }; - const prc = linux.poll(&pfds, pfds.len, -1); - switch (linux.errno(prc)) { - .SUCCESS => {}, - .INTR, .AGAIN => continue, - else => return error.Io, - } - if (pfds[1].revents != 0 and pfds[0].revents == 0) return error.Stopped; - } const rc = linux.read(fd, buf.ptr, buf.len); switch (linux.errno(rc)) { .SUCCESS => return rc, @@ -339,6 +428,9 @@ fn readSome(fd: i32, stop_fd: i32, buf: []u8) Session.Error!usize { pub fn enameToErrno(ename: []const u8) linux.E { const Rule = struct { needle: []const u8, err: linux.E }; const rules = [_]Rule{ + // A server that answers a flushed request with an error (Pardes says + // "Interrupted system call") should look like a flush to the caller. + .{ .needle = "interrupt", .err = .INTR }, .{ .needle = "not exist", .err = .NOENT }, .{ .needle = "not found", .err = .NOENT }, .{ .needle = "no such", .err = .NOENT }, @@ -461,6 +553,8 @@ test { } test "ename → errno mapping" { + try testing.expectEqual(linux.E.INTR, enameToErrno("Interrupted system call")); + try testing.expectEqual(linux.E.INTR, enameToErrno("read interrupted")); try testing.expectEqual(linux.E.NOENT, enameToErrno("file does not exist")); try testing.expectEqual(linux.E.NOENT, enameToErrno("No Such File")); try testing.expectEqual(linux.E.NOENT, enameToErrno("directory entry not found")); @@ -524,10 +618,19 @@ const FakeFile = struct { calls: usize = 0, max_count: u32 = 0, short_write_at: ?usize = null, + /// Fail the call with `error.Interrupted` once this many calls were made. + interrupt_at: ?usize = null, + /// Report an armed interrupt once this many calls were made. + armed_at: ?usize = null, scratch: [4096]u8 = undefined, + fn interruptArmed(f: *const FakeFile) bool { + return if (f.armed_at) |at| f.calls >= at else false; + } + fn rpc(f: *FakeFile, req: cloud9.Client.Request) Session.Error!cloud9.Client.Result { f.calls += 1; + if (f.interrupt_at) |at| if (f.calls > at) return error.Interrupted; switch (req) { .read => |r| { f.max_count = @max(f.max_count, r.count); @@ -583,20 +686,60 @@ test "write chunks and stops at a short write" { try testing.expectEqual(@as(usize, 2), f.calls); } +test "interrupted chunk loops: partial count if data moved, Interrupted otherwise" { + var buf: [4000]u8 = undefined; + // The second chunk's rpc is interrupted: the first chunk is returned. + var f: FakeFile = .{ .len = 2500, .interrupt_at = 1 }; + try testing.expectEqual(@as(usize, 1000), try readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000)); + // The first chunk's rpc is interrupted: nothing was transferred. + f = .{ .len = 2500, .interrupt_at = 0 }; + try testing.expectError(error.Interrupted, readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000)); + // A reply that raced the flush arms the interrupt: stop before the next chunk. + f = .{ .len = 2500, .armed_at = 2 }; + try testing.expectEqual(@as(usize, 2000), try readWith(&f, FakeFile.rpc, 7, 0, &buf, 1000)); + try testing.expectEqual(@as(usize, 2), f.calls); + // Same for writes. + var data: [2500]u8 = undefined; + for (&data, 0..) |*b, i| b.* = @truncate(i); + f = .{ .len = 0, .interrupt_at = 2 }; + try testing.expectEqual(@as(usize, 2000), try writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000)); + f = .{ .len = 0, .interrupt_at = 0 }; + try testing.expectError(error.Interrupted, writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000)); + f = .{ .len = 0, .armed_at = 1 }; + try testing.expectEqual(@as(usize, 1000), try writeWith(&f, FakeFile.rpc, 7, 0, &data, 1000)); +} + // -- in-process server test --------------------------------------------------------- /// A tiny 9P2000 backend on a cloud9.Server: answers version/attach/walk/stat/open/ /// read/clunk/remove with canned data. Runs in its own thread over a socketpair. -const FakeServer = struct { +/// Test support only (bridge.zig's tests use it too). +pub const FakeServer = struct { fd: i32, msize: u32, max_read_count: u32 = 0, file_len: usize, + /// A Tread at this offset is never answered (a blocked stream read); the + /// server keeps serving whatever else arrives, notably a Tflush. + hang_offset: ?u64 = null, + /// What a Tflush gets: the connection dropped (`.hangup`), an Rflush, or + /// first the Rread the flush was aimed at and then the Rflush (the race). + on_flush: enum { hangup, rflush, reply_then_rflush } = .hangup, + /// Observed by the test thread: number of Tflush seen and the last oldtag. + flushes: std.atomic.Value(u32) = .init(0), + flush_oldtag: std.atomic.Value(u32) = .init(0xFFFF), + /// The tag of the hung Tread, for the test to compare with `flush_oldtag`. + hung_tag: std.atomic.Value(u32) = .init(0xFFFF), + /// When set, the server injects that source's INTERRUPT the moment a read + /// hangs, so the cancellation provably arrives while the wait is on. + on_hang_inject: ?*PipeInterrupt = null, + /// Delay before every Rread, so a test can be sure the client is waiting. + read_delay_ns: u64 = 0, - const file_qid: cloud9.Qid = .{ .type = 0, .version = 3, .path = 0x1234 }; - const dir_qid: cloud9.Qid = .{ .type = cloud9.qtdir, .version = 1, .path = 0x1 }; + pub const file_qid: cloud9.Qid = .{ .type = 0, .version = 3, .path = 0x1234 }; + pub const dir_qid: cloud9.Qid = .{ .type = cloud9.qtdir, .version = 1, .path = 0x1 }; - fn run(fs: *FakeServer) void { + pub fn run(fs: *FakeServer) void { fs.loop() catch |e| std.debug.print("fake server: {s}\n", .{@errorName(e)}); _ = linux.close(fs.fd); } @@ -610,6 +753,7 @@ const FakeServer = struct { var srv: cloud9.Server = .init(.{ .in = in, .out = out }); var tmp: [4096]u8 = undefined; var data: [8192]u8 = undefined; + var hung: ?struct { tag: u16, offset: u64, count: u32 } = null; while (true) { while (try srv.receive()) |req| { const tag = req.tag; @@ -649,18 +793,41 @@ const FakeServer = struct { .topen => |m| try srv.reply(tag, .{ .ropen = .{ .qid = file_qid, .iounit = if (m.mode == cloud9.owrite) 700 else 0 } }), .tread => |m| { fs.max_read_count = @max(fs.max_read_count, m.count); - var n: usize = 0; - if (m.offset < fs.file_len) n = @min(@as(usize, m.count), fs.file_len - @as(usize, @intCast(m.offset))); - n = @min(n, data.len); - for (data[0..n], 0..) |*b, i| b.* = @truncate(m.offset + i); - try srv.reply(tag, .{ .rread = .{ .data = data[0..n] } }); + if (fs.hang_offset != null and fs.hang_offset.? == m.offset) { + hung = .{ .tag = tag, .offset = m.offset, .count = m.count }; + fs.hung_tag.store(tag, .seq_cst); + if (fs.on_hang_inject) |p| try p.inject(p.unique); + } else { + if (fs.read_delay_ns != 0) { + const ts: linux.timespec = .{ .sec = @intCast(fs.read_delay_ns / std.time.ns_per_s), .nsec = @intCast(fs.read_delay_ns % std.time.ns_per_s) }; + _ = linux.nanosleep(&ts, null); + } + try srv.reply(tag, .{ .rread = .{ .data = fs.fill(&data, m.offset, m.count) } }); + } }, .twrite => |m| try srv.reply(tag, .{ .rwrite = .{ .count = @intCast(m.data.len) } }), .tclunk => try srv.reply(tag, .rclunk), .tremove => try srv.reply(tag, .{ .rerror = .{ .ename = "permission denied" } }), .twstat => try srv.reply(tag, .rwstat), - // A flush is the test's "hang up now" signal. - .tflush => return, + .tflush => |m| { + fs.flush_oldtag.store(m.oldtag, .seq_cst); + _ = fs.flushes.fetchAdd(1, .seq_cst); + switch (fs.on_flush) { + // The test's "hang up now" signal. + .hangup => return, + .rflush => { + if (hung != null and hung.?.tag == m.oldtag) hung = null; + try srv.reply(tag, .rflush); + }, + .reply_then_rflush => { + if (hung) |h| if (h.tag == m.oldtag) { + try srv.reply(h.tag, .{ .rread = .{ .data = fs.fill(&data, h.offset, h.count) } }); + hung = null; + }; + try srv.reply(tag, .rflush); + }, + } + }, else => try srv.reply(tag, .{ .rerror = .{ .ename = "not supported" } }), } srv.release(); @@ -677,8 +844,199 @@ const FakeServer = struct { if (srv.push(tmp[0..rc]) != rc) return error.Overflow; } } + + /// File contents: byte i == i & 0xff, `file_len` bytes long. + fn fill(fs: *const FakeServer, data: []u8, offset: u64, count: u32) []const u8 { + var n: usize = 0; + if (offset < fs.file_len) n = @min(@as(usize, count), fs.file_len - @as(usize, @intCast(offset))); + n = @min(n, data.len); + for (data[0..n], 0..) |*b, i| b.* = @truncate(offset + i); + return data[0..n]; + } + + /// A connected session (fid 0 attached, fid 1 walked to "file" and opened + /// for reading) plus the server thread; `close` when done. + pub const Pair = struct { + server: *FakeServer, + session: Session, + thread: std.Thread, + + pub fn close(p: *Pair) void { + p.session.deinit(); + p.thread.join(); + } + }; + + pub fn start(fs: *FakeServer) !Pair { + var fds: [2]i32 = undefined; + if (linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &fds)) != .SUCCESS) return error.Io; + fs.fd = fds[1]; + const th = try std.Thread.spawn(.{}, FakeServer.run, .{fs}); + var s = try Session.connect(testing.allocator, .{ .fd = fds[0] }, fs.msize); + errdefer s.deinit(); + _ = try s.attach(0, "me", ""); + const fid = s.allocFid(); + _ = try s.walk(0, fid, &.{"file"}); + _ = try s.open(fid, cloud9.oread); + return .{ .server = fs, .session = s, .thread = th }; + } +}; + +/// A fake interrupt source for the rpc wait loop: a SOCK_SEQPACKET pair stands +/// in for the FUSE descriptor (one datagram per request, like /dev/fuse +/// delivers one request per read). Mirrors the bridge's rules: an INTERRUPT for +/// `unique` arms and cancels, anything else is consumed and ignored. +pub const PipeInterrupt = struct { + read_end: i32, + write_end: i32, + unique: u64, + armed_flag: bool = false, + /// Requests consumed that were not the matching INTERRUPT. + ignored: usize = 0, + + const fuse_interrupt_opcode: u32 = 36; + + pub fn init(unique: u64) !PipeInterrupt { + var fds: [2]i32 = undefined; + const flags = linux.SOCK.SEQPACKET | linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK; + if (linux.errno(linux.socketpair(linux.AF.UNIX, flags, 0, &fds)) != .SUCCESS) return error.Io; + return .{ .read_end = fds[0], .write_end = fds[1], .unique = unique }; + } + + pub fn deinit(p: *PipeInterrupt) void { + _ = linux.close(p.read_end); + _ = linux.close(p.write_end); + } + + pub fn interface(p: *PipeInterrupt) Interrupt { + return .{ .ctx = p, .watch = watch, .onReadable = onReadable, .armed = armed }; + } + + /// Writes a FUSE_INTERRUPT request (InHeader + InterruptIn) naming `target`. + pub fn inject(p: *PipeInterrupt, target: u64) !void { + var wire: [48]u8 = undefined; + std.mem.writeInt(u32, wire[0..4], 48, .little); // len + std.mem.writeInt(u32, wire[4..8], fuse_interrupt_opcode, .little); // opcode + std.mem.writeInt(u64, wire[8..16], 0x8000_0000_0000_0001, .little); // the interrupt's own unique + @memset(wire[16..40], 0); // nodeid, uid, gid, pid, extlen, padding + std.mem.writeInt(u64, wire[40..48], target, .little); // InterruptIn.unique + if (linux.write(p.write_end, &wire, wire.len) != wire.len) return error.Io; + } + + fn watch(ctx: *anyopaque) i32 { + const p: *PipeInterrupt = @ptrCast(@alignCast(ctx)); + return p.read_end; + } + + fn onReadable(ctx: *anyopaque) Session.Error!bool { + const p: *PipeInterrupt = @ptrCast(@alignCast(ctx)); + var buf: [4096]u8 = undefined; + const rc = linux.read(p.read_end, &buf, buf.len); + if (linux.errno(rc) != .SUCCESS or rc != 48) return error.Io; + const opcode = std.mem.readInt(u32, buf[4..8], .little); + const target = std.mem.readInt(u64, buf[40..48], .little); + if (opcode == fuse_interrupt_opcode and target == p.unique) { + p.armed_flag = true; + return true; + } + p.ignored += 1; + return false; + } + + fn armed(ctx: *anyopaque) bool { + const p: *PipeInterrupt = @ptrCast(@alignCast(ctx)); + return p.armed_flag; + } }; +test "rpc wait loop: INTERRUPT → Tflush → Rflush → error.Interrupted; session still usable" { + var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush }; + var pair = try fs.start(); + defer pair.close(); + const s = &pair.session; + var pi = try PipeInterrupt.init(77); + defer pi.deinit(); + s.interrupt = pi.interface(); + defer s.interrupt = null; + + // An INTERRUPT for some other request is consumed and ignored: the wait + // goes on, and the reply (a read past the hang offset) arrives normally. + var buf: [100]u8 = undefined; + try pi.inject(78); + try testing.expectEqual(@as(usize, 40), try s.read(1, 10, &buf)); + try testing.expect(!pi.armed_flag); + + // The read at offset 0 hangs; the INTERRUPT for our request cancels it + // (the stray one above is consumed along the way if the reply beat it). + try pi.inject(77); + try testing.expectError(error.Interrupted, s.read(1, 0, &buf)); + try testing.expect(pi.armed_flag); + try testing.expectEqual(@as(usize, 1), pi.ignored); + try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst)); + try testing.expectEqual(fs.hung_tag.load(.seq_cst), fs.flush_oldtag.load(.seq_cst)); + + // Both tags are free again: further rpcs work. + const st = try s.stat(1); + try testing.expectEqual(@as(u64, 50), st.length); + try testing.expectEqual(@as(usize, 40), try s.read(1, 10, &buf)); + try testing.expectEqual(@as(usize, 0), s.client.pending()); +} + +test "rpc wait loop: the reply beats the Rflush → data returned, stray Rflush swallowed" { + var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .reply_then_rflush }; + var pair = try fs.start(); + defer pair.close(); + const s = &pair.session; + var pi = try PipeInterrupt.init(5); + defer pi.deinit(); + s.interrupt = pi.interface(); + defer s.interrupt = null; + + var buf: [40]u8 = undefined; + try pi.inject(5); + try testing.expectEqual(@as(usize, 30), try s.read(1, 0, buf[0..30])); + for (buf[0..30], 0..) |b, i| try testing.expectEqual(@as(u8, @truncate(i)), b); + try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst)); + try testing.expectEqual(fs.hung_tag.load(.seq_cst), fs.flush_oldtag.load(.seq_cst)); + + // The Rflush is still in flight (or already buffered): the next rpcs must + // step over it, and afterwards nothing is pending in the client. + pi.armed_flag = false; + const st = try s.stat(1); + try testing.expectEqual(@as(u64, 50), st.length); + try testing.expectEqual(@as(usize, 40), try s.read(1, 10, &buf)); + try testing.expectEqual(@as(usize, 0), s.client.pending()); +} + +test "rpc wait loop: a read interrupted after some data is a short read" { + // Chunks of maxRead = 1013 (msize 1024); the third chunk (offset 2026) hangs + // and the server fires the INTERRUPT at that moment. + var pi = try PipeInterrupt.init(9); + defer pi.deinit(); + var fs: FakeServer = .{ .fd = -1, .msize = 1024, .file_len = 5000, .hang_offset = 2026, .on_flush = .rflush, .on_hang_inject = &pi }; + var pair = try fs.start(); + defer pair.close(); + const s = &pair.session; + s.interrupt = pi.interface(); + defer s.interrupt = null; + + var buf: [4000]u8 = undefined; + try testing.expectEqual(@as(usize, 2026), try s.read(1, 0, &buf)); + for (buf[0..2026], 0..) |b, i| try testing.expectEqual(@as(u8, @truncate(i)), b); + try testing.expect(pi.armed_flag); + try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst)); + try testing.expectEqual(@as(usize, 0), s.client.pending()); + + // An INTERRUPT already waiting when the read starts: the first chunk's + // reply races the flush and wins, the armed flag then stops the loop. + pi.armed_flag = false; + try pi.inject(9); + try testing.expectEqual(@as(usize, 1013), try s.read(1, 0, &buf)); + // The stray Rflush is consumed by the next call. + _ = try s.stat(1); + try testing.expectEqual(@as(usize, 0), s.client.pending()); +} + 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/src/ns.zig b/9ns/src/ns.zig index 6da2c7f..b4342e5 100644 --- a/9ns/src/ns.zig +++ b/9ns/src/ns.zig @@ -198,25 +198,47 @@ fn buildArgv(gpa: Allocator, argv: []const []const u8) ![:null]?[*:0]const u8 { /// Make sure `path` is a directory, inside the *current* mount namespace: /// /// * already a directory → done; -/// * else `mkdir`; on `EACCES`/`EPERM`/`EROFS` shadow the parent directory -/// with a tmpfs that re-exposes every existing entry (bind mounts for -/// directories and files, recreated symlinks) and `mkdir` inside it; +/// * else find the deepest existing ancestor and `mkdir` the missing +/// components under it one by one (`mkdir -p`); the first of them failing +/// with `EACCES`/`EPERM`/`EROFS` (the normal case for `/mnt/9p/<name>` as +/// a plain user) means **shadow that ancestor**: mount a `tmpfs` over it +/// that re-exposes every existing entry (bind mounts for directories and +/// files, recreated symlinks), then create the missing components inside; /// * anything else fails with the errno and a hint. /// +/// So `/mnt/9p/x` on a host without `/mnt/9p` shadows `/mnt` and creates +/// `9p/x`; with a root-owned `/mnt/9p` it shadows `/mnt/9p`; inside a 9ns +/// namespace, where `/mnt/9p` is ours, it just creates `x`. `/` and `/proc` +/// are never shadowed, nor a directory with more than `max_shadow_entries`. +/// /// Every failure prints `9ns: <step> <path>: E<errno>` to stderr before /// returning. Meant to be called in the child of `spawn` (or from a /// throwaway namespace: `unshare -Urm`). pub fn ensureMountpoint(gpa: Allocator, path: [:0]const u8) !void { - if (fileType(linux.AT.FDCWD, path, false)) |ft| { + if (try existingKind(path)) |ft| { if (ft == .dir) return; std.debug.print("9ns: mountpoint {s}: exists but is not a directory\n", .{path}); return error.Mountpoint; } - if (fileType(linux.AT.FDCWD, path, true) == .symlink) { - std.debug.print("9ns: mountpoint {s}: dangling symlink\n", .{path}); - return error.Mountpoint; + // Deepest existing ancestor: walk up until something is there. + var base: []const u8 = path; + while (true) { + base = std.fs.path.dirname(base) orelse "/"; + var base_buf: [path_max]u8 = undefined; + const base_z = std.fmt.bufPrintZ(&base_buf, "{s}", .{base}) catch { + std.debug.print("9ns: mountpoint {s}: path too long\n", .{path}); + return error.Mountpoint; + }; + const kind = try existingKind(base_z) orelse continue; + if (kind != .dir) { + std.debug.print("9ns: mountpoint {s}: {s} is not a directory\n", .{ path, base }); + return error.Mountpoint; + } + break; } - const mk = linux.errno(linux.mkdirat(linux.AT.FDCWD, path, 0o755)); + // Missing components, deepest ancestor first. + const missing = path[base.len..]; + const mk = mkdirComponents(path, base.len, missing); switch (mk) { .SUCCESS => return, .ACCES, .PERM, .ROFS => {}, @@ -225,21 +247,20 @@ pub fn ensureMountpoint(gpa: Allocator, path: [:0]const u8) !void { return error.Mountpoint; }, } - const parent = std.fs.path.dirname(path) orelse "/"; - if (std.mem.eql(u8, parent, "/") or isSameDirectory(parent, "/")) { + if (std.mem.eql(u8, base, "/") or isSameDirectory(base, "/")) { std.debug.print("9ns: mkdir {s}: E{t}; refusing to shadow / (pass --mount an existing directory)\n", .{ path, mk }); return error.Mountpoint; } // The shadow rebuilds entries from /proc/self/fd/<fd>/<name>; a tmpfs // over /proc (or a subtree of it) would take that away from itself. - if (std.mem.eql(u8, parent, "/proc") or std.mem.startsWith(u8, parent, "/proc/")) { - std.debug.print("9ns: mkdir {s}: E{t}; refusing to shadow {s} (pass --mount an existing directory)\n", .{ path, mk, parent }); + if (std.mem.eql(u8, base, "/proc") or std.mem.startsWith(u8, base, "/proc/")) { + std.debug.print("9ns: mkdir {s}: E{t}; refusing to shadow {s} (pass --mount an existing directory)\n", .{ path, mk, base }); return error.Mountpoint; } - const parent_z = try gpa.dupeZ(u8, parent); - defer gpa.free(parent_z); - try shadowDirectory(gpa, parent_z); - switch (linux.errno(linux.mkdirat(linux.AT.FDCWD, path, 0o755))) { + const base_z = try gpa.dupeZ(u8, base); + defer gpa.free(base_z); + try shadowDirectory(gpa, base_z); + switch (mkdirComponents(path, base.len, missing)) { .SUCCESS => {}, else => |e| { std.debug.print("9ns: mkdir {s} (in shadow tmpfs): E{t}\n", .{ path, e }); @@ -248,6 +269,36 @@ pub fn ensureMountpoint(gpa: Allocator, path: [:0]const u8) !void { } } +/// What `path` is, following symlinks: null when nothing is there; an +/// error (reported) for a dangling symlink. +fn existingKind(path: [*:0]const u8) !?FileType { + if (fileType(linux.AT.FDCWD, path, false)) |ft| return ft; + if (fileType(linux.AT.FDCWD, path, true) == .symlink) { + std.debug.print("9ns: mountpoint {s}: dangling symlink\n", .{std.mem.span(path)}); + return error.Mountpoint; + } + return null; +} + +/// `mkdir` each component of `missing` (which is `path[base_len..]`, so +/// it starts with '/') under the existing prefix `path[0..base_len]`, in +/// order. Returns the errno of the first failure (`.SUCCESS` when all were +/// created); `EEXIST` on a component is fine (another process, or a retry). +fn mkdirComponents(path: []const u8, base_len: usize, missing: []const u8) E { + var buf: [path_max]u8 = undefined; + var end: usize = base_len; + var it = std.mem.tokenizeScalar(u8, missing, '/'); + while (it.next()) |comp| { + end += 1 + comp.len; + const prefix = std.fmt.bufPrintZ(&buf, "{s}", .{path[0..end]}) catch return .NAMETOOLONG; + switch (linux.errno(linux.mkdirat(linux.AT.FDCWD, prefix, 0o755))) { + .SUCCESS, .EXIST => {}, + else => |e| return e, + } + } + return .SUCCESS; +} + const FileType = enum { dir, symlink, other }; /// True when both paths resolve (following symlinks, including magic ones diff --git a/9ns/test/adv_bridge_hostile.py b/9ns/test/adv_bridge_hostile.py index d353541..f676d2f 100755 --- a/9ns/test/adv_bridge_hostile.py +++ b/9ns/test/adv_bridge_hostile.py @@ -62,7 +62,8 @@ name_huge / has an entry with a 60000-byte name name_dots / lists "." and ".." too rerror_big Twalk to nope: Rerror with 65535 bytes of text extra_reply Rread on /f: an unsolicited Rclunk (tag 9) precedes the real reply -never Tread on /f: never reply (hang) +never Tread on /f: never reply (hang); the server stops reading, so a Tflush is never seen +never_flush Tread on /f: never reply, but keep serving: log every Tflush and answer Rflush close_mid Tread on /f: close the socket without replying renegotiate Tread on /f: an unsolicited Rversion precedes the real reply length_max Tstat on /f: length = 2**64-1 @@ -142,6 +143,9 @@ class Server: parent.children[name] = n return n + def log(self, text): + print(text, flush=True) + # -- framing --------------------------------------------------------- def frame(self, typ, tag, body): return s32(7 + len(body)) + s8(typ) + s16(tag) + body @@ -196,6 +200,8 @@ class Server: self.fids[fid] = [self.root, False] return self.frame(Rattach, tag, self.root.qid(self)) if typ == Tflush: + oldtag = r.u16() + self.log('Tflush tag=%d oldtag=%d' % (tag, oldtag)) return self.frame(Rflush, tag, b'') if typ == Twalk: fid, newfid, nw = r.u32(), r.u32(), r.u16() @@ -292,6 +298,9 @@ class Server: if m == 'never': time.sleep(3600) return None + if m == 'never_flush': + self.log('Tread tag=%d fid=%d hang' % (tag, fid)) + return b'' if m == 'close_mid': return None if m == 'renegotiate': diff --git a/9ns/test/adv_bridge_hostile.sh b/9ns/test/adv_bridge_hostile.sh index f5ba871..20c05bf 100755 --- a/9ns/test/adv_bridge_hostile.sh +++ b/9ns/test/adv_bridge_hostile.sh @@ -37,7 +37,7 @@ start_server() { # mode run() { local mode=$1 script=$2; shift 2 start_server "$mode" - OUT=$(timeout 30 "$NS" --unix "$SOCK" "$@" -- sh -c "$script" 2>"$TMP/stderr") + OUT=$(timeout 30 "$NS" --unix "$SOCK" --mount "$M" "$@" -- sh -c "$script" 2>"$TMP/stderr") RC=$? STDERR=$(cat "$TMP/stderr") } @@ -146,7 +146,7 @@ echo "# a server that never replies" start_server never # SIGTERM is forwarded to the child; once the child is gone 9ns must leave the # pending 9P reply behind and exit even though the server stays silent. -timeout -s TERM 3 "$NS" --unix "$SOCK" -- sh -c "cat $M/f; echo unreachable" >"$TMP/never.out" 2>"$TMP/never.err" & +timeout -s TERM 3 "$NS" --unix "$SOCK" --mount "$M" -- sh -c "cat $M/f; echo unreachable" >"$TMP/never.out" 2>"$TMP/never.err" & TPID=$! sleep 4 if kill -0 "$TPID" 2>/dev/null; then @@ -156,7 +156,7 @@ else fi wait "$TPID" 2>/dev/null # Without a signal the mount hangs (documented v1 limitation) until the server dies. -timeout 30 "$NS" --unix "$SOCK" -- sh -c "cat $M/f; echo unreachable" >"$TMP/never.out" 2>"$TMP/never.err" & +timeout 30 "$NS" --unix "$SOCK" --mount "$M" -- sh -c "cat $M/f; echo unreachable" >"$TMP/never.out" 2>"$TMP/never.err" & TPID=$! sleep 1.5 if kill -0 "$TPID" 2>/dev/null; then diff --git a/9ns/test/adv_bridge_interrupt.sh b/9ns/test/adv_bridge_interrupt.sh new file mode 100755 index 0000000..e364143 --- /dev/null +++ b/9ns/test/adv_bridge_interrupt.sh @@ -0,0 +1,136 @@ +#!/usr/bin/env bash +# Interrupt tests: a reader blocked in a 9P read the server never answers must +# be releasable. Killing or Ctrl-C-ing it makes the kernel send FUSE_INTERRUPT, +# which the bridge turns into a Tflush; a server that answers Rflush (the +# hostile server's `never_flush` mode) unblocks the reader with EINTR and the +# mount stays usable; a server that ignores everything (`never`) still blocks +# the mount, but SIGTERM to 9ns ends the session as before. +# Usage: bash 9ns/test/adv_bridge_interrupt.sh <9ns> <9proc-demo> (9proc unused; part of zig build 9ns-adv) +set -u +NS=$(realpath "${1:?path to 9ns}") +HERE=$(cd "$(dirname "$0")" && pwd) +SRV=$HERE/adv_bridge_hostile.py +TMP=$(mktemp -d "${TMPDIR:-/tmp}/9ns-int.XXXXXX") +M=/mnt/9p +FAILED=0 +PASSED=0 +SRVPID= + +cleanup() { [ -n "$SRVPID" ] && kill "$SRVPID" 2>/dev/null; pkill -f "adv_bridge_hostile.py $TMP" 2>/dev/null; rm -rf "$TMP"; } +trap cleanup EXIT + +if ! unshare -Urm true 2>/dev/null || [ ! -c /dev/fuse ]; then echo "SKIP: no user namespaces or /dev/fuse"; 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() { if [ "$2" = "$3" ]; then pass "$1"; else fail "$1" "expected: $(printf %q "$2")" "actual: $(printf %q "$3")"; fi; } +expect_contains() { case "$3" in *"$2"*) pass "$1" ;; *) fail "$1" "missing: $(printf %q "$2")" "in: $(printf %q "$3")" ;; esac; } +# expect_lt NAME ACTUAL LIMIT +expect_lt() { if [ "$2" -lt "$3" ]; then pass "$1 ($2 < $3)"; else fail "$1" "expected < $3, got $2"; fi; } + +start_server() { # mode + [ -n "$SRVPID" ] && { kill "$SRVPID" 2>/dev/null; wait "$SRVPID" 2>/dev/null; } + SOCK=$TMP/$1.sock + SRVLOG=$TMP/$1.srv.out + rm -f "$SOCK" + python3 "$SRV" "$SOCK" "$1" >"$SRVLOG" 2>&1 </dev/null & + SRVPID=$! + for _ in $(seq 1 100); do [ -S "$SOCK" ] && return 0; sleep 0.02; done + echo "server for $1 did not start"; cat "$SRVLOG"; exit 1 +} + +# run MODE SCRIPT [extra 9ns args...]: starts the server, runs 9ns; sets OUT, RC, STDERR. +run() { + local mode=$1 script=$2; shift 2 + start_server "$mode" + OUT=$(timeout 60 "$NS" --unix "$SOCK" --mount "$M" "$@" -- sh -c "$script" 2>"$TMP/stderr") + RC=$? + STDERR=$(cat "$TMP/stderr") +} + +no_crash() { # name + if [ "$RC" -ge 128 ] || [ "$RC" -eq 124 ]; then fail "$1: 9ns exit $RC" "$STDERR"; return; fi + case "$STDERR" in *panic*|*"Segmentation"*|*"integer overflow"*|*"reached unreachable"*|*"index out of bounds"*) fail "$1: crash text in stderr" "$STDERR";; *) pass "$1: no crash (exit $RC)";; esac +} + +# The shell snippet that times a command: prints "rc=N" and "ms=N". +timed() { # command... + printf 's=$(date +%%s%%N); %s; echo rc=$?; e=$(date +%%s%%N); echo ms=$(( (e - s) / 1000000 ))' "$*" +} +field() { printf '%s\n' "$2" | sed -n "s/^$1=//p" | tail -1; } +# Server-side fid count before and after the interrupted operation, measured in +# the same run. Everything the scripts touch is looked up first, so the inode +# fids the kernel keeps for f, d and g are in both samples and the comparison +# is exact: any difference is a fid an interrupted operation left behind. +BEFORE="stat $M/f $M/d/g >/dev/null; ls $M >/dev/null; echo before=\$(cat $M/fids)" +AFTER="sleep 0.2; echo after=\$(cat $M/fids)" +fids_same() { # name + local b a; b=$(field before "$OUT"); a=$(field after "$OUT") + case "$b" in ''|*[!0-9]*) fail "$1: fid count before is not numeric: '$b'"; return;; esac + expect_eq "$1: server-side fid count unchanged ($b)" "$b" "$a" +} +# Same for the bridge's own counter (--debug prints fids=N per request): the +# first and the last READ traced are the two `cat fids`. +bridge_fids_same() { # name + local reads; reads=$(printf '%s\n' "$STDERR" | sed -n 's/^9ns: <- read .*(fids=\([0-9]*\) .*/\1/p') + expect_eq "$1: bridge fid counter unchanged ($(printf '%s\n' "$reads" | head -1))" "$(printf '%s\n' "$reads" | head -1)" "$(printf '%s\n' "$reads" | tail -1)" +} + +echo "# (a) SIGINT to a reader blocked in a read the server never answers" +run never_flush "$BEFORE; $(timed "timeout -s INT 2 cat $M/f"); ls $M | tr '\n' ' '; echo; cat $M/d/g; $AFTER" --debug +no_crash "never_flush/SIGINT" +expect_eq "never_flush/SIGINT: cat was killed by the signal (timeout reports 124)" "124" "$(field rc "$OUT")" +expect_lt "never_flush/SIGINT: the reader was released promptly" "$(field ms "$OUT")" 2500 +expect_contains "never_flush/SIGINT: the mount is still usable (ls)" "big d f fids" "$OUT" +expect_contains "never_flush/SIGINT: the mount is still usable (read another file)" "in d" "$OUT" +fids_same "never_flush/SIGINT" +hung_tag=$(sed -n 's/^Tread tag=\([0-9]*\) .*hang$/\1/p' "$SRVLOG" | head -1) +flush_oldtag=$(sed -n 's/^Tflush tag=[0-9]* oldtag=\([0-9]*\)$/\1/p' "$SRVLOG" | head -1) +expect_eq "never_flush/SIGINT: exactly one Tflush reached the server" "1" "$(grep -c '^Tflush ' "$SRVLOG")" +expect_eq "never_flush/SIGINT: Tflush.oldtag is the hung Tread's tag ($hung_tag)" "$hung_tag" "$flush_oldtag" +expect_contains "never_flush/SIGINT: --debug shows the interrupt being forwarded" "sending Tflush" "$STDERR" +expect_contains "never_flush/SIGINT: --debug shows EINTR going back to the kernel" "error EINTR" "$STDERR" +bridge_fids_same "never_flush/SIGINT" + +echo "# (c) SIGKILL to the blocked reader" +run never_flush "$BEFORE; $(timed "cat $M/f & p=\$!; sleep 0.5; kill -9 \$p; wait \$p"); cat $M/d/g; $AFTER" +no_crash "never_flush/SIGKILL" +expect_eq "never_flush/SIGKILL: reader died of SIGKILL (137)" "137" "$(field rc "$OUT")" +expect_lt "never_flush/SIGKILL: released promptly" "$(field ms "$OUT")" 2000 +expect_contains "never_flush/SIGKILL: mount still usable" "in d" "$OUT" +fids_same "never_flush/SIGKILL" +expect_eq "never_flush/SIGKILL: one Tflush" "1" "$(grep -c '^Tflush ' "$SRVLOG")" + +echo "# a second blocked read after the first was interrupted" +run never_flush "$BEFORE; $(timed "timeout -s INT 1 cat $M/f"); $(timed "timeout -s INT 1 cat $M/f"); cat $M/d/g; $AFTER" +no_crash "never_flush/twice" +expect_eq "never_flush/twice: both readers killed" "124 124" "$(printf '%s\n' "$OUT" | sed -n 's/^rc=//p' | tr '\n' ' ' | sed 's/ $//')" +expect_eq "never_flush/twice: two Tflush, no tag confusion" "2" "$(grep -c '^Tflush ' "$SRVLOG")" +expect_contains "never_flush/twice: mount still usable" "in d" "$OUT" +fids_same "never_flush/twice" + +echo "# an interrupted open+read through a lookup (walk+stat) leaves no fid behind" +# `f` is looked up fresh each time (cache 0), so the LOOKUP's walk+stat and the +# OPEN's clone+open all run before the read blocks; all their fids must go. +run never_flush "$BEFORE; $(timed "timeout -s INT 1 cat $M/f"); $AFTER" --cache 0 --debug +no_crash "never_flush/cache0" +expect_eq "never_flush/cache0: reader killed" "124" "$(field rc "$OUT")" +fids_same "never_flush/cache0" +bridge_fids_same "never_flush/cache0" + +echo "# (b) a server that ignores Tflush too: the reader stays blocked, SIGTERM to 9ns still ends the session" +start_server never +timeout -s TERM 3 "$NS" --unix "$SOCK" --mount "$M" -- sh -c "timeout -s INT 1 cat $M/f; echo unreachable-rc=\$?" >"$TMP/never.out" 2>"$TMP/never.err" & +TPID=$! +sleep 4 +if kill -0 "$TPID" 2>/dev/null; then + fail "never: SIGTERM did not end 9ns while the reader was stuck"; kill -9 "$TPID" +else + pass "never: SIGTERM ends 9ns even though the server ignores the Tflush" +fi +wait "$TPID" 2>/dev/null +expect_eq "never: the reader never came back (server ignores Tflush)" "" "$(grep unreachable "$TMP/never.out")" + +echo +echo "passed=$PASSED failed=$FAILED" +[ "$FAILED" -eq 0 ] diff --git a/9ns/test/adv_bridge_semantics.sh b/9ns/test/adv_bridge_semantics.sh index b38ce9c..526ebc8 100755 --- a/9ns/test/adv_bridge_semantics.sh +++ b/9ns/test/adv_bridge_semantics.sh @@ -24,7 +24,7 @@ SOCK=$TMP/i.sock SRVPID=$! for _ in $(seq 1 100); do [ -S "$SOCK" ] && break; sleep 0.02; done # Each run is a fresh session; state persists in the server, so tests clean up after themselves. -run() { OUT=$(timeout 120 "$NS" --unix "$SOCK" "${EXTRA[@]}" -- sh -c "$1" 2>"$TMP/stderr"); RC=$?; STDERR=$(cat "$TMP/stderr"); } +run() { OUT=$(timeout 120 "$NS" --unix "$SOCK" --mount "$M" "${EXTRA[@]}" -- sh -c "$1" 2>"$TMP/stderr"); RC=$?; STDERR=$(cat "$TMP/stderr"); } EXTRA=() py() { run "python3 - <<'PYEOF' $1 @@ -141,8 +141,9 @@ for n in os.listdir(d): os.unlink(f'{d}/{n}') os.rmdir(d) " expect_eq "readdir of a directory that changes mid-iteration" "a True True" "$OUT" -run "ls $M/.. > /dev/null && echo ok; stat -c %i $M $M/. $M/scratch/..; cd $M/scratch && ls .. | grep -c scratch" -expect_eq ".. of the root and of a subdir" $'ok\n1\n1\n1\n1' "$OUT" +# The root's inode number is its 9P qid.path (not a fixed 1): all three views must agree. +run "ls $M/.. > /dev/null && echo ok; stat -c %i $M $M/. $M/scratch/.. | sort -u | wc -l; cd $M/scratch && ls .. | grep -c scratch" +expect_eq ".. of the root and of a subdir" $'ok\n1\n1' "$OUT" run "cd $M && find . -type d | wc -l && find . -type f | head -1 && find $S -type f | wc -l" expect_contains "find -type works" "./README" "$OUT" run "stat -f -c '%T %S %l' $M; df -P $M | tail -1 | awk '{print \$1}'; sync -f $M && echo synced; sync && echo synced2" diff --git a/9ns/test/adv_bridge_stress.sh b/9ns/test/adv_bridge_stress.sh index b306b28..d49f453 100755 --- a/9ns/test/adv_bridge_stress.sh +++ b/9ns/test/adv_bridge_stress.sh @@ -21,7 +21,7 @@ SOCK=$TMP/i.sock "$PROC" --unix "$SOCK" >"$TMP/srv.out" 2>&1 & SRVPID=$! for _ in $(seq 1 100); do [ -S "$SOCK" ] && break; sleep 0.02; done -run() { OUT=$(timeout 600 "$NS" --unix "$SOCK" -- sh -c "$1" 2>"$TMP/stderr"); RC=$?; STDERR=$(cat "$TMP/stderr"); } +run() { OUT=$(timeout 600 "$NS" --unix "$SOCK" --mount "$M" -- sh -c "$1" 2>"$TMP/stderr"); RC=$?; STDERR=$(cat "$TMP/stderr"); } echo "# 100k+ 9P operations in one session; RSS must plateau" # Each iteration: create+write+close, open+read+close, unlink, plus a failing lookup: ~15 RPCs. diff --git a/9ns/test/adversarial.sh b/9ns/test/adversarial.sh index 71c6f44..8575dcf 100755 --- a/9ns/test/adversarial.sh +++ b/9ns/test/adversarial.sh @@ -8,7 +8,7 @@ NS=${1:?path to 9ns} PROC=${2:?path to 9proc-demo} HERE=$(cd "$(dirname "$0")" && pwd) status=0 -for suite in adv_ns_process adv_bridge_hostile adv_bridge_semantics adv_bridge_stress; do +for suite in adv_ns_process adv_bridge_hostile adv_bridge_interrupt adv_bridge_semantics adv_bridge_stress; do echo "### $suite" if bash "$HERE/$suite.sh" "$NS" "$PROC"; then echo "### $suite: ok"; else echo "### $suite: FAILED"; status=1; fi done diff --git a/9ns/test/integration.sh b/9ns/test/integration.sh index 5e0ad13..32d132a 100755 --- a/9ns/test/integration.sh +++ b/9ns/test/integration.sh @@ -10,6 +10,8 @@ TMP=$(mktemp -d "${TMPDIR:-/tmp}/9ns-itest.XXXXXX") PIDS=() FAILED=0 PASSED=0 +# The suites pin the mountpoint with --mount; the default is /mnt/9p/<name> +# (see the "# mount names" section below). M=/mnt/9p cleanup() { @@ -39,8 +41,8 @@ wait_socket() { # path return 1 } -# run_in "<shell script>" — run inside a namespace with the current transport ($TRANSPORT array). -run_in() { timeout 60 "$NS" "${TRANSPORT[@]}" -- sh -c "$1" 2>"$TMP/stderr"; } +# run_in "<shell script>" — run inside a namespace with the current transport ($TRANSPORT array), mounted on $M. +run_in() { timeout 60 "$NS" "${TRANSPORT[@]}" --mount "$M" -- sh -c "$1" 2>"$TMP/stderr"; } # --- scratch battery: works against any writable 9P tree rooted at $1 (relative to mount) --- scratch_battery() { # label scratchdir @@ -111,6 +113,26 @@ mkdir -p "$TMP/mnt" expect_eq "--mount existing dir" "ok" "$(timeout 60 "$NS" --unix "$SOCK" --mount "$TMP/mnt" -- sh -c "[ -f $TMP/mnt/README ] && echo ok")" expect_eq "--mount relative" "ok" "$(cd "$TMP" && timeout 60 "$NS" --unix "$SOCK" --mount rel -- sh -c "[ -f $TMP/rel/README ] && echo ok")" expect_eq "--mount missing under /" "125" "$(timeout 60 "$NS" --unix "$SOCK" --mount /nonexistent-9ns-dir -- true 2>/dev/null; echo $?)" +expect_eq "--mount beats --name" "$TMP/mnt" "$(timeout 60 "$NS" --unix "$SOCK" --name zz --mount "$TMP/mnt" -- sh -c 'echo $NINE_MOUNT')" + +echo "# mount names (default mountpoint /mnt/9p/<name>)" +mkdir -p "$TMP/n" +"$PROC" --unix "$TMP/n/foo.sock" & +PIDS+=($!) +wait_socket "$TMP/n/foo.sock" || echo "foo.sock missing" +named() { timeout 60 "$NS" --unix "$TMP/n/foo.sock" "$@"; } +expect_eq "unix socket foo.sock -> /mnt/9p/foo" "/mnt/9p/foo" "$(named -- sh -c 'echo $NINE_MOUNT')" +expect_eq "default mount is served" "ok" "$(named -- sh -c '[ -f /mnt/9p/foo/README ] && grep -q "^9ns /mnt/9p/foo fuse" /proc/self/mounts && echo ok')" +expect_eq "default mount keeps /mnt entries" "$(ls -A /mnt | sort | tr '\n' ' ')" "$(named -- sh -c "ls -A /mnt | grep -v '^9p\$' | sort | tr '\n' ' '")" +expect_eq "--name bar -> /mnt/9p/bar" "/mnt/9p/bar ok" "$(named --name bar -- sh -c 'echo $NINE_MOUNT; [ -f /mnt/9p/bar/README ] && echo ok' | tr '\n' ' ' | sed 's/ $//')" +expect_eq "--name=bar form" "/mnt/9p/bar" "$(named --name=bar -- sh -c 'echo $NINE_MOUNT')" +expect_eq "--name a/b is a usage error" "125" "$(named --name a/b -- true 2>/dev/null; echo $?)" +expect_eq "--name . is a usage error" "125" "$(named --name . -- true 2>/dev/null; echo $?)" +expect_eq "--name '' is a usage error" "125" "$(named --name '' -- true 2>/dev/null; echo $?)" +expect_eq "nested, two names: both under /mnt/9p" "a b both /mnt/9p/b" "$(named --name a -- "$NS" --unix "$TMP/n/foo.sock" --name b -- sh -c 'ls /mnt/9p | tr "\n" " "; [ -f /mnt/9p/a/README ] && [ -f /mnt/9p/b/README ] && echo both; echo $NINE_MOUNT' 2>/dev/null | tr '\n' ' ' | sed 's/ $//')" +expect_eq "nested, same name: inner shadows outer" "2 /mnt/9p/a" "$(named --name a -- "$NS" --unix "$TMP/n/foo.sock" --name a -- sh -c 'grep -c "^9ns /mnt/9p/a fuse" /proc/self/mounts; [ -f /mnt/9p/a/README ] && echo $NINE_MOUNT' 2>/dev/null | tr '\n' ' ' | sed 's/ $//')" +expect_eq "--spawn default name is the command basename" "/mnt/9p/9proc-demo" "$(timeout 60 "$NS" --spawn "$PROC --stdio" -- sh -c 'echo $NINE_MOUNT')" +expect_eq "host mount table untouched by named mounts" "no" "$(grep -q " /mnt/9p" /proc/self/mountinfo && echo yes || echo no)" echo "# lifecycle" START=$(date +%s) @@ -131,6 +153,7 @@ PIDS+=($!) sleep 0.3 TRANSPORT=(--tcp "127.0.0.1:$PORT") expect_eq "tcp: zig_version" "$(zig version)" "$(run_in "cat $M/build/zig_version")" +expect_eq "tcp: default name" "/mnt/9p/tcp-127.0.0.1-$PORT" "$(timeout 60 "$NS" --tcp "127.0.0.1:$PORT" -- sh -c 'echo $NINE_MOUNT')" expect_eq "tcp: scratch" "tcp" "$(run_in "echo tcp > $M/scratch/t && cat $M/scratch/t && rm $M/scratch/t")" echo "# server death" |
