summaryrefslogtreecommitdiff
path: root/9ns
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-22 11:18:05 -0300
committerGabriel Schneider <[email protected]>2026-09-22 11:39:25 -0300
commitb7fc01550c7bde290cf14276d94193b5b4031dc8 (patch)
tree696e8fbf26c819b28ce5552df2e7dc943a6b2a8e /9ns
parent1f3aff78702b65c328384bc5b422c448751e809b (diff)
downloadcloud9-b7fc01550c7bde290cf14276d94193b5b4031dc8.tar.gz
cloud9-b7fc01550c7bde290cf14276d94193b5b4031dc8.zip
9ns --mntgen: a server that never answers stalls only its own name
Opening a fish (self-wrapped in `9ns --mntgen`) and running an agent in it would sometimes freeze the whole session: no input reached it and nothing under /mnt/9p answered, until the shell was killed from outside. The cause was one posted server that accepted a connection and then never spoke 9P — pardes, answering its 9P from the same loop that was walking its own mount, was the one on this machine, but any wedged or half-dead server does it. Three things conspired, and each is fixed on its own: * The dispatcher dialed. A LOOKUP of an undialed name ran connect, Tversion, Tattach and Tstat on the one thread that reads /dev/fuse, so while that server kept quiet no request for any name was read, and no FUSE_INTERRUPT either. Now the dispatcher makes a Mount without touching the network and queues the walk to the mount's worker, which dials while serving it. The dial is the request in flight, so an interrupt of the walk abandons it at once (`Session.abort_on_cancel`: nothing to flush before a session exists) and the walk answers EINTR; a failed dial leaves the mount undialed for the next walk to retry; a full listen backlog (the server stopped accepting) is retried for 5s and then EIO. Every later LOOKUP of the name goes through the same queue and is answered from the remembered root attr, so the dispatcher never holds a session at all. * Once the dispatcher had read a request the process behind it was unkillable (FUSE waits out a request userspace has taken), and an INTERRUPT for a request still sitting in a mount's queue was dropped. The dispatcher now takes a queued request out and answers EINTR itself, and forwards only in-flight ones to the worker; queue and in-flight unique are read under the mount's mutex, where the worker moves a request from one to the other. A Tflush the server never answers is given 3s (`Session.flush_grace_ms`) and then the session is declared wedged: the request answers EINTR, the mount dies, the next walk makes a new one. * The kernel serialized the directory. Without FUSE_PARALLEL_DIROPS in the INIT reply every LOOKUP and READDIR in a directory takes its inode lock, so one parked walk held up every other name under /mnt/9p however free the dispatcher was (`cat` sat in fuse_lock_inode). The flag is now negotiated when the kernel offers it. What remains is the kernel's own serialization of lookups of one *name*: a second walker into the parked name waits for the first walk to end, and only then proceeds (and can be interrupted in its turn). An adversarial review of the above found three more things, fixed here: the single-connection bridge's one-slot stash stopped polling the FUSE fd while a second request was parked, so an INTERRUPT could not arrive (and parallel dirops make a second request routine) — the stash is now a queue of copies and the fd is always watched; a dead or wedged mount kept its socket open until exit, where a late-answering single-threaded server could block on it — the session is closed when the mount dies; and teardown after DESTROY or ENODEV (the child still alive, so stop_fd says nothing) could join a worker parked on a mute server forever — the sockets are shut down before the join. The flush grace is a deadline now, not a timer restarted on every wakeup. A black-box run against the binary (hostile servers: mute, garbage, close-after-accept, full backlog, 100 mute names, interrupt storms, 300 deaths of one server) found that a dead mount kept its socket, its interrupt pipe and a megabyte of buffers until exit — three descriptors per death — so `retire` now frees all of it and keeps only the slot; descriptors, threads and RSS stay flat across 400 deaths. The 4096-slot cap per process remains and is documented. Reproduced with a socket that accepts and never writes, posted beside 9agents in a scratch registry: before, `cat /mnt/9p/agents/pid` parked behind `stat /mnt/9p/hang` and SIGINT did nothing; after, it answers at once, the parked walker dies of its signal within milliseconds, and a server that answers the handshake but ignores reads and Tflush releases its reader after the grace. mntgen.sh and adv_bridge_interrupt.sh now check exactly that; nine.zig gains unit tests for the grace and the abort. Also in this change: the uncommitted ESTALE-on-death and FUSE_NOTIFY_INVAL_ENTRY work from the working copy, which the dead-mount path here builds on. Co-Authored-By: Claude Fable 5.1 <[email protected]>
Diffstat (limited to '9ns')
-rw-r--r--9ns/docs/DESIGN.md180
-rw-r--r--9ns/src/bridge.zig647
-rw-r--r--9ns/src/fuse.zig45
-rw-r--r--9ns/src/nine.zig167
-rwxr-xr-x9ns/test/adv_bridge_interrupt.sh26
-rwxr-xr-x9ns/test/mntgen.sh49
6 files changed, 812 insertions, 302 deletions
diff --git a/9ns/docs/DESIGN.md b/9ns/docs/DESIGN.md
index 6b1b80c..32c548c 100644
--- a/9ns/docs/DESIGN.md
+++ b/9ns/docs/DESIGN.md
@@ -144,8 +144,8 @@ listing that needs one meanwhile answers EIO rather than waiting.
in new userns │
/mnt/9p ─FUSE─▶ kernel ─▶ │ dispatcher (main thread)
alpha/ beta/ │ ├─ synthetic root (node 1): lists the registry
- ...each a server │ ├─ LOOKUP(alpha) ── dial+attach+stat ─▶ worker 1 ─ bridge ─ 9P session ─ server alpha
- │ └─ LOOKUP(beta) ── dial+attach+stat ─▶ worker 2 ─ bridge ─ 9P session ─ server beta
+ ...each a server │ ├─ LOOKUP(alpha) ─▶ worker 1 ─ dial+attach+stat, then bridge ─ 9P session ─ server alpha
+ │ └─ LOOKUP(beta) ─▶ worker 2 ─ dial+attach+stat, then bridge ─ 9P session ─ server beta
```
Threads and node ids:
@@ -153,17 +153,30 @@ Threads and node ids:
* **Dispatcher** (the main thread) is the only reader of `/dev/fuse`. It
answers INIT/DESTROY, serves the **synthetic root** (node 1) itself, and
routes every other request by the node id's top bits — the **mount
- index** — to the owning server's mount; a FUSE_INTERRUPT is routed by
- scanning the mounts for the one currently serving the interrupted unique
- and dropping the target into its interrupt pipe.
-* A **mount** is a dialed server: its own `nine.Session`, its own bridge
- state (inode table, handles — the ordinary single-server translation,
- unchanged) and a **worker thread**. The worker pops copied requests off a
- queue and serves them one at a time, exactly the single-connection
- contract; the dispatcher keeps reading `/dev/fuse` meanwhile, so one slow
- server never blocks the other names. Replies go straight back on the FUSE
- fd (one `writev` per reply; the kernel processes each write as one
- message).
+ index** — to the owning server's mount. It waits on no server, ever:
+ the only blocking it does is the registry directory (a tmpfs) and the
+ mount queues' locks. A FUSE_INTERRUPT is routed by scanning the mounts:
+ a request still sitting in a mount's queue is taken out and answered
+ `EINTR` on the spot (the kernel sends an INTERRUPT once, and a request
+ that only runs later would otherwise run to the end with nobody left
+ wanting it); one in flight gets its unique dropped into that mount's
+ interrupt pipe. Queue and in-flight unique are read under the mount's
+ mutex, which is also where the worker moves a request from one to the
+ other, so an interrupt cannot fall between them.
+* A **mount** is a posted name walked into: the socket to dial, its own
+ `nine.Session` and bridge state (inode table, handles — the ordinary
+ single-server translation, unchanged) once dialed, and a **worker
+ thread** that owns all of it. The dispatcher makes a mount without
+ touching the network and hands it the walk; the worker dials while
+ serving that first request. It pops copied requests off a queue and
+ serves them one at a time, exactly the single-connection contract; the
+ dispatcher keeps reading `/dev/fuse` meanwhile, so one slow server never
+ blocks the other names. Replies go straight back on the FUSE fd (one
+ `writev` per reply; the kernel processes each write as one message).
+ The queue is a plain list under a mutex and condition rather than an
+ `std.Io.Queue` because an interrupt has to find and remove a request by
+ unique from the middle of it — the same reason a 9P server keeps its
+ pending requests in a list a Tflush can search.
* Node id layout: `nodeid = (mount_index << 32) | local`. Index 0 is the
synthetic root; per-mount local ids start at 1 (the server's 9P root) and
never exceed 2^32 (a bridge never reuses one). Mount indexes are
@@ -172,40 +185,61 @@ Threads and node ids:
inode of its replacement. Reported `st_ino` mixes the index into the
qid.path (`(index+1) << 48` XOR), so two servers handing out the same
qid.path (two ramfs instances) still get distinct inode numbers.
-* **Lazy dial**: LOOKUP of an unmounted name checks the registry, dials,
- attaches, stats the root and spawns the worker — all on the dispatcher
- thread, in service of the walk that triggered it. No eager connection is
- ever made: `ls` of the root reads the registry directory only (a plain
- file dropped there is listed too — and yields EIO on the walk, never
- deleted). A walk into a **stale** entry (socket present, connect refused)
- answers EIO. The whole dial watches `stop_fd` (the session is built with
- `nine.Session.connectWatched`, so the `Tversion` exchange is covered too):
- when the program exits while a walk is parked in a dial, the dial fails
- with `Stopped`, the pending LOOKUP answers EIO and 9ns follows the program
- out. There is no dial timeout of our own (a slow server delays the walk,
- like it would delay any 9P client), and a server that accepts but never
- answers `Tversion` still parks the dispatcher until the program exits —
- including the unkillable corner where the *blocked walk itself* is the
- only thing keeping the program alive (the task sits in D state until the
- filesystem answers; a same-user self-DoS, accepted with the pinned
- "dispatcher dials, no concurrent dial" design).
+* **Lazy dial, on the worker**: LOOKUP of an unmounted name stats the
+ registry entry (dispatcher), makes the mount and queues the LOOKUP to it;
+ the worker dials (connect, `Tversion`, `Tattach`, `Tstat` of the root) as
+ the first thing it does for that request, then answers it from the root
+ stat. Every later LOOKUP of the name is queued the same way and answered
+ from the remembered root attr, so the dispatcher never holds a session.
+ No eager connection is ever made: `ls` of the root reads the registry
+ directory only (a plain file dropped there is listed too — and yields EIO
+ on the walk, never deleted). A dial that fails answers the walk that
+ asked — `ENOENT` when the entry vanished, `EIO` for a stale entry
+ (connect refused), a full backlog or a server that will not speak 9P —
+ and leaves the mount undialed, so the next walk simply tries again.
+ The dial is the request in flight, so it ends the way any request does:
+ `stop_fd` (the program exited) fails it with `Stopped`; an interrupt of
+ the walk abandons it at once, with nothing sent (`abort_on_cancel`: there
+ is no session yet to flush anything out of) and the walk answers
+ `EINTR`. A server that accepts but never answers `Tversion` therefore
+ costs exactly the walks into its own name, and Ctrl-C ends those.
+ A server whose listen backlog is full (it stopped accepting) makes the
+ connect report `EAGAIN`; that is retried for 5s, polling `stop_fd` and
+ the interrupt pipe between tries, then answers `EIO`. What no design can
+ fix is the kernel side: the VFS serializes lookups of one *name*, so a
+ second walker into the parked name waits in `d_wait_lookup` until the
+ first walk ends — interrupt the first, and the second proceeds (and can
+ be interrupted in its turn).
+* **FUSE_PARALLEL_DIROPS** is negotiated in the INIT reply. Without it the
+ kernel takes the directory inode's lock around every LOOKUP and READDIR
+ in it, so one parked walk would still hold up every other name under the
+ same directory — the whole registry root, for a mntgen mount — however
+ free the dispatcher is.
* **Death and re-dial**: when a worker's session dies mid-request, dispatch
- has already answered that request EIO, the mount is marked dead, and
- everything further routed to that subtree answers EIO (a FORGET is
- dropped). Nothing reconnects eagerly. Because synthetic-root entries are
- served with zero entry-validity, the next walk into the name LOOKUPs it
- again; a dead mount is skipped and the name is dialed afresh — a new
- mount under a new index, so kernel-held inodes of the corpse keep
- answering EIO until forgotten. `ls` still lists the dead name (listing
- connects to nothing). Death is discovered lazily: the first walk after a
- silent death answers EIO (it marks the mount dead), the next walk re-dials.
+ has already answered that request (`ESTALE` for a lost connection, so the
+ VFS redoes the path walk instead of failing; `EIO` for a protocol error),
+ the mount is marked dead, its entry is invalidated in the kernel's dentry
+ cache (`FUSE_NOTIFY_INVAL_ENTRY`, best effort), and everything further
+ routed to that subtree answers `ESTALE` (a FORGET is dropped). Nothing
+ reconnects eagerly. The next walk into the name LOOKUPs it again; a dead
+ mount is skipped and the name gets a new mount under a new index, so
+ kernel-held inodes of the corpse keep answering `ESTALE` until forgotten.
+ `ls` still lists the dead name (listing connects to nothing). Death is
+ discovered lazily: the first walk after a silent death takes the error
+ (it marks the mount dead), the next walk re-dials.
* **Interrupts** work per mount: the worker's session polls the mount's
interrupt pipe while a 9P reply is outstanding; the dispatcher writes the
interrupted request's unique into it and the usual `Tflush` dance
- (see *Interrupts*) follows. `stop_fd` (the child's death) is watched by
- every session, so no worker can stay blocked on a hung server past the
- program's exit; teardown wakes every worker, joins them, and tears down
- their sessions and bridge state.
+ (see *Interrupts*) follows — with a grace: a server that answers neither
+ the request nor the `Tflush` within 3s (`nine.Session.flush_grace_ms`)
+ is declared gone, the request answers `EINTR`, the session is wedged and
+ the mount dies with it (the next walk makes a new one). The protocol
+ says a client waits for the Rflush; a server that has not managed one in
+ that long is not going to, and the process behind the interrupt is
+ unkillable until we stop waiting. `stop_fd` (the child's death) is
+ watched by every session, so no worker can stay blocked on a hung server
+ past the program's exit; teardown wakes every worker, joins them, and
+ tears down their sessions and bridge state.
* The synthetic root is read-only (`dr-xr-xr-x`, like `/srv`): services are
posted and unposted by their servers (`cloud9.post`'s
`post`/`listenPosted`/`unpost`), not created and removed through files.
@@ -290,7 +324,8 @@ setupmapping = 48, removemapping = 49, syncfs = 50, tmpfile = 51, statx = 52, _
Constants: `kernel_version = 7`, `kernel_minor = 31` (what we answer; the
kernel adapts to the lower minor), `FOPEN_DIRECT_IO = 1`, `FOPEN_KEEP_CACHE = 2`,
-`FOPEN_NONSEEKABLE = 4`, `FUSE_ASYNC_READ = 1`, `FUSE_MAX_PAGES = 1<<22`,
+`FOPEN_NONSEEKABLE = 4`, `FUSE_ASYNC_READ = 1`, `FUSE_PARALLEL_DIROPS = 1<<18`,
+`FUSE_MAX_PAGES = 1<<22`,
`FATTR_MODE=1, FATTR_UID=2, FATTR_GID=4, FATTR_SIZE=8, FATTR_ATIME=16,
FATTR_MTIME=32, FATTR_FH=64, FATTR_ATIME_NOW=128, FATTR_MTIME_NOW=256,
FATTR_LOCKOWNER=512, FATTR_CTIME=1024`. `root_id = 1`.
@@ -409,7 +444,7 @@ pub fn serve(gpa: std.mem.Allocator, fuse_fd: i32, nine: *nine.Session, root_fid
// mntgen (see "mntgen: one mount, many servers" under Process model):
pub const MntgenOptions = struct {
- io: std.Io, // dispatcher-thread only: post.posted / post.dial
+ io: std.Io, // dispatcher-thread only: post.posted / registry stats
env: post.Env, // XDG_RUNTIME_DIR names the registry
uname: []const u8, aname: []const u8 = "", msize: u32 = 131072,
};
@@ -449,7 +484,7 @@ Op mapping (9P2000 has no symlinks, links, xattrs, locks, mknod):
| FUSE | 9P |
|---|---|
-| INIT | reply `InitOut{ major=7, minor=31, max_readahead=in.max_readahead, flags = FUSE_ASYNC_READ \| FUSE_ATOMIC_O_TRUNC \| FUSE_AUTO_INVAL_DATA \| FUSE_BIG_WRITES (plus FUSE_MAX_PAGES with max_pages=256 if offered), max_background=16, congestion_threshold=12, max_write=1 MiB, time_gran=1 }`. Atomic O_TRUNC matters: without it the kernel truncates via a separate SETATTR(size=0) that synthetic control files reject; with it `O_TRUNC` becomes 9P `OTRUNC` inside the open |
+| INIT | reply `InitOut{ major=7, minor=31, max_readahead=in.max_readahead, flags = FUSE_ASYNC_READ \| FUSE_ATOMIC_O_TRUNC \| FUSE_AUTO_INVAL_DATA \| FUSE_BIG_WRITES (plus FUSE_MAX_PAGES with max_pages=256, and FUSE_PARALLEL_DIROPS, each if offered), max_background=16, congestion_threshold=12, max_write=1 MiB, time_gran=1 }`. Atomic O_TRUNC matters: without it the kernel truncates via a separate SETATTR(size=0) that synthetic control files reject; with it `O_TRUNC` becomes 9P `OTRUNC` inside the open |
| LOOKUP(parent,name) | `walk(parent.fid → newfid, [name])`; `stat(newfid)`; dedupe by qid; `EntryOut` |
| FORGET / BATCH_FORGET | `nlookup -= n`; at 0 `clunk` and drop (no reply) |
| GETATTR | `stat(inode.fid)` → `AttrOut` |
@@ -512,9 +547,10 @@ is outstanding, including the initial root stat:
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.
+ copied into a queue that `serve` dispatches, in arrival order, before it
+ polls again; the fd stays watched throughout, so an INTERRUPT is never
+ stuck behind a parked request (a copy that cannot be allocated answers
+ `ENOMEM` on the spot).
* 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
@@ -532,11 +568,14 @@ is outstanding, including the initial root stat:
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).
+Limits: a server that ignores Tflush is given `flush_grace_ms` (3s) and then
+declared gone — flush(5) says the server "should answer the Tflush message
+immediately", and a client that waits longer leaves an unkillable process
+behind — so the reader gets `EINTR` and the session ends (the hostile
+`never` mode; in single-connection mode that is the mount, in mntgen that
+one name). A slow-but-alive server that cannot answer a flush inside its own
+blocking read pays the same price; the alternative was a hang that only
+SIGKILL of 9ns could end.
### `src/ns.zig` — namespace and process plumbing
@@ -676,26 +715,27 @@ ReleaseSafe.
(event files) stalls that server's subtree while it is outstanding (but
not past the child's exit). It can be interrupted: killing or Ctrl-C-ing
the reader sends `FUSE_INTERRUPT`, which becomes `Tflush`; servers that
- honour it unblock immediately, servers that don't still block that
- subtree until they answer.
-* mntgen: the dispatcher dials on the main thread, in service of the walk
- that triggered it: a server that accepts the connection but never
- answers `Tversion` parks the whole mount for as long as the walk's
- program keeps running (as it would delay any 9P client). The dial
- watches `stop_fd`, so the program exiting ends it (`Stopped`, EIO to the
- pending LOOKUP); no concurrent dial, no dial timeout of our own. The
- remaining corner is a same-user self-DoS: the program blocked *in that
- very walk* sits in D state until the filesystem answers and cannot be
- killed to fire `stop_fd` — only killing the mute server (EOF) ends it.
+ honour it unblock immediately, servers that don't are given 3s and then
+ declared gone: the reader gets `EINTR` and the session (single
+ connection: the mount; mntgen: that name) ends.
+* mntgen: a name's dial runs on that name's worker, inside the walk that
+ asked, with no timeout of its own beyond the 5s backlog retry: a server
+ that accepts and never answers `Tversion` parks the walks into its own
+ name, and only those, until each is interrupted or the program exits.
+ The one wait nothing here can cut short is the kernel's own
+ serialization of lookups of a single name (a second walker into the
+ parked name waits for the first walk to end, uninterruptibly).
* mntgen: mount indexes are ordinals and never reused (cap 4096 dials per
process); per-mount node ids cap at 2^32 lookups. `st_ino` mixes the
index with the qid.path via XOR — collisions remain theoretically
possible, just not the practical ones (identical qid.paths across two
- servers). A dead mount keeps its slot — and its session socket and
- interrupt-pipe descriptors — until process exit (early close would race
- the dispatcher's INTERRUPT scan against fd reuse), so with a small
- `ulimit -n` a re-dial storm exhausts descriptors before the index cap;
- dials then fail cleanly with EIO.
+ servers). A dead mount releases its socket, its interrupt pipe and its
+ bridge state at once (the pipe under the mount's mutex, where the
+ dispatcher writes it) and keeps only its slot, so a server that dies and
+ comes back costs one index per death and nothing else; after 4096 of
+ them in one 9ns process no further name can be walked into (EIO) until
+ the shell is restarted. Reclaiming a slot once the kernel has forgotten
+ every node of the corpse is the next step, not taken here.
* mntgen: no posting/unposting through the mount (the synthetic root is
read-only; servers manage their registry entries through
`cloud9.post`), and no per-name mount options: one set of
diff --git a/9ns/src/bridge.zig b/9ns/src/bridge.zig
index 457c883..327cde1 100644
--- a/9ns/src/bridge.zig
+++ b/9ns/src/bridge.zig
@@ -9,7 +9,8 @@
//! 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.
+//! sends meanwhile is copied into a queue and served next, so the FUSE fd
+//! is always being read and no INTERRUPT waits behind a parked request.
const std = @import("std");
const cloud9 = @import("cloud9");
const post = cloud9.post;
@@ -132,13 +133,12 @@ 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;
+ // Requests that arrived while a 9P reply was outstanding go first,
+ // in the order the kernel sent them.
+ if (b.stash.items.len != 0) {
+ const raw = b.stash.orderedRemove(0);
+ defer gpa.free(raw);
+ if (!try b.dispatch(requestOf(raw))) return;
continue;
}
if (b.fuse_gone) return;
@@ -181,15 +181,21 @@ const Bridge = struct {
/// same qid.path (two ramfs instances, say) cannot share an inode number
/// inside one mntgen mount. 0 in single-connection mode.
ino_xor: u64 = 0,
+ /// mntgen only: answer a lost 9P connection with ESTALE instead of EIO,
+ /// so the VFS redoes the path walk and re-dials the replacement server
+ /// (see the retired-mount branch in `routeToMount`). False in
+ /// single-connection mode, where there is no second server to find and
+ /// the retry would only turn one EIO into one ESTALE.
+ stale_on_death: bool = false,
req_buf: []align(8) u8 = &.{},
- /// Second request buffer: what the interrupt poll reads into. Holds the
- /// stashed request until `serve` swaps it in.
+ /// Second request buffer: what the interrupt poll reads into.
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,
+ /// Non-INTERRUPT requests read while a 9P reply was outstanding, copied
+ /// out of `spare_buf` in arrival order; `serve` dispatches them before
+ /// reading the fd again. Single-connection mode only (mntgen's
+ /// dispatcher owns the fd and queues to the workers instead).
+ stash: std.ArrayList([]align(8) u8) = .empty,
/// `unique` of the FUSE request being served, if any (0 = none): the only
/// one an INTERRUPT may cancel. Atomic: in mntgen mode the dispatcher
/// thread scans it to route FUSE_INTERRUPTs to the right server.
@@ -221,6 +227,8 @@ const Bridge = struct {
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);
+ for (b.stash.items) |raw| b.gpa.free(raw);
+ b.stash.deinit(b.gpa);
}
// -- interrupt source (polled by nine.Session while a reply is outstanding) -----
@@ -229,16 +237,17 @@ const Bridge = struct {
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.
+ /// Poll the FUSE fd for as long as it is alive: an INTERRUPT must be
+ /// able to arrive whatever else is queued.
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;
+ return if (!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.
+ /// dropped (the kernel expects no reply); anything else is copied into
+ /// the stash, or answered ENOMEM if it cannot be.
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) {
@@ -269,7 +278,16 @@ const Bridge = struct {
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;
+ const raw = b.spare_buf[0..h.len];
+ const copy = b.gpa.alignedAlloc(u8, .@"8", raw.len) catch {
+ if (wantsReply(raw)) b.replyError(h.unique, .NOMEM) catch {};
+ return false;
+ };
+ @memcpy(copy, raw);
+ b.stash.append(b.gpa, copy) catch {
+ b.gpa.free(copy);
+ if (wantsReply(raw)) b.replyError(h.unique, .NOMEM) catch {};
+ };
return false;
}
@@ -315,7 +333,8 @@ const Bridge = struct {
error.OutOfMemory => .NOMEM,
error.TooLarge => .NAMETOOLONG,
error.BadDir => .IO,
- error.Closed, error.Protocol, error.Io, error.Stopped => .IO,
+ error.Closed, error.Io => if (b.stale_on_death) .STALE else .IO,
+ error.Protocol, error.Stopped => .IO,
error.Interrupted => .INTR,
error.FuseIo => return error.FuseIo,
};
@@ -326,6 +345,10 @@ const Bridge = struct {
else => {},
}
};
+ // An EINTR that came from the server going mute (a Tflush unanswered
+ // past its grace) is the last thing this session says: the
+ // connection is gone with it.
+ if (b.nine.wedged) return error.Closed;
return true;
}
@@ -832,15 +855,22 @@ const Bridge = struct {
//
// `9ns --mntgen` serves one FUSE mount whose synthetic root lists the posted
// 9P services in `$XDG_RUNTIME_DIR/9p` (cloud9.post's registry; no connection
-// is made to list). A walk into a name dials that server lazily and starts a
-// per-server worker thread running the ordinary bridge translation above.
+// is made to list). A walk into a name makes a mount for it and hands the
+// walk to the mount's worker thread, which dials the server (the only
+// thread that ever talks to it) and then runs the ordinary bridge
+// translation above for everything under the name. The dispatcher thread
+// reads /dev/fuse, routes, and serves the registry's own directories; it
+// never waits on a server, so one server that never answers holds up only
+// the walks into its own name — and those, being ordinary requests on a
+// worker, an interrupt can still end.
//
// Node ids carry the server in the top bits: a request for `(index, local)`
// is routed by `index` to that server's mount. Indexes are ordinals, never
// reused, so a stale kernel-side inode of a dead server can never be
-// conflated with a fresh inode of its replacement. A dead server answers EIO
-// on its whole subtree until the next walk into its name re-dials it (a new
-// mount, a new index); nothing reconnects eagerly.
+// conflated with a fresh inode of its replacement. A dead server answers
+// ESTALE on its whole subtree until the next walk into its name makes a new
+// mount (a new index); nothing reconnects eagerly. A dial that fails leaves
+// the mount as it was, undialed: the next walk simply tries again.
/// Bit position of the mount index inside a FUSE node id; the low bits are
/// one server's bridge node ids, the top bits name the server (0 = the
@@ -889,8 +919,9 @@ pub fn nameIno(name: []const u8) u64 {
pub const MntgenOptions = struct {
/// Used only on the dispatcher (main) thread: registry listing and
- /// dial. The worker threads never call io (their locks use the
- /// uncancelable futex paths, which are thread-safe globals).
+ /// stats. The worker threads never call io (their locks use the
+ /// uncancelable futex paths, which are thread-safe globals; their
+ /// sockets are raw syscalls).
io: std.Io,
/// Environment block (post.Env); XDG_RUNTIME_DIR names the registry.
env: post.Env,
@@ -925,35 +956,63 @@ const SynthDir = struct {
}
};
-/// One dialed server: its 9P session, its bridge state and its worker
-/// thread, plus the queue the dispatcher feeds requests through.
+/// One posted name walked into: the socket to dial, the 9P session and
+/// bridge that come of dialing it, the worker thread that owns those, and
+/// the queue the dispatcher feeds it. The dispatcher makes a mount without
+/// touching the network; the worker dials while serving the first walk
+/// into the name, so a server that never answers costs exactly the walks
+/// into its own name, each of them interruptible, and nothing else.
const Mount = struct {
gpa: std.mem.Allocator,
io: std.Io,
+ mo: *const MntgenOptions,
debug: bool,
+ /// Registry-relative key ("agents", or "sub/dir/agents").
name: []u8,
index: u32,
/// This server's root node id (index in the top bits, 1 below).
root_node: u64,
- /// Attr of the server's 9P root, cached from the dial-time stat: the
- /// dispatcher answers LOOKUP of the name from it without an rpc (the
- /// session belongs to the worker thread).
- root_attr: fuse.Attr,
- session: *nine.Session,
+ /// Where this mount hangs in the synthetic tree: the node id of the
+ /// directory holding its name (the mntgen root, or a synthetic registry
+ /// subdirectory) and the single name component under it — the last
+ /// component of `name`, which is what the kernel caches the dentry
+ /// under. Together they are what a FUSE_NOTIFY_INVAL_ENTRY needs when
+ /// the server dies.
+ parent_node: u64,
+ entry_name: []u8,
+ /// The registry socket, NUL-terminated.
+ sock_buf: [post.sun_path_len]u8 = undefined,
+ sock_len: u16 = 0,
+ /// The program's exit ends a dial or an rpc in progress.
+ stop_fd: i32,
+ /// The 9P session, from a successful dial until the connection dies
+ /// (or teardown). Null before the dial: a walk into the name (the only
+ /// request a mount without nodes can receive) dials first. Written by
+ /// the worker under `mutex`, so the dispatcher can shut the socket down
+ /// at teardown; `root_attr` and the bridge's root inode exist with it.
+ session: ?nine.Session = null,
+ /// Attr of the server's 9P root from the dial-time stat: what a
+ /// LOOKUP of the name answers.
+ root_attr: fuse.Attr = undefined,
b: *Bridge,
+ /// Guards `queue`, `stopping`, `session`'s existence, and the moment a
+ /// request leaves the queue for `b.cur_unique` (see `routeInterrupt`).
mutex: std.Io.Mutex = .init,
cond: std.Io.Condition = .init,
queue: std.ArrayList([]align(8) u8) = .empty,
- /// The 9P connection died: the subtree answers EIO; a walk into the
- /// name re-dials as a new mount. Set by the worker, read by all.
+ /// The 9P connection died: the subtree answers ESTALE; a walk into the
+ /// name makes a new mount. Set by the worker, read by all.
dead: std.atomic.Value(bool) = .init(false),
/// Draining and exiting (child gone, FUSE device gone or DESTROY).
- /// Under `mutex`.
stopping: bool = false,
/// The dispatcher drops a FUSE_INTERRUPT's target unique in here; the
/// worker's session polls it while a 9P reply is outstanding.
int_pipe: [2]i32,
thread: std.Thread,
+
+ fn sockPath(m: *const Mount) [:0]const u8 {
+ return m.sock_buf[0..m.sock_len :0];
+ }
};
/// Runs the mntgen dispatcher on the calling thread until the FUSE fd
@@ -1043,6 +1102,11 @@ const Mntgen = struct {
const m = slot orelse continue;
m.mutex.lockUncancelable(mg.io);
m.stopping = true;
+ // A worker parked in an rpc on a mute server (DESTROY and
+ // ENODEV come with the child still alive, so stop_fd says
+ // nothing) reads EOF instead and comes out through the death
+ // path; the fd stays the worker's to close.
+ if (m.session) |*sess| _ = linux.shutdown(sess.fd, linux.SHUT.RDWR);
m.mutex.unlock(mg.io);
m.cond.signal(mg.io);
}
@@ -1051,13 +1115,15 @@ const Mntgen = struct {
m.thread.join();
for (m.queue.items) |buf| mg.gpa.free(buf);
m.queue.deinit(mg.gpa);
- _ = linux.close(m.int_pipe[0]);
- _ = linux.close(m.int_pipe[1]);
- m.b.deinit();
- m.session.deinit();
+ // A mount that died released all of this itself (`retire`).
+ if (!m.dead.load(.seq_cst)) {
+ _ = linux.close(m.int_pipe[0]);
+ _ = linux.close(m.int_pipe[1]);
+ m.b.deinit();
+ if (m.session) |*sess| sess.deinit();
+ }
mg.gpa.free(m.name);
mg.gpa.destroy(m.b);
- mg.gpa.destroy(m.session);
mg.gpa.destroy(m);
}
for (&mg.root_dirs) |*slot| {
@@ -1149,10 +1215,23 @@ const Mntgen = struct {
return true;
}
// A retired mount (its server died and a later walk re-dialed under
- // a new index) or an unknown node: the dead subtree answers EIO, a
- // FORGET is simply dropped.
+ // a new index) or an unknown node; a FORGET is simply dropped.
+ //
+ // ESTALE rather than EIO, because it is both truer and useful: the
+ // node id named a file on a server that is gone, which is precisely
+ // a stale handle, and the VFS answers ESTALE by redoing the path walk
+ // with LOOKUP_REVAL instead of failing. The mount's death already
+ // invalidated its dentry, so that second walk re-LOOKUPs the name,
+ // re-dials the server and succeeds — a restarted server costs a
+ // retry inside one syscall rather than a visible error.
+ //
+ // A read(2) on an fd opened before the death still fails: there is no
+ // path left to re-walk, and inventing one would be a lie about which
+ // file the caller holds. It is open(2) — where the path is still in
+ // hand — that recovers, which is what a caller re-running `cat` or a
+ // program reopening its config actually needs.
mg.trace(" nodeid={d} has no live mount (index {d})", .{ h.nodeid, idx });
- if (wants_reply) mg.replyError(h.unique, .IO) catch {};
+ if (wants_reply) mg.replyError(h.unique, .STALE) catch {};
return true;
}
@@ -1172,8 +1251,10 @@ const Mntgen = struct {
@memcpy(buf, raw);
var reject: ?linux.E = null;
m.mutex.lockUncancelable(mg.io);
- if (m.dead.load(.seq_cst) or m.stopping) {
- reject = .IO;
+ if (m.dead.load(.seq_cst)) {
+ reject = .STALE; // raced with the death; recoverable, as in routeToMount
+ } else if (m.stopping) {
+ reject = .IO; // draining for good: no replacement is coming
} else if (m.queue.append(mg.gpa, buf)) |_| {
m.cond.signal(mg.io);
} else |_| {
@@ -1249,22 +1330,45 @@ const Mntgen = struct {
}
/// A FUSE_INTERRUPT names the request it wants cancelled; FUSE uniques
- /// are unique across the whole connection, so the mount whose bridge is
- /// currently serving that unique gets the packet and its session turns
- /// it into a Tflush (see the worker's interrupt source below).
+ /// are unique across the whole connection, so exactly one mount holds
+ /// it. Still queued there, it has not started: it is taken out and
+ /// answered EINTR here, because the kernel sends an INTERRUPT once and a
+ /// request that only runs later would otherwise run to the end with
+ /// nobody left wanting it. In flight, the worker gets the packet and its
+ /// session turns it into a Tflush (see the worker's interrupt source
+ /// below). The queue and `cur_unique` are read under the mount's mutex,
+ /// which is also where the worker moves a request from one to the
+ /// other, so an interrupt cannot fall between them.
fn routeInterrupt(mg: *Mntgen, target: u64) void {
for (mg.mounts) |slot| {
const m = slot orelse continue;
if (m.dead.load(.seq_cst)) continue;
- if (m.b.cur_unique.load(.seq_cst) == target) {
- mg.trace(" interrupt for unique={d}: forwarding to '{s}'", .{ target, m.name });
+ m.mutex.lockUncancelable(mg.io);
+ for (m.queue.items, 0..) |raw, i| {
+ // Nothing waits on a FORGET, so nothing interrupts one; a
+ // unique that names one anyway is not ours to take out.
+ if (uniqueOf(raw) != target or !wantsReply(raw)) continue;
+ _ = m.queue.orderedRemove(i);
+ m.mutex.unlock(mg.io);
+ mg.trace(" interrupt for unique={d}: still queued at '{s}'; answered EINTR", .{ target, m.name });
+ mg.replyError(target, .INTR) catch {};
+ mg.gpa.free(raw);
+ return;
+ }
+ const in_flight = m.b.cur_unique.load(.seq_cst) == target;
+ if (in_flight and m.int_pipe[1] >= 0) {
+ // Under the mutex, because the worker closes the pipe there
+ // when the mount dies. The pipe is small and nonblocking; a
+ // dropped packet only means one interrupt missed its window
+ // (the kernel does not retry INTERRUPTs, but the child's
+ // exit ends the session through stop_fd regardless).
var packet: [8]u8 = undefined;
std.mem.writeInt(u64, &packet, target, .little);
- // The pipe is small and nonblocking; a dropped packet only
- // means one interrupt missed its window (the kernel does
- // not retry INTERRUPTs, but the child's exit ends the
- // session through stop_fd regardless).
_ = linux.write(m.int_pipe[1], &packet, packet.len);
+ }
+ m.mutex.unlock(mg.io);
+ if (in_flight) {
+ mg.trace(" interrupt for unique={d}: forwarded to '{s}'", .{ target, m.name });
return;
}
}
@@ -1286,14 +1390,6 @@ const Mntgen = struct {
};
}
- fn rootEntryOut(mg: *Mntgen, nodeid: u64, attr: fuse.Attr) fuse.EntryOut {
- // No entry caching for synthetic-root entries: a walk re-LOOKUPs the
- // name, which is what notices a dead server and re-dials it. No
- // invalidation machinery needed, and nothing to invalidate.
- _ = mg;
- return .{ .nodeid = nodeid, .generation = 0, .attr = attr };
- }
-
fn handleRoot(mg: *Mntgen, req: fuse.Request) error{FuseIo}!void {
const u = req.header.unique;
switch (req.header.op()) {
@@ -1361,21 +1457,18 @@ const Mntgen = struct {
const u = req.header.unique;
const name = fuse.nameAfter(void, req) catch return mg.replyError(u, .INVAL);
if (std.mem.eql(u8, name, ".") or std.mem.eql(u8, name, "..")) {
- const out = mg.rootEntryOut(fuse.root_id, mg.rootAttr());
+ const out = entryOut(fuse.root_id, mg.rootAttr());
return mg.reply(u, &.{std.mem.asBytes(&out)});
}
if (!post.legalName(name)) return mg.replyError(u, .NOENT);
if (findMount(mg.mounts, name)) |m| {
- // Live: answer from the dial-time snapshot. No 9P rpc (the
- // session belongs to the worker thread); no connection made.
- mg.trace(" lookup '{s}': mount {d} already live", .{ name, m.index });
- const out = mg.rootEntryOut(m.root_node, m.root_attr);
- return mg.reply(u, &.{std.mem.asBytes(&out)});
+ mg.trace(" lookup '{s}': mount {d}", .{ name, m.index });
+ return mg.enqueue(m, mg.bytesOf(req));
}
// What the entry is decides what a walk into it becomes: a socket
- // dials (the original behavior), a directory is served like the
- // root itself (its sockets dial on walk, its directories recurse),
- // anything else answers EIO.
+ // becomes a mount (dialed by its worker), a directory is served like
+ // the root itself (its sockets dial on walk, its directories
+ // recurse), anything else answers EIO.
var path_buf: [post.sun_path_len]u8 = undefined;
const entry_path = post.registryPath(mg.mo.env, name, &path_buf) catch
return mg.replyError(u, .NOENT);
@@ -1390,28 +1483,27 @@ const Mntgen = struct {
return mg.replyError(u, .IO);
};
const node = sd.nodeid();
- const out = mg.rootEntryOut(node, mg.synthAttr(node));
+ const out = entryOut(node, mg.synthAttr(node));
return mg.reply(u, &.{std.mem.asBytes(&out)});
},
- .unix_domain_socket => {},
+ .unix_domain_socket => {
+ const m = mg.newMount(entry_path, name, fuse.root_id) catch |e| {
+ mg.trace(" lookup '{s}': no mount: {t}", .{ name, e });
+ return mg.replyError(u, .IO);
+ };
+ mg.trace(" lookup '{s}': new mount {d}", .{ name, m.index });
+ mg.enqueue(m, mg.bytesOf(req));
+ },
else => {
mg.trace(" lookup '{s}': registry entry is not a socket", .{name});
return mg.replyError(u, .IO);
},
}
- const m = mg.dialMount(name) catch |e| switch (e) {
- error.Stale => {
- mg.trace(" lookup '{s}': registry entry is stale (no server behind it)", .{name});
- return mg.replyError(u, .IO);
- },
- else => {
- mg.trace(" lookup '{s}': dial failed: {t}", .{ name, e });
- return mg.replyError(u, .IO);
- },
- };
- mg.trace(" lookup '{s}': dialed as mount {d}", .{ name, m.index });
- const out = mg.rootEntryOut(m.root_node, m.root_attr);
- try mg.reply(u, &.{std.mem.asBytes(&out)});
+ }
+
+ /// The bytes of the request being routed, as read from the FUSE fd.
+ fn bytesOf(mg: *const Mntgen, req: fuse.Request) []const u8 {
+ return mg.req_buf[0..req.header.len];
}
/// One OPENDIR of the synthetic root: `.` and `..` plus every registry
@@ -1522,7 +1614,7 @@ const Mntgen = struct {
const u = req.header.unique;
const name = fuse.nameAfter(void, req) catch return mg.replyError(u, .INVAL);
if (std.mem.eql(u8, name, ".")) {
- const out = mg.rootEntryOut(sd.nodeid(), mg.synthAttr(sd.nodeid()));
+ const out = entryOut(sd.nodeid(), mg.synthAttr(sd.nodeid()));
return mg.reply(u, &.{std.mem.asBytes(&out)});
}
// The kernel resolves ".." from its own dentry tree and a LOOKUP of
@@ -1530,7 +1622,7 @@ const Mntgen = struct {
// onto itself either: answer with the parent's node id (the root's
// for a top-level subdirectory — its attr is the same shape).
if (std.mem.eql(u8, name, "..")) {
- const out = mg.rootEntryOut(sd.parent, mg.synthAttr(sd.parent));
+ const out = entryOut(sd.parent, mg.synthAttr(sd.parent));
return mg.reply(u, &.{std.mem.asBytes(&out)});
}
if (!post.legalName(name)) return mg.replyError(u, .NOENT);
@@ -1543,9 +1635,8 @@ const Mntgen = struct {
const child = std.fmt.bufPrintSentinel(&path_buf, "{s}/{s}", .{ sd.path(), name }, 0) catch
return mg.replyError(u, .NOTNAM);
if (findMount(mg.mounts, key)) |m| {
- mg.trace(" lookup '{s}': mount {d} already live", .{ key, m.index });
- const out = mg.rootEntryOut(m.root_node, m.root_attr);
- return mg.reply(u, &.{std.mem.asBytes(&out)});
+ mg.trace(" lookup '{s}': mount {d}", .{ key, m.index });
+ return mg.enqueue(m, mg.bytesOf(req));
}
const st = std.Io.Dir.statFile(.cwd(), mg.mo.io, child, .{}) catch |e| switch (e) {
error.FileNotFound => {
@@ -1564,17 +1655,16 @@ const Mntgen = struct {
return mg.replyError(u, .IO);
};
const node = child_sd.nodeid();
- const out = mg.rootEntryOut(node, mg.synthAttr(node));
+ const out = entryOut(node, mg.synthAttr(node));
return mg.reply(u, &.{std.mem.asBytes(&out)});
},
.unix_domain_socket => {
- const m = mg.dialMountAt(child, key) catch |e| {
- mg.trace(" lookup '{s}': dial failed: {t}", .{ key, e });
+ const m = mg.newMount(child, key, sd.nodeid()) catch |e| {
+ mg.trace(" lookup '{s}': no mount: {t}", .{ key, e });
return mg.replyError(u, .IO);
};
- mg.trace(" lookup '{s}': dialed as mount {d}", .{ key, m.index });
- const out = mg.rootEntryOut(m.root_node, m.root_attr);
- try mg.reply(u, &.{std.mem.asBytes(&out)});
+ mg.trace(" lookup '{s}': new mount {d}", .{ key, m.index });
+ mg.enqueue(m, mg.bytesOf(req));
},
// Not a service and not a directory: the entry answers EIO on
// walk, like a plain file in the registry itself.
@@ -1633,74 +1723,32 @@ const Mntgen = struct {
return sd;
}
- // -- dialing ---------------------------------------------------------------
-
- /// Dials `name` out of the registry, attaches, stats the server root and
- /// starts its worker thread. Runs on the dispatcher thread, in service
- /// of the LOOKUP that triggered it (so a hung server delays that walk,
- /// like it would delay any 9P client).
- fn dialMount(mg: *Mntgen, name: []const u8) !*Mount {
- const stream = try post.dial(mg.mo.io, mg.mo.env, name);
- return mg.mountStream(stream, name);
- }
-
- /// Dials the socket at `path` (a registry subdirectory entry) and
- /// mounts it under `key`, the registry-relative path — the
- /// subdirectory analogue of `dialMount`.
- fn dialMountAt(mg: *Mntgen, path: [:0]const u8, key: []const u8) !*Mount {
- const stream = try post.dialPath(mg.mo.io, path);
- return mg.mountStream(stream, key);
- }
+ // -- mounts --------------------------------------------------------------
- /// The shared dial tail: session, attach, stat, bridge and worker.
- fn mountStream(mg: *Mntgen, stream: std.Io.net.Stream, key: []const u8) !*Mount {
+ /// A mount for the socket at `sock_path`, keyed by the registry-relative
+ /// `key`, under `parent_node` (the mntgen root or a synthetic
+ /// subdirectory). Nothing is dialed here: the bridge, the pipe and the
+ /// worker are made, and the worker dials when the first walk reaches
+ /// it. Runs on the dispatcher thread and blocks on nothing.
+ fn newMount(mg: *Mntgen, sock_path: [:0]const u8, key: []const u8, parent_node: u64) !*Mount {
if (mg.next_index >= max_mounts) return error.TooManyMounts;
+ if (sock_path.len >= post.sun_path_len) return error.NameTooLong;
const index: u32 = mg.next_index;
- const fd: i32 = @intCast(stream.socket.handle);
- // The session below does blocking I/O: make sure a dial that left
- // the descriptor nonblocking cannot spin its read loop on EAGAIN,
- // and keep the descriptor out of any future exec (the running
- // program already forked, but hygiene is free).
- _ = linux.fcntl(fd, linux.F.SETFL, 0);
- _ = linux.fcntl(fd, linux.F.SETFD, linux.FD_CLOEXEC);
-
- const session = mg.gpa.create(nine.Session) catch |e| {
- _ = linux.close(fd);
- return e;
- };
- errdefer mg.gpa.destroy(session);
- // The child's exit must unblock the whole dial — Tversion included,
- // not just the attach/stat after it: a silent server must not pin
- // the dispatcher past the program.
- session.* = nine.Session.connectWatched(mg.gpa, .{ .fd = fd }, mg.mo.msize, mg.stop_fd) catch |e| {
- _ = linux.close(fd);
- return e;
- };
- errdefer session.deinit();
- _ = try session.attach(0, mg.mo.uname, mg.mo.aname);
- const st = try session.stat(0);
+ const root_node = mountNode(index, fuse.root_id);
const b = try mg.gpa.create(Bridge);
errdefer mg.gpa.destroy(b);
- const root_node = mountNode(index, fuse.root_id);
b.* = .{
.gpa = mg.gpa,
.fuse_fd = mg.fuse_fd,
- .nine = session,
+ .nine = undefined, // the mount's session, once the worker has dialed
.opts = mg.opts,
.root_id = root_node,
.ino_xor = @as(u64, index + 1) << 48,
+ .stale_on_death = true,
};
errdefer b.deinit();
b.data_buf = try mg.gpa.alloc(u8, max_write);
- // The kernel-visible root of this server's subtree: global node id
- // (index included), fid 0, its qid from the stat above. From here on
- // the bridge is an ordinary single-server bridge, just with node
- // ids that already carry the index.
- try b.inodes.put(mg.gpa, root_node, .{ .fid = 0, .qid = st.qid, .nlookup = 1, .parent = root_node });
- try b.by_qid.put(mg.gpa, st.qid.path, root_node);
- b.root_path = st.qid.path;
- b.next_node = root_node + 1;
var pipes: [2]i32 = undefined;
if (linux.errno(linux.pipe2(&pipes, .{ .CLOEXEC = true, .NONBLOCK = true })) != .SUCCESS) {
@@ -1716,30 +1764,66 @@ const Mntgen = struct {
m.* = .{
.gpa = mg.gpa,
.io = mg.io,
+ .mo = &mg.mo,
.debug = mg.opts.debug,
.name = try mg.gpa.dupe(u8, key),
.index = index,
.root_node = root_node,
- .root_attr = attrFromStat(st, inoFromPath(st.qid.path) ^ b.ino_xor, mg.opts.uid, mg.opts.gid),
- .session = session,
+ .parent_node = parent_node,
+ .entry_name = undefined, // a slice of `name`, set below
+ .stop_fd = mg.stop_fd,
.b = b,
.int_pipe = pipes,
.thread = undefined,
};
errdefer mg.gpa.free(m.name);
+ m.entry_name = lastComponent(m.name);
+ @memcpy(m.sock_buf[0..sock_path.len], sock_path);
+ m.sock_buf[sock_path.len] = 0;
+ m.sock_len = @intCast(sock_path.len);
- // The interrupt source must be installed before the thread starts.
- session.interrupt = mountInterrupt(m);
- m.thread = std.Thread.spawn(.{}, workerMain, .{m}) catch |e| {
- session.interrupt = null;
- return e;
- };
+ m.thread = try std.Thread.spawn(.{}, workerMain, .{m});
mg.mounts[index] = m;
mg.next_index += 1;
return m;
}
};
+/// A request over its copied bytes: the header, and the body after it.
+fn requestOf(raw: []align(8) u8) fuse.Request {
+ const header = std.mem.bytesToValue(fuse.InHeader, raw[0..@sizeOf(fuse.InHeader)]);
+ return .{ .header = header, .body = raw[@sizeOf(fuse.InHeader)..header.len] };
+}
+
+/// The FUSE `unique` of a raw request (InHeader bytes 8..16).
+fn uniqueOf(raw: []const u8) u64 {
+ return std.mem.readInt(u64, raw[8..16], .little);
+}
+
+/// Whether the kernel expects an answer to a raw request: everything but
+/// the forgets.
+fn wantsReply(raw: []const u8) bool {
+ return switch (@as(fuse.Opcode, @enumFromInt(std.mem.readInt(u32, raw[4..8], .little)))) {
+ .forget, .batch_forget => false,
+ else => true,
+ };
+}
+
+/// A LOOKUP answer for an entry of the synthetic tree: a mount's root or a
+/// synthetic directory. No entry caching for these: a walk re-LOOKUPs the
+/// name, which is what notices a dead server and makes a new mount. No
+/// invalidation machinery needed, and nothing to invalidate.
+fn entryOut(nodeid: u64, attr: fuse.Attr) fuse.EntryOut {
+ return .{ .nodeid = nodeid, .generation = 0, .attr = attr };
+}
+
+/// The last path component of a registry-relative key: the name the kernel
+/// caches the dentry under. "agents" and "sub/dir/agents" both give "agents".
+fn lastComponent(key: []u8) []u8 {
+ if (std.mem.lastIndexOfScalar(u8, key, '/')) |i| return key[i + 1 ..];
+ return key;
+}
+
/// The first live mount posted under `name`, skipping dead ones (a walk into
/// a name whose server died dials it afresh rather than reuse the corpse).
fn findMount(mounts: []const ?*Mount, name: []const u8) ?*Mount {
@@ -1751,10 +1835,11 @@ fn findMount(mounts: []const ?*Mount, name: []const u8) ?*Mount {
return null;
}
-/// A server's worker: pops copied requests off its queue and runs them
-/// through the ordinary bridge dispatch, one at a time (same concurrency
-/// contract as single-connection 9ns). The dispatcher keeps reading /dev/fuse
-/// meanwhile, so a slow server never blocks the other names.
+/// A mount's worker: takes requests off the queue and runs them through
+/// the ordinary bridge dispatch, one at a time (same concurrency contract
+/// as single-connection 9ns), dialing the server first if the mount has
+/// not been. The dispatcher keeps reading /dev/fuse meanwhile, so a slow
+/// server never blocks the other names.
fn workerMain(m: *Mount) void {
while (true) {
m.mutex.lockUncancelable(m.io);
@@ -1762,6 +1847,9 @@ fn workerMain(m: *Mount) void {
m.cond.waitUncancelable(m.io, &m.mutex);
}
const buf: ?[]align(8) u8 = if (m.queue.items.len != 0) m.queue.orderedRemove(0) else null;
+ // In flight from the moment it leaves the queue, under the same
+ // lock: an INTERRUPT finds it in one place or the other.
+ if (buf) |bytes| m.b.cur_unique.store(uniqueOf(bytes), .seq_cst);
m.mutex.unlock(m.io);
if (buf) |bytes| {
defer m.gpa.free(bytes);
@@ -1775,13 +1863,39 @@ fn workerMain(m: *Mount) void {
fn serveQueued(m: *Mount, bytes: []align(8) u8) void {
const header = std.mem.bytesToValue(fuse.InHeader, bytes[0..@sizeOf(fuse.InHeader)]);
const req = fuse.Request{ .header = header, .body = bytes[@sizeOf(fuse.InHeader)..header.len] };
- if (m.debug) std.debug.print("9ns: [{s}] <- {s} unique={d} nodeid={d}\n", .{ m.name, opName(req.header.op()), req.header.unique, req.header.nodeid });
- const keep_going = m.b.dispatch(req) catch {
+ const b = m.b;
+ defer b.cur_unique.store(0, .seq_cst);
+ b.interrupted = false;
+ if (m.debug) std.debug.print("9ns: [{s}] <- {s} unique={d} nodeid={d}\n", .{ m.name, opName(header.op()), header.unique, header.nodeid });
+ if (m.dead.load(.seq_cst)) {
+ // Queued behind the death: a stale handle, like everything the
+ // dispatcher answers for this mount from now on.
+ if (wantsReply(bytes)) b.replyError(header.unique, .STALE) catch {};
+ return;
+ }
+ if (m.session == null) {
+ dial(m) catch |e| {
+ // The walk that asked gets the verdict; the mount stays undialed
+ // and the next walk tries again.
+ if (m.debug) std.debug.print("9ns: [{s}] dial failed: {t}\n", .{ m.name, e });
+ if (wantsReply(bytes)) b.replyError(header.unique, dialErrno(e)) catch {};
+ return;
+ };
+ if (m.debug) std.debug.print("9ns: [{s}] dialed as mount {d}\n", .{ m.name, m.index });
+ }
+ // A LOOKUP whose node is not one of ours is the walk into our name
+ // (from the mntgen root or a synthetic directory): the server's root.
+ if (header.op() == .lookup and mountIndex(header.nodeid) != m.index) {
+ const out = entryOut(m.root_node, m.root_attr);
+ b.reply(header.unique, &.{std.mem.asBytes(&out)}) catch {};
+ return;
+ }
+ const keep_going = b.dispatch(req) catch {
// The 9P session (or the FUSE device) died mid-request: dispatch has
- // already replied EIO for this one. Everything still queued answers
- // EIO just as fast, and no new request is routed here again.
- m.dead.store(true, .seq_cst);
- if (m.debug) std.debug.print("9ns: [{s}] server connection lost; subtree now answers EIO\n", .{m.name});
+ // already replied for this one. Everything still queued answers
+ // ESTALE just as fast, and no new request is routed here again.
+ if (m.debug) std.debug.print("9ns: [{s}] server connection lost; subtree now answers ESTALE\n", .{m.name});
+ retire(m);
return;
};
if (!keep_going) {
@@ -1792,10 +1906,164 @@ fn serveQueued(m: *Mount, bytes: []align(8) u8) void {
}
}
+/// How long a walk waits for a server whose listen backlog is full before
+/// the walk answers EIO, and how often it looks for a stop or an
+/// interrupt meanwhile.
+const connect_grace_ms: i32 = 5000;
+const connect_retry_ms: i32 = 100;
+
+const DialError = error{ NotPosted, Stale, Busy, SystemResources, Io } || nine.Session.Error || std.mem.Allocator.Error;
+
+/// Connects, negotiates, attaches and stats the server root: the worker's
+/// half of a mount, run inside the walk that asked, which is the request
+/// in flight. The program's exit ends it (stop_fd), and so does an
+/// interrupt of that walk — at once and with nothing sent, since there is
+/// no session yet to flush a request out of (`abort_on_cancel`).
+fn dial(m: *Mount) DialError!void {
+ const fd = try connectSocket(m);
+ var fresh = nine.Session.dial(m.gpa, .{ .fd = fd }, m.mo.msize, m.stop_fd) catch |e| {
+ _ = linux.close(fd);
+ // With an `.fd` address the only thing left to fail is the buffers.
+ return switch (e) {
+ error.OutOfMemory => error.OutOfMemory,
+ else => error.Io,
+ };
+ };
+ fresh.interrupt = mountInterrupt(m);
+ fresh.abort_on_cancel = true;
+ m.mutex.lockUncancelable(m.io);
+ m.session = fresh;
+ m.mutex.unlock(m.io);
+ errdefer dropSession(m);
+ const sess = &m.session.?;
+ m.b.nine = sess;
+ try sess.version();
+ _ = try sess.attach(0, m.mo.uname, m.mo.aname);
+ const st = try sess.stat(0);
+ sess.abort_on_cancel = false;
+
+ // The kernel-visible root of this server's subtree: the mount's root
+ // node id, fid 0, its qid from the stat. From here on the bridge is an
+ // ordinary single-server bridge, just with node ids that already carry
+ // the index.
+ const b = m.b;
+ try b.inodes.put(b.gpa, m.root_node, .{ .fid = 0, .qid = st.qid, .nlookup = 1, .parent = m.root_node });
+ try b.by_qid.put(b.gpa, st.qid.path, m.root_node);
+ b.root_path = st.qid.path;
+ b.next_node = m.root_node + 1;
+ m.root_attr = attrFromStat(st, inoFromPath(st.qid.path) ^ b.ino_xor, b.opts.uid, b.opts.gid);
+}
+
+/// Closes the session and its socket; the mount is back to having none.
+fn dropSession(m: *Mount) void {
+ m.mutex.lockUncancelable(m.io);
+ if (m.session) |*sess| sess.deinit();
+ m.session = null;
+ m.mutex.unlock(m.io);
+}
+
+/// The mount is dead: its subtree answers ESTALE from now on and a walk
+/// into the name makes a new mount. Everything it held goes now, not at
+/// exit — the socket (a server that answers late must not fill a buffer
+/// nobody reads and block on it), the interrupt pipe (closed under the
+/// mutex, where the dispatcher writes it) and the bridge's tables and
+/// buffers — so a server that dies and comes back a thousand times costs
+/// a thousand slots, not a thousand descriptors and megabytes. The slot
+/// itself stays: indexes are never reused (see `max_mounts`).
+fn retire(m: *Mount) void {
+ m.dead.store(true, .seq_cst);
+ retireEntry(m);
+ dropSession(m);
+ m.mutex.lockUncancelable(m.io);
+ _ = linux.close(m.int_pipe[0]);
+ _ = linux.close(m.int_pipe[1]);
+ m.int_pipe = .{ -1, -1 };
+ m.mutex.unlock(m.io);
+ m.b.deinit();
+}
+
+/// Connects the mount's socket. A Unix stream connect completes on the
+/// spot, so a nonblocking one either succeeds, is refused, or reports EAGAIN
+/// when the server's backlog is full — a server that has stopped accepting.
+/// That case is retried for `connect_grace_ms`, watching stop_fd and the
+/// interrupt pipe in between, so a wedged server cannot hold the walk into
+/// its name past the walker's patience or the program's exit. The
+/// descriptor comes back blocking (the session reads it that way).
+fn connectSocket(m: *Mount) DialError!i32 {
+ const path = m.sockPath();
+ const rc = linux.socket(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK, 0);
+ if (linux.errno(rc) != .SUCCESS) return error.SystemResources;
+ const fd: i32 = @intCast(rc);
+ errdefer _ = linux.close(fd);
+ var addr: linux.sockaddr.un = .{ .path = @splat(0) };
+ @memcpy(addr.path[0..path.len], path);
+ var waited: i32 = 0;
+ while (true) {
+ switch (linux.errno(linux.connect(fd, @ptrCast(&addr), @sizeOf(linux.sockaddr.un)))) {
+ .SUCCESS => break,
+ .AGAIN => {},
+ .INTR => continue,
+ .NOENT, .NOTDIR => return error.NotPosted,
+ .CONNREFUSED => return error.Stale,
+ else => return error.Io,
+ }
+ if (waited >= connect_grace_ms) return error.Busy;
+ var pfds = [_]linux.pollfd{
+ .{ .fd = m.stop_fd, .events = linux.POLL.IN, .revents = 0 },
+ .{ .fd = m.int_pipe[0], .events = linux.POLL.IN, .revents = 0 },
+ };
+ switch (linux.errno(linux.poll(&pfds, pfds.len, connect_retry_ms))) {
+ .SUCCESS, .INTR, .AGAIN => {},
+ else => return error.Io,
+ }
+ if (pfds[0].revents != 0) return error.Stopped;
+ if (pfds[1].revents != 0 and try interruptHit(m)) return error.Interrupted;
+ waited += connect_retry_ms;
+ }
+ _ = linux.fcntl(fd, linux.F.SETFL, 0);
+ return fd;
+}
+
+/// The errno a walk gets when its dial fails.
+fn dialErrno(e: DialError) linux.E {
+ return switch (e) {
+ error.NotPosted => .NOENT,
+ error.Interrupted => .INTR,
+ error.OutOfMemory => .NOMEM,
+ // A stale entry, a full backlog, a server that will not speak 9P:
+ // the name is there and the server behind it is not. EIO, as for
+ // one that died.
+ else => .IO,
+ };
+}
+
+/// Tells the kernel to forget the dentry this dead mount was reached
+/// through, so the next access of the path re-LOOKUPs the name instead of
+/// reusing node ids that belong to the corpse.
+///
+/// Without this the walk still recovers, but only once the kernel's entry
+/// cache expires (`--cache`, 1s by default): until then every path under the
+/// name routes to the retired index and answers ESTALE, so a server restart
+/// surfaces as one spurious error to whoever touches the mount first.
+/// Invalidating the entry closes that window — the new mount happens inside
+/// the next LOOKUP, and the caller never sees the corpse.
+///
+/// Best effort, and deliberately not fatal: a notification the kernel
+/// rejects leaves exactly the old behaviour (ESTALE until the cache
+/// expires), which is degraded, not broken. Writing to /dev/fuse from this
+/// thread is safe — a notification carries `unique = 0`, and the workers
+/// already write their own replies to the same fd.
+fn retireEntry(m: *Mount) void {
+ fuse.notifyInvalEntry(m.b.fuse_fd, m.parent_node, m.entry_name) catch |e| {
+ if (m.debug) std.debug.print("9ns: [{s}] could not invalidate its entry: {t}\n", .{ m.name, e });
+ };
+}
+
// -- a mount's interrupt source -------------------------------------------------
//
-// The worker's session polls `int_pipe[0]` while a 9P reply is outstanding;
-// the dispatcher writes the interrupted request's unique into it. Unlike the
+// The worker's session polls `int_pipe[0]` while a 9P reply is outstanding
+// (and `connectSocket` while it waits on a full backlog); the dispatcher
+// writes the interrupted request's unique into it. Unlike the
// single-connection source, no FUSE reading happens here — the dispatcher
// owns /dev/fuse.
@@ -1810,6 +2078,12 @@ fn mountWatch(ctx: *anyopaque) i32 {
fn mountOnReadable(ctx: *anyopaque) nine.Session.Error!bool {
const m: *Mount = @ptrCast(@alignCast(ctx));
+ return interruptHit(m);
+}
+
+/// Drains the interrupt pipe; true when one of the packets named the
+/// request in flight, which is then marked interrupted.
+fn interruptHit(m: *Mount) nine.Session.Error!bool {
var packet: [8]u8 = undefined;
var hit = false;
while (true) {
@@ -1827,7 +2101,7 @@ fn mountOnReadable(ctx: *anyopaque) nine.Session.Error!bool {
}
}
if (hit) {
- if (m.debug) std.debug.print("9ns: [{s}] interrupt for unique={d} (in flight): sending Tflush\n", .{ m.name, m.b.cur_unique.load(.seq_cst) });
+ if (m.debug) std.debug.print("9ns: [{s}] interrupt for unique={d} (in flight): cancelling\n", .{ m.name, m.b.cur_unique.load(.seq_cst) });
m.b.interrupted = true;
return true;
}
@@ -2135,14 +2409,15 @@ test "interrupt source: INTERRUPT in flight → Tflush, other INTERRUPTs ignored
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);
+ try testing.expectEqual(@as(usize, 0), b.stash.items.len);
// The session is intact: a clunk-style cleanup rpc and a further read work.
b.interrupted = false;
b.cur_unique.store(8, .seq_cst);
try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40]));
try testing.expectEqual(@as(usize, 0), s.client.pending());
- // A FORGET arriving during a wait is stashed, and the fd is then not watched.
+ // A FORGET arriving during a wait is copied into the stash, and the fd
+ // stays watched: an INTERRUPT can still arrive behind it.
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));
@@ -2152,24 +2427,28 @@ test "interrupt source: INTERRUPT in flight → Tflush, other INTERRUPTs ignored
b.cur_unique.store(9, .seq_cst);
try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40]));
try testing.expect(!b.interrupted);
- const stashed = b.stash orelse return error.TestUnexpectedResult;
+ try testing.expectEqual(@as(usize, 1), b.stash.items.len);
+ const stashed = requestOf(b.stash.items[0]);
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 testing.expectEqual(pi.read_end, Bridge.interruptWatch(&b));
+ // With a request queued, an INTERRUPT for the one in flight is still
+ // consumed: the Tflush goes out and arms the operation even though the
+ // reply wins the race.
try pi.inject(9);
+ 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);
- 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).
+ try testing.expect(b.interrupted);
+ try testing.expectEqual(@as(u32, 2), fs.flushes.load(.seq_cst));
+ try testing.expectEqual(@as(usize, 1), b.stash.items.len);
+ // The late Rflush is swallowed by the next call, which is undisturbed.
+ b.interrupted = false;
b.cur_unique.store(10, .seq_cst);
fs.read_delay_ns = 30 * std.time.ns_per_ms;
try testing.expectEqual(@as(usize, 40), try b.read(1, 10, buf[0..40]));
try testing.expect(!b.interrupted);
- try testing.expect(b.stash == null);
+ try testing.expectEqual(@as(usize, 1), b.stash.items.len);
}
test "DirList frees its names" {
diff --git a/9ns/src/fuse.zig b/9ns/src/fuse.zig
index 216e616..e45abc2 100644
--- a/9ns/src/fuse.zig
+++ b/9ns/src/fuse.zig
@@ -37,6 +37,12 @@ pub const FUSE_BIG_WRITES: u32 = 1 << 5;
/// Required for `--no-direct-io` correctness: 9P sizes change under us, and
/// without this the kernel trusts a stale cached size and truncates reads.
pub const FUSE_AUTO_INVAL_DATA: u32 = 1 << 12;
+/// Without this the kernel takes the directory's inode lock around every
+/// LOOKUP and READDIR in it, so one walk parked on a server that never
+/// answers holds up every other name in the same directory — the whole
+/// registry root, for a mntgen mount. With it, walks into different names
+/// proceed side by side and only the parked one waits.
+pub const FUSE_PARALLEL_DIROPS: u32 = 1 << 18;
pub const FUSE_MAX_PAGES: u32 = 1 << 22;
pub const FATTR_MODE: u32 = 1 << 0;
@@ -389,6 +395,44 @@ pub fn replyError(fd: i32, unique: u64, err: linux.E) Error!void {
return writeAll(fd, &iov, 1, @sizeOf(OutHeader));
}
+/// Notification codes, sent to the kernel unsolicited: they ride an
+/// `OutHeader` with `unique = 0` and the code (positive) in `error`.
+pub const notify_inval_entry: i32 = 3;
+
+/// The body of a `FUSE_NOTIFY_INVAL_ENTRY`, followed by the name and a NUL.
+pub const NotifyInvalEntryOut = extern struct {
+ parent: u64,
+ namelen: u32,
+ padding: u32 = 0,
+};
+
+/// Drops the kernel's cached dentry for `name` under `parent`, so the next
+/// access of that path comes back as a fresh LOOKUP instead of reusing a
+/// node id the server no longer knows.
+///
+/// `unique = 0` marks the message as a notification rather than a reply, so
+/// it is safe to interleave with replies on the same fd, from any thread.
+///
+/// Best-effort by nature: the kernel answers ENOENT when it had nothing
+/// cached under that name (`writeAll` already treats that as success) and
+/// EINVAL when it does not support the notification, which surfaces here as
+/// `error.Io`. Callers ignore both — a notification that does not land
+/// leaves them exactly where they were without it.
+pub fn notifyInvalEntry(fd: i32, parent: u64, name: []const u8) Error!void {
+ if (name.len == 0 or name.len > std.math.maxInt(u32) - 1) return error.Protocol;
+ const out = NotifyInvalEntryOut{ .parent = parent, .namelen = @intCast(name.len) };
+ const total = @sizeOf(OutHeader) + @sizeOf(NotifyInvalEntryOut) + name.len + 1;
+ const header = OutHeader{ .len = @intCast(total), .@"error" = notify_inval_entry, .unique = 0 };
+ const nul = [_]u8{0};
+ var iov = [_]std.posix.iovec_const{
+ .{ .base = @ptrCast(&header), .len = @sizeOf(OutHeader) },
+ .{ .base = @ptrCast(&out), .len = @sizeOf(NotifyInvalEntryOut) },
+ .{ .base = name.ptr, .len = name.len },
+ .{ .base = &nul, .len = 1 },
+ };
+ return writeAll(fd, &iov, iov.len, total);
+}
+
fn writeAll(fd: i32, iov: [*]const std.posix.iovec_const, count: usize, total: usize) Error!void {
while (true) {
const rc = linux.writev(fd, iov, count);
@@ -462,6 +506,7 @@ pub fn initReply(in: *const InitIn, max_write: u32) InitOut {
out.flags |= FUSE_MAX_PAGES;
out.max_pages = 256;
}
+ if (in.flags & FUSE_PARALLEL_DIROPS != 0) out.flags |= FUSE_PARALLEL_DIROPS;
return out;
}
diff --git a/9ns/src/nine.zig b/9ns/src/nine.zig
index cf9c49a..6f1fe09 100644
--- a/9ns/src/nine.zig
+++ b/9ns/src/nine.zig
@@ -71,6 +71,20 @@ pub const Session = struct {
/// Optional interrupt source (the bridge's FUSE descriptor) consulted while
/// a reply is outstanding; see `Interrupt`.
interrupt: ?Interrupt = null,
+ /// While set, a cancellation ends the rpc at once with `error.Interrupted`
+ /// and wedges the session, with no Tflush: for the handshake (version,
+ /// attach, the root stat), where there is no session yet to flush a
+ /// request out of, and the honest answer to "stop waiting" is to hang up.
+ abort_on_cancel: bool = false,
+ /// Milliseconds a Tflush may go unanswered before the server is declared
+ /// wedged: the rpc fails with `error.Interrupted` and the session with it.
+ /// The protocol says a client waits for the Rflush; a server that has not
+ /// managed one in this long is not going to, and the process behind the
+ /// interrupt is unkillable until we stop waiting. 0 waits forever.
+ flush_grace_ms: i32 = 3000,
+ /// The server is gone as far as this session is concerned (see
+ /// `abort_on_cancel`, `flush_grace_ms`); every rpc answers `error.Closed`.
+ wedged: bool = false,
/// Connect to `address`, then negotiate the protocol version.
/// `msize` is the maximum message size to ask for (0 = the buffers' size).
@@ -81,23 +95,33 @@ pub const Session = struct {
/// `connect`, with `stop_fd` watched for the whole handshake (the version
/// rpc included). A server that accepts the connection and then never
/// answers the Tversion would otherwise pin the caller in a blocking read
- /// with no way out: 9ns's mntgen dispatcher dials on the strength of the
- /// program's walk, so it must come back when that program is gone. The
- /// field stays set on the returned session, so the attach and stat that
- /// follow a dial keep watching it too; -1 disables the watch.
+ /// with no way out; -1 disables the watch. The field stays set on the
+ /// returned session. An `.fd` address is the caller's to close on failure.
pub fn connectWatched(gpa: std.mem.Allocator, address: Address, msize: u32, stop_fd: i32) !Session {
+ var s = try dial(gpa, address, msize, stop_fd);
+ s.version() catch |e| {
+ s.freeBuffers();
+ if (address != .fd) _ = linux.close(s.fd);
+ return e;
+ };
+ return s;
+ }
+
+ /// The transport and the buffers, no handshake: for a caller that wants
+ /// its interrupt source in place before `version()` (9ns's mntgen
+ /// worker, so a walk interrupted mid-dial can abandon the dial). Owns
+ /// the descriptor from here: `deinit` closes it.
+ pub fn dial(gpa: std.mem.Allocator, address: Address, msize: u32, stop_fd: i32) !Session {
const want: u32 = if (msize == 0) 8192 else @max(msize, 24);
const fd = try openTransport(address);
errdefer if (address != .fd) {
_ = linux.close(fd);
};
-
const in_buf = try gpa.alloc(u8, want);
errdefer gpa.free(in_buf);
const out_buf = try gpa.alloc(u8, want);
errdefer gpa.free(out_buf);
-
- var s: Session = .{
+ return .{
.gpa = gpa,
.fd = fd,
.client = .init(.{ .in = in_buf, .out = out_buf }),
@@ -106,19 +130,26 @@ pub const Session = struct {
.msize = want,
.stop_fd = stop_fd,
};
- const r = try s.rpc(.{ .version = .{ .msize = want } });
+ }
+
+ /// Negotiates the protocol version with the msize `dial` was given.
+ pub fn version(s: *Session) Error!void {
+ const r = try s.rpc(.{ .version = .{ .msize = s.msize } });
if (!std.mem.eql(u8, r.version.version, "9P2000")) return error.Protocol;
s.msize = r.version.msize;
- return s;
}
- /// Closes the descriptor and frees the buffers. Fids are not clunked.
- pub fn deinit(s: *Session) void {
- _ = linux.close(s.fd);
+ fn freeBuffers(s: *Session) void {
s.free_fids.deinit(s.gpa);
s.iounits.deinit(s.gpa);
s.gpa.free(s.in_buf);
s.gpa.free(s.out_buf);
+ }
+
+ /// Closes the descriptor and frees the buffers. Fids are not clunked.
+ pub fn deinit(s: *Session) void {
+ _ = linux.close(s.fd);
+ s.freeBuffers();
s.* = undefined;
}
@@ -154,8 +185,12 @@ pub const Session = struct {
/// 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).
+ /// `error.Interrupted`; the server has dropped the request), or neither
+ /// within `flush_grace_ms` (→ `error.Interrupted`, and the session is
+ /// wedged: the server stopped talking). Under `abort_on_cancel` the
+ /// cancellation itself wedges the session, with nothing sent.
pub fn rpc(s: *Session, req: cloud9.Client.Request) Error!cloud9.Client.Result {
+ if (s.wedged) return error.Closed;
s.ename_len = 0;
const tag = s.client.submit(req) catch |e| switch (e) {
error.NoTags, error.Handshake, error.Dead => return error.Protocol,
@@ -167,6 +202,8 @@ pub const Session = struct {
};
try s.flush();
var flush_tag: ?u16 = null;
+ // Monotonic ms by which the Tflush must have been answered.
+ var flush_deadline: ?i64 = null;
var tmp: [64 * 1024]u8 = undefined;
while (true) {
while (s.client.take()) |done| {
@@ -190,16 +227,28 @@ pub const Session = struct {
// 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;
- switch (try s.wait()) {
+ const grace: i32 = if (flush_deadline) |d| @intCast(@max(d - nowMs(), 0)) else -1;
+ switch (try s.wait(grace)) {
.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();
+ .cancel => {
+ if (s.abort_on_cancel) {
+ s.wedged = true;
+ return error.Interrupted;
+ }
+ if (flush_tag == null) {
+ flush_tag = s.client.submit(.{ .flush = .{ .oldtag = tag } }) catch return error.Protocol;
+ try s.flush();
+ if (s.flush_grace_ms != 0) flush_deadline = nowMs() + s.flush_grace_ms;
+ }
+ },
+ .timeout => {
+ s.wedged = true;
+ return error.Interrupted;
},
}
}
@@ -306,13 +355,14 @@ pub const Session = struct {
return chunkSize(s.client.maxWrite(), s.iounits.get(fid) orelse 0);
}
- const Ready = enum { socket, cancel };
+ const Ready = enum { socket, cancel, timeout };
/// Blocks until the socket is readable (`.socket`), the interrupt source
- /// wants the request in flight cancelled (`.cancel`), or `stop_fd` fires
+ /// wants the request in flight cancelled (`.cancel`), `timeout_ms` passes
+ /// with neither (`.timeout`; -1 waits forever), 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 {
+ fn wait(s: *Session, timeout_ms: i32) Error!Ready {
while (true) {
var pfds: [3]linux.pollfd = undefined;
var n: usize = 0;
@@ -329,13 +379,14 @@ pub const Session = struct {
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);
+ if (n == 1 and timeout_ms < 0) return .socket;
+ const prc = linux.poll(&pfds, @intCast(n), timeout_ms);
switch (linux.errno(prc)) {
.SUCCESS => {},
.INTR, .AGAIN => continue,
else => return error.Io,
}
+ if (prc == 0) return .timeout;
// A reply that is already there wins over everything else.
if (pfds[0].revents != 0) return .socket;
if (stop_at) |i| {
@@ -368,6 +419,13 @@ pub const Session = struct {
}
};
+/// The monotonic clock in milliseconds: deadlines, not timestamps.
+fn nowMs() i64 {
+ var ts: linux.timespec = undefined;
+ _ = linux.clock_gettime(.MONOTONIC, &ts);
+ return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);
+}
+
fn chunkSize(max: u32, iounit: u32) u32 {
if (iounit != 0 and iounit < max) return iounit;
return max;
@@ -723,6 +781,62 @@ test "interrupted chunk loops: partial count if data moved, Interrupted otherwis
// -- in-process server test ---------------------------------------------------------
+test "flush grace: a Tflush the server never answers wedges the session after the grace" {
+ var pi = try PipeInterrupt.init(7);
+ defer pi.deinit();
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .ignore, .on_hang_inject = &pi };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ s.interrupt = pi.interface();
+ s.flush_grace_ms = 100;
+ // The read at offset 0 hangs; the injected INTERRUPT sends a Tflush; the
+ // server ignores it; the grace runs out.
+ var buf: [100]u8 = undefined;
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ try testing.expect(s.wedged);
+ // From here on the session is closed for business, without another
+ // byte to the server.
+ try testing.expectError(error.Closed, s.read(1, 10, &buf));
+ try testing.expectError(error.Closed, s.clunk(1));
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ // The socket is still the session's to close: the server's thread ends
+ // when `close` drops it.
+}
+
+test "flush grace: zero waits for the Rflush, however late" {
+ var pi = try PipeInterrupt.init(7);
+ defer pi.deinit();
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush, .on_hang_inject = &pi, .read_delay_ns = 0 };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ s.interrupt = pi.interface();
+ s.flush_grace_ms = 0;
+ var buf: [100]u8 = undefined;
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expect(!s.wedged);
+ // The session lives: the Rflush released the tag and reads go on.
+ try testing.expectEqual(@as(usize, 40), try s.read(1, 10, buf[0..40]));
+}
+
+test "abort_on_cancel: a cancellation ends the rpc at once, sends no Tflush, wedges the session" {
+ var pi = try PipeInterrupt.init(7);
+ defer pi.deinit();
+ var fs: 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;
+ s.interrupt = pi.interface();
+ s.abort_on_cancel = true;
+ var buf: [100]u8 = undefined;
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expect(s.wedged);
+ try testing.expectEqual(@as(u32, 0), fs.flushes.load(.seq_cst));
+ try testing.expectError(error.Closed, s.read(1, 10, &buf));
+}
+
/// 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.
/// Test support only (bridge.zig's tests use it too).
@@ -734,9 +848,11 @@ pub const FakeServer = struct {
/// 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,
+ /// What a Tflush gets: the connection dropped (`.hangup`), an Rflush,
+ /// first the Rread the flush was aimed at and then the Rflush (the
+ /// race), or nothing at all (`.ignore`: counted, never answered, the
+ /// read stays hung — a wedged server).
+ on_flush: enum { hangup, rflush, reply_then_rflush, ignore } = .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),
@@ -827,6 +943,7 @@ pub const FakeServer = struct {
switch (fs.on_flush) {
// The test's "hang up now" signal.
.hangup => return,
+ .ignore => {},
.rflush => {
if (hung != null and hung.?.tag == m.oldtag) hung = null;
try srv.reply(tag, .rflush);
diff --git a/9ns/test/adv_bridge_interrupt.sh b/9ns/test/adv_bridge_interrupt.sh
index e364143..ba76f78 100755
--- a/9ns/test/adv_bridge_interrupt.sh
+++ b/9ns/test/adv_bridge_interrupt.sh
@@ -118,18 +118,20 @@ 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 "# (b) a server that ignores Tflush too: the reader comes back after the flush grace, and the session is gone with it"
+# `never` stops reading once the Tread hangs, so the Tflush is never even
+# seen. The protocol says wait for the Rflush; a server that has not managed
+# one in 3s (nine.Session.flush_grace_ms) is not going to, and the reader
+# behind the interrupt is unkillable until we stop waiting. So the read fails
+# EINTR after the grace, the session is wedged, and 9ns ends the mount: the
+# next access answers ENOTCONN instead of parking another process forever.
+run never "$(timed "timeout -s INT 1 cat $M/f"); cat $M/d/g 2>&1; echo after-rc=\$?"
+no_crash "never/grace"
+expect_eq "never/grace: reader killed by the signal (124)" "124" "$(field rc "$OUT")"
+NEVER_MS=$(field ms "$OUT")
+if [ "$NEVER_MS" -ge 3500 ] && [ "$NEVER_MS" -lt 9000 ]; then pass "never/grace: released after the 3s flush grace (${NEVER_MS}ms)"; else fail "never/grace: release time out of range" "ms=$NEVER_MS (expected 3500..9000)"; fi
+expect_eq "never/grace: the mount is gone afterwards (not a hang)" "after-rc=1" "$(printf '%s\n' "$OUT" | grep '^after-rc=')"
+expect_contains "never/grace: 9ns reports the closed session" "connection closed" "$STDERR"
echo
echo "passed=$PASSED failed=$FAILED"
diff --git a/9ns/test/mntgen.sh b/9ns/test/mntgen.sh
index aade76c..7dada53 100755
--- a/9ns/test/mntgen.sh
+++ b/9ns/test/mntgen.sh
@@ -254,14 +254,18 @@ echo "# --debug and --no-direct-io reach the mntgen dispatcher"
DEBUG_ERR=$(timeout 60 "$NS" --mntgen --debug -- sh -c "ls $M/alpha >/dev/null" 2>&1 >/dev/null)
expect_contains "--debug traces the dispatcher" "lookup 'alpha'" "$DEBUG_ERR"
-echo "# a dial parked on a mute server unwedges when the program exits"
-# A server that accepts the connection but never answers Tversion parks the
-# dispatcher in the dial (pinned: no dial timeout, no concurrent dial). The
-# dial must watch stop_fd through the whole handshake: when the program's
-# main flow exits while a background walk is parked there, 9ns must follow it
-# out instead of wedging forever (pre-fix it survived SIGTERM).
+echo "# a mute server costs only the walks into its own name"
+# A server that accepts the connection but never answers Tversion. The dial
+# runs on that name's worker, inside the walk that asked, so: every other
+# name (and the root) keeps answering; a signal releases the parked walker
+# at once (the dial is abandoned, EINTR); a fresh walk parks again and is
+# just as interruptible; and when the program exits with a walk still
+# parked, 9ns follows it out. Background jobs of a non-interactive sh ignore
+# SIGINT, so the parked walker is sent SIGTERM; the foreground `timeout -s
+# INT` case covers Ctrl-C. Pre-fix the dial ran on the dispatcher thread and
+# parked the whole mount, unkillably, until the program exited.
if command -v python3 > /dev/null; then
- python3 - "$REG/mute" <<'PYEOF' &
+ python3 - "$REG/mute" <<'MUTEEOF' &
import socket, sys, os
path = sys.argv[1]
try:
@@ -273,22 +277,45 @@ s.bind(path)
s.listen(8)
while True:
conn, _ = s.accept() # accept, then never say a word
-PYEOF
+MUTEEOF
MUTEPID=$!
PIDS+=($MUTEPID)
sleep 0.3
+ MUTE_OUT=$(timeout 60 "$NS" --mntgen -- sh -c '
+ stat '"$M"'/mute >/dev/null 2>&1 & W=$!
+ sleep 0.3
+ echo "alpha=$(timeout 5 cat '"$M"'/alpha/build/zig_version)"
+ echo "root=$(timeout 5 ls '"$M"' | grep -c .)"
+ echo "readdir_alpha=$(timeout 5 ls '"$M"'/alpha | grep -c .)"
+ s=$(date +%s%N); kill -TERM $W; wait $W; echo "term_rc=$? term_ms=$(( ($(date +%s%N) - s) / 1000000 ))"
+ s=$(date +%s%N); timeout -s INT 2 stat '"$M"'/mute >/dev/null 2>&1; echo "int_rc=$? int_ms=$(( ($(date +%s%N) - s) / 1000000 ))"
+ echo "alpha_again=$(timeout 5 cat '"$M"'/alpha/build/zig_version)"
+ ' 2>"$TMP/mute.err")
+ expect_eq "another name is served while a walk into mute is parked" "alpha=$(zig version)" "$(printf '%s\n' "$MUTE_OUT" | grep '^alpha=')"
+ ROOT_N=$(printf '%s\n' "$MUTE_OUT" | sed -n 's/^root=//p')
+ [ -n "$ROOT_N" ] && [ "$ROOT_N" -ge 2 ] && pass "the root lists while a walk into mute is parked ($ROOT_N entries)" || fail "the root listing did not answer" "$MUTE_OUT"
+ READDIR_N=$(printf '%s\n' "$MUTE_OUT" | sed -n 's/^readdir_alpha=//p')
+ [ -n "$READDIR_N" ] && [ "$READDIR_N" -ge 1 ] && pass "a readdir on another name is served meanwhile" || fail "readdir on alpha did not answer" "$MUTE_OUT"
+ expect_eq "SIGTERM releases the parked walker (died of the signal)" "term_rc=143" "$(printf '%s\n' "$MUTE_OUT" | sed -n 's/^\(term_rc=[0-9]*\) .*/\1/p')"
+ TERM_MS=$(printf '%s\n' "$MUTE_OUT" | sed -n 's/.*term_ms=//p')
+ [ -n "$TERM_MS" ] && [ "$TERM_MS" -lt 1500 ] && pass "released promptly (${TERM_MS}ms)" || fail "release of the parked walker was slow or missing" "term_ms=$TERM_MS"
+ expect_eq "a fresh walk into mute is interruptible (Ctrl-C after 2s)" "int_rc=124" "$(printf '%s\n' "$MUTE_OUT" | sed -n 's/^\(int_rc=[0-9]*\) .*/\1/p')"
+ INT_MS=$(printf '%s\n' "$MUTE_OUT" | sed -n 's/.*int_ms=//p')
+ [ -n "$INT_MS" ] && [ "$INT_MS" -ge 1900 ] && [ "$INT_MS" -lt 4000 ] && pass "the interrupted walk came back on the signal (${INT_MS}ms)" || fail "interrupted walk timing off" "int_ms=$INT_MS"
+ expect_eq "alpha still served after all that" "alpha_again=$(zig version)" "$(printf '%s\n' "$MUTE_OUT" | grep '^alpha_again=')"
+
HANG_START=$(date +%s)
timeout 20 "$NS" --mntgen -- bash -c "(stat $M/mute >/dev/null 2>&1) & sleep 1" >/dev/null 2>&1
HANG_RC=$?
HANG_SECONDS=$(( $(date +%s) - HANG_START ))
if [ "$HANG_RC" -eq 124 ] || [ "$HANG_SECONDS" -ge 15 ]; then
- fail "9ns unwedges after the program exits a parked dial" "rc=$HANG_RC after ${HANG_SECONDS}s (wedged)"
+ fail "9ns exits with a walk still parked in a dial" "rc=$HANG_RC after ${HANG_SECONDS}s (wedged)"
else
- pass "9ns unwedges after the program exits a parked dial (rc=$HANG_RC after ${HANG_SECONDS}s)"
+ pass "9ns exits with a walk still parked in a dial (rc=$HANG_RC after ${HANG_SECONDS}s)"
fi
rm -f "$REG/mute"
else
- echo "SKIP: mute-server dial-hang check needs python3"
+ echo "SKIP: mute-server checks need python3"
fi
echo