From b7fc01550c7bde290cf14276d94193b5b4031dc8 Mon Sep 17 00:00:00 2001 From: Gabriel Schneider Date: Tue, 22 Sep 2026 11:18:05 -0300 Subject: 9ns --mntgen: a server that never answers stalls only its own name MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- 9ns/src/bridge.zig | 647 ++++++++++++++++++++++++++++++++++++++--------------- 1 file changed, 463 insertions(+), 184 deletions(-) (limited to '9ns/src/bridge.zig') 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" { -- cgit v1.3