summaryrefslogtreecommitdiff
path: root/9ns/src
diff options
context:
space:
mode:
Diffstat (limited to '9ns/src')
-rw-r--r--9ns/src/bridge.zig647
-rw-r--r--9ns/src/fuse.zig45
-rw-r--r--9ns/src/nine.zig167
3 files changed, 650 insertions, 209 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" {
diff --git a/9ns/src/fuse.zig b/9ns/src/fuse.zig
index 216e616..e45abc2 100644
--- a/9ns/src/fuse.zig
+++ b/9ns/src/fuse.zig
@@ -37,6 +37,12 @@ pub const FUSE_BIG_WRITES: u32 = 1 << 5;
/// Required for `--no-direct-io` correctness: 9P sizes change under us, and
/// without this the kernel trusts a stale cached size and truncates reads.
pub const FUSE_AUTO_INVAL_DATA: u32 = 1 << 12;
+/// Without this the kernel takes the directory's inode lock around every
+/// LOOKUP and READDIR in it, so one walk parked on a server that never
+/// answers holds up every other name in the same directory — the whole
+/// registry root, for a mntgen mount. With it, walks into different names
+/// proceed side by side and only the parked one waits.
+pub const FUSE_PARALLEL_DIROPS: u32 = 1 << 18;
pub const FUSE_MAX_PAGES: u32 = 1 << 22;
pub const FATTR_MODE: u32 = 1 << 0;
@@ -389,6 +395,44 @@ pub fn replyError(fd: i32, unique: u64, err: linux.E) Error!void {
return writeAll(fd, &iov, 1, @sizeOf(OutHeader));
}
+/// Notification codes, sent to the kernel unsolicited: they ride an
+/// `OutHeader` with `unique = 0` and the code (positive) in `error`.
+pub const notify_inval_entry: i32 = 3;
+
+/// The body of a `FUSE_NOTIFY_INVAL_ENTRY`, followed by the name and a NUL.
+pub const NotifyInvalEntryOut = extern struct {
+ parent: u64,
+ namelen: u32,
+ padding: u32 = 0,
+};
+
+/// Drops the kernel's cached dentry for `name` under `parent`, so the next
+/// access of that path comes back as a fresh LOOKUP instead of reusing a
+/// node id the server no longer knows.
+///
+/// `unique = 0` marks the message as a notification rather than a reply, so
+/// it is safe to interleave with replies on the same fd, from any thread.
+///
+/// Best-effort by nature: the kernel answers ENOENT when it had nothing
+/// cached under that name (`writeAll` already treats that as success) and
+/// EINVAL when it does not support the notification, which surfaces here as
+/// `error.Io`. Callers ignore both — a notification that does not land
+/// leaves them exactly where they were without it.
+pub fn notifyInvalEntry(fd: i32, parent: u64, name: []const u8) Error!void {
+ if (name.len == 0 or name.len > std.math.maxInt(u32) - 1) return error.Protocol;
+ const out = NotifyInvalEntryOut{ .parent = parent, .namelen = @intCast(name.len) };
+ const total = @sizeOf(OutHeader) + @sizeOf(NotifyInvalEntryOut) + name.len + 1;
+ const header = OutHeader{ .len = @intCast(total), .@"error" = notify_inval_entry, .unique = 0 };
+ const nul = [_]u8{0};
+ var iov = [_]std.posix.iovec_const{
+ .{ .base = @ptrCast(&header), .len = @sizeOf(OutHeader) },
+ .{ .base = @ptrCast(&out), .len = @sizeOf(NotifyInvalEntryOut) },
+ .{ .base = name.ptr, .len = name.len },
+ .{ .base = &nul, .len = 1 },
+ };
+ return writeAll(fd, &iov, iov.len, total);
+}
+
fn writeAll(fd: i32, iov: [*]const std.posix.iovec_const, count: usize, total: usize) Error!void {
while (true) {
const rc = linux.writev(fd, iov, count);
@@ -462,6 +506,7 @@ pub fn initReply(in: *const InitIn, max_write: u32) InitOut {
out.flags |= FUSE_MAX_PAGES;
out.max_pages = 256;
}
+ if (in.flags & FUSE_PARALLEL_DIROPS != 0) out.flags |= FUSE_PARALLEL_DIROPS;
return out;
}
diff --git a/9ns/src/nine.zig b/9ns/src/nine.zig
index cf9c49a..6f1fe09 100644
--- a/9ns/src/nine.zig
+++ b/9ns/src/nine.zig
@@ -71,6 +71,20 @@ pub const Session = struct {
/// Optional interrupt source (the bridge's FUSE descriptor) consulted while
/// a reply is outstanding; see `Interrupt`.
interrupt: ?Interrupt = null,
+ /// While set, a cancellation ends the rpc at once with `error.Interrupted`
+ /// and wedges the session, with no Tflush: for the handshake (version,
+ /// attach, the root stat), where there is no session yet to flush a
+ /// request out of, and the honest answer to "stop waiting" is to hang up.
+ abort_on_cancel: bool = false,
+ /// Milliseconds a Tflush may go unanswered before the server is declared
+ /// wedged: the rpc fails with `error.Interrupted` and the session with it.
+ /// The protocol says a client waits for the Rflush; a server that has not
+ /// managed one in this long is not going to, and the process behind the
+ /// interrupt is unkillable until we stop waiting. 0 waits forever.
+ flush_grace_ms: i32 = 3000,
+ /// The server is gone as far as this session is concerned (see
+ /// `abort_on_cancel`, `flush_grace_ms`); every rpc answers `error.Closed`.
+ wedged: bool = false,
/// Connect to `address`, then negotiate the protocol version.
/// `msize` is the maximum message size to ask for (0 = the buffers' size).
@@ -81,23 +95,33 @@ pub const Session = struct {
/// `connect`, with `stop_fd` watched for the whole handshake (the version
/// rpc included). A server that accepts the connection and then never
/// answers the Tversion would otherwise pin the caller in a blocking read
- /// with no way out: 9ns's mntgen dispatcher dials on the strength of the
- /// program's walk, so it must come back when that program is gone. The
- /// field stays set on the returned session, so the attach and stat that
- /// follow a dial keep watching it too; -1 disables the watch.
+ /// with no way out; -1 disables the watch. The field stays set on the
+ /// returned session. An `.fd` address is the caller's to close on failure.
pub fn connectWatched(gpa: std.mem.Allocator, address: Address, msize: u32, stop_fd: i32) !Session {
+ var s = try dial(gpa, address, msize, stop_fd);
+ s.version() catch |e| {
+ s.freeBuffers();
+ if (address != .fd) _ = linux.close(s.fd);
+ return e;
+ };
+ return s;
+ }
+
+ /// The transport and the buffers, no handshake: for a caller that wants
+ /// its interrupt source in place before `version()` (9ns's mntgen
+ /// worker, so a walk interrupted mid-dial can abandon the dial). Owns
+ /// the descriptor from here: `deinit` closes it.
+ pub fn dial(gpa: std.mem.Allocator, address: Address, msize: u32, stop_fd: i32) !Session {
const want: u32 = if (msize == 0) 8192 else @max(msize, 24);
const fd = try openTransport(address);
errdefer if (address != .fd) {
_ = linux.close(fd);
};
-
const in_buf = try gpa.alloc(u8, want);
errdefer gpa.free(in_buf);
const out_buf = try gpa.alloc(u8, want);
errdefer gpa.free(out_buf);
-
- var s: Session = .{
+ return .{
.gpa = gpa,
.fd = fd,
.client = .init(.{ .in = in_buf, .out = out_buf }),
@@ -106,19 +130,26 @@ pub const Session = struct {
.msize = want,
.stop_fd = stop_fd,
};
- const r = try s.rpc(.{ .version = .{ .msize = want } });
+ }
+
+ /// Negotiates the protocol version with the msize `dial` was given.
+ pub fn version(s: *Session) Error!void {
+ const r = try s.rpc(.{ .version = .{ .msize = s.msize } });
if (!std.mem.eql(u8, r.version.version, "9P2000")) return error.Protocol;
s.msize = r.version.msize;
- return s;
}
- /// Closes the descriptor and frees the buffers. Fids are not clunked.
- pub fn deinit(s: *Session) void {
- _ = linux.close(s.fd);
+ fn freeBuffers(s: *Session) void {
s.free_fids.deinit(s.gpa);
s.iounits.deinit(s.gpa);
s.gpa.free(s.in_buf);
s.gpa.free(s.out_buf);
+ }
+
+ /// Closes the descriptor and frees the buffers. Fids are not clunked.
+ pub fn deinit(s: *Session) void {
+ _ = linux.close(s.fd);
+ s.freeBuffers();
s.* = undefined;
}
@@ -154,8 +185,12 @@ pub const Session = struct {
/// and the wait continues until either the original reply arrives (the flush
/// lost the race; the result is returned as if nothing happened and the
/// Rflush is swallowed by a later call) or the Rflush does (→
- /// `error.Interrupted`; the server has dropped the request).
+ /// `error.Interrupted`; the server has dropped the request), or neither
+ /// within `flush_grace_ms` (→ `error.Interrupted`, and the session is
+ /// wedged: the server stopped talking). Under `abort_on_cancel` the
+ /// cancellation itself wedges the session, with nothing sent.
pub fn rpc(s: *Session, req: cloud9.Client.Request) Error!cloud9.Client.Result {
+ if (s.wedged) return error.Closed;
s.ename_len = 0;
const tag = s.client.submit(req) catch |e| switch (e) {
error.NoTags, error.Handshake, error.Dead => return error.Protocol,
@@ -167,6 +202,8 @@ pub const Session = struct {
};
try s.flush();
var flush_tag: ?u16 = null;
+ // Monotonic ms by which the Tflush must have been answered.
+ var flush_deadline: ?i64 = null;
var tmp: [64 * 1024]u8 = undefined;
while (true) {
while (s.client.take()) |done| {
@@ -190,16 +227,28 @@ pub const Session = struct {
// space is at least what the pending frame still needs.
const room = s.client.in.len - s.client.in_len;
if (room == 0) return error.Protocol;
- switch (try s.wait()) {
+ const grace: i32 = if (flush_deadline) |d| @intCast(@max(d - nowMs(), 0)) else -1;
+ switch (try s.wait(grace)) {
.socket => {
const n = try readSocket(s.fd, tmp[0..@min(room, tmp.len)]);
if (n == 0) return error.Closed;
const pushed = s.client.push(tmp[0..n]);
if (pushed != n) return error.Protocol;
},
- .cancel => if (flush_tag == null) {
- flush_tag = s.client.submit(.{ .flush = .{ .oldtag = tag } }) catch return error.Protocol;
- try s.flush();
+ .cancel => {
+ if (s.abort_on_cancel) {
+ s.wedged = true;
+ return error.Interrupted;
+ }
+ if (flush_tag == null) {
+ flush_tag = s.client.submit(.{ .flush = .{ .oldtag = tag } }) catch return error.Protocol;
+ try s.flush();
+ if (s.flush_grace_ms != 0) flush_deadline = nowMs() + s.flush_grace_ms;
+ }
+ },
+ .timeout => {
+ s.wedged = true;
+ return error.Interrupted;
},
}
}
@@ -306,13 +355,14 @@ pub const Session = struct {
return chunkSize(s.client.maxWrite(), s.iounits.get(fid) orelse 0);
}
- const Ready = enum { socket, cancel };
+ const Ready = enum { socket, cancel, timeout };
/// Blocks until the socket is readable (`.socket`), the interrupt source
- /// wants the request in flight cancelled (`.cancel`), or `stop_fd` fires
+ /// wants the request in flight cancelled (`.cancel`), `timeout_ms` passes
+ /// with neither (`.timeout`; -1 waits forever), or `stop_fd` fires
/// (`error.Stopped`). Anything the interrupt source consumes without asking
/// for a cancellation simply resumes the wait.
- fn wait(s: *Session) Error!Ready {
+ fn wait(s: *Session, timeout_ms: i32) Error!Ready {
while (true) {
var pfds: [3]linux.pollfd = undefined;
var n: usize = 0;
@@ -329,13 +379,14 @@ pub const Session = struct {
pfds[n] = .{ .fd = ifd, .events = linux.POLL.IN, .revents = 0 };
n += 1;
}
- if (n == 1) return .socket;
- const prc = linux.poll(&pfds, @intCast(n), -1);
+ if (n == 1 and timeout_ms < 0) return .socket;
+ const prc = linux.poll(&pfds, @intCast(n), timeout_ms);
switch (linux.errno(prc)) {
.SUCCESS => {},
.INTR, .AGAIN => continue,
else => return error.Io,
}
+ if (prc == 0) return .timeout;
// A reply that is already there wins over everything else.
if (pfds[0].revents != 0) return .socket;
if (stop_at) |i| {
@@ -368,6 +419,13 @@ pub const Session = struct {
}
};
+/// The monotonic clock in milliseconds: deadlines, not timestamps.
+fn nowMs() i64 {
+ var ts: linux.timespec = undefined;
+ _ = linux.clock_gettime(.MONOTONIC, &ts);
+ return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);
+}
+
fn chunkSize(max: u32, iounit: u32) u32 {
if (iounit != 0 and iounit < max) return iounit;
return max;
@@ -723,6 +781,62 @@ test "interrupted chunk loops: partial count if data moved, Interrupted otherwis
// -- in-process server test ---------------------------------------------------------
+test "flush grace: a Tflush the server never answers wedges the session after the grace" {
+ var pi = try PipeInterrupt.init(7);
+ defer pi.deinit();
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .ignore, .on_hang_inject = &pi };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ s.interrupt = pi.interface();
+ s.flush_grace_ms = 100;
+ // The read at offset 0 hangs; the injected INTERRUPT sends a Tflush; the
+ // server ignores it; the grace runs out.
+ var buf: [100]u8 = undefined;
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ try testing.expect(s.wedged);
+ // From here on the session is closed for business, without another
+ // byte to the server.
+ try testing.expectError(error.Closed, s.read(1, 10, &buf));
+ try testing.expectError(error.Closed, s.clunk(1));
+ try testing.expectEqual(@as(u32, 1), fs.flushes.load(.seq_cst));
+ // The socket is still the session's to close: the server's thread ends
+ // when `close` drops it.
+}
+
+test "flush grace: zero waits for the Rflush, however late" {
+ var pi = try PipeInterrupt.init(7);
+ defer pi.deinit();
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush, .on_hang_inject = &pi, .read_delay_ns = 0 };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ s.interrupt = pi.interface();
+ s.flush_grace_ms = 0;
+ var buf: [100]u8 = undefined;
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expect(!s.wedged);
+ // The session lives: the Rflush released the tag and reads go on.
+ try testing.expectEqual(@as(usize, 40), try s.read(1, 10, buf[0..40]));
+}
+
+test "abort_on_cancel: a cancellation ends the rpc at once, sends no Tflush, wedges the session" {
+ var pi = try PipeInterrupt.init(7);
+ defer pi.deinit();
+ var fs: FakeServer = .{ .fd = -1, .msize = 8192, .file_len = 50, .hang_offset = 0, .on_flush = .rflush, .on_hang_inject = &pi };
+ var pair = try fs.start();
+ defer pair.close();
+ const s = &pair.session;
+ s.interrupt = pi.interface();
+ s.abort_on_cancel = true;
+ var buf: [100]u8 = undefined;
+ try testing.expectError(error.Interrupted, s.read(1, 0, &buf));
+ try testing.expect(s.wedged);
+ try testing.expectEqual(@as(u32, 0), fs.flushes.load(.seq_cst));
+ try testing.expectError(error.Closed, s.read(1, 10, &buf));
+}
+
/// A tiny 9P2000 backend on a cloud9.Server: answers version/attach/walk/stat/open/
/// read/clunk/remove with canned data. Runs in its own thread over a socketpair.
/// Test support only (bridge.zig's tests use it too).
@@ -734,9 +848,11 @@ pub const FakeServer = struct {
/// A Tread at this offset is never answered (a blocked stream read); the
/// server keeps serving whatever else arrives, notably a Tflush.
hang_offset: ?u64 = null,
- /// What a Tflush gets: the connection dropped (`.hangup`), an Rflush, or
- /// first the Rread the flush was aimed at and then the Rflush (the race).
- on_flush: enum { hangup, rflush, reply_then_rflush } = .hangup,
+ /// What a Tflush gets: the connection dropped (`.hangup`), an Rflush,
+ /// first the Rread the flush was aimed at and then the Rflush (the
+ /// race), or nothing at all (`.ignore`: counted, never answered, the
+ /// read stays hung — a wedged server).
+ on_flush: enum { hangup, rflush, reply_then_rflush, ignore } = .hangup,
/// Observed by the test thread: number of Tflush seen and the last oldtag.
flushes: std.atomic.Value(u32) = .init(0),
flush_oldtag: std.atomic.Value(u32) = .init(0xFFFF),
@@ -827,6 +943,7 @@ pub const FakeServer = struct {
switch (fs.on_flush) {
// The test's "hang up now" signal.
.hangup => return,
+ .ignore => {},
.rflush => {
if (hung != null and hung.?.tag == m.oldtag) hung = null;
try srv.reply(tag, .rflush);