diff options
Diffstat (limited to '9ns/src/bridge.zig')
| -rw-r--r-- | 9ns/src/bridge.zig | 647 |
1 files changed, 463 insertions, 184 deletions
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" { |
