summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--9ns/docs/DESIGN.md25
-rw-r--r--9ns/src/bridge.zig372
-rwxr-xr-x9ns/test/mntgen.sh74
-rw-r--r--9proc/src/linux/probe.zig7
-rwxr-xr-x9proc/test/adv_linux_probe.py34
-rw-r--r--src/post.zig138
-rw-r--r--src/serve.zig224
7 files changed, 865 insertions, 9 deletions
diff --git a/9ns/docs/DESIGN.md b/9ns/docs/DESIGN.md
index 1d66387..6b1b80c 100644
--- a/9ns/docs/DESIGN.md
+++ b/9ns/docs/DESIGN.md
@@ -114,6 +114,31 @@ posted name reaches that server's whole 9P tree, and nothing is connected
until something walks. XDG_RUNTIME_DIR unset is fatal before anything is
forked (the registry is not guessable; no `/tmp` fallback).
+A registry entry that is a **directory** is served the same way the root
+is: a synthetic directory (its node id under the reserved index
+`synth_index`, the slot in the low bits) listing the real directory's
+entries, dialing the sockets found inside on walk and recursing into
+further directories — up to `max_synth_depth` (8) levels, bounded by
+`max_synth_dirs` (64) synthetic nodes per 9ns process. This is how
+multi-service providers organize themselves (zmx posts its sessions
+under `zmx/<name>`), the plan9port `mntgen` shape: one tree, many
+mounts, each entry a mount point. Non-socket, non-directory entries
+inside a directory answer EIO on walk, exactly like a plain file in the
+registry itself; a directory's slots are freed when the kernel forgets
+the dentry.
+
+A synthetic node's slot comes back through FORGET, and the kernel sends
+most of them as `BATCH_FORGET`, whose header `nodeid` is 0 and whose
+body carries one `(nodeid, nlookup)` per forgotten node — for any mix of
+owners. The dispatcher therefore cannot route a batch by its header the
+way it routes every other request: `distributeForgets` unpacks the body
+and hands each entry to its owner (synthetic root, synthetic
+subdirectory or per-mount bridge). Routing the batch whole instead loses
+every entry in it, so the 64 slots leak and a subdirectory served once
+answers EIO forever. The forgets themselves arrive on the kernel's
+schedule, not at `close`, so a slot may take a moment to return; a
+listing that needs one meanwhile answers EIO rather than waiting.
+
```
program 9ns parent
in new userns │
diff --git a/9ns/src/bridge.zig b/9ns/src/bridge.zig
index 76327c3..457c883 100644
--- a/9ns/src/bridge.zig
+++ b/9ns/src/bridge.zig
@@ -855,6 +855,15 @@ pub const max_mounts: usize = 4096;
pub const max_root_dirs: usize = 64;
/// The staged buffer handed to `post.posted` for one root listing.
pub const stage_len: usize = 8192;
+/// Registry entries that are directories are served like the root itself
+/// (a synthetic directory mirroring the real one, dialing the sockets
+/// found inside); this bounds the synthetic directories one 9ns serves.
+pub const max_synth_dirs: usize = 64;
+/// How deep those registry subdirectories nest.
+pub const max_synth_depth: u8 = 8;
+/// The node-id index reserved for synthetic registry subdirectories;
+/// mounts use ordinals below 4096, so this never collides with one.
+pub const synth_index: u32 = 0xFFFF_FFFF;
/// The node id a server's subtree lives under: `index` in the top bits,
/// `local` (1 = that server's 9P root) below.
@@ -891,6 +900,31 @@ pub const MntgenOptions = struct {
msize: u32 = 131072,
};
+/// One registry subdirectory served as a synthetic directory: the real
+/// path it mirrors, the registry-relative key its children dial under,
+/// its slot (which names its node id) and its depth from the registry.
+const SynthDir = struct {
+ slot: usize = 0,
+ parent: u64 = 0,
+ depth: u8 = 0,
+ path_buf: [post.sun_path_len]u8 = undefined,
+ path_len: u16 = 0,
+ rel_buf: [post.sun_path_len]u8 = undefined,
+ rel_len: u16 = 0,
+
+ fn path(sd: *const SynthDir) [:0]const u8 {
+ return sd.path_buf[0..sd.path_len :0];
+ }
+
+ fn rel(sd: *const SynthDir) []const u8 {
+ return sd.rel_buf[0..sd.rel_len];
+ }
+
+ fn nodeid(sd: *const SynthDir) u64 {
+ return mountNode(synth_index, sd.slot + 1);
+ }
+};
+
/// One dialed server: its 9P session, its bridge state and its worker
/// thread, plus the queue the dispatcher feeds requests through.
const Mount = struct {
@@ -998,6 +1032,8 @@ const Mntgen = struct {
next_index: u32 = 1,
/// Open directory handles of the synthetic root (dispatcher-owned).
root_dirs: [max_root_dirs]?*DirList = @splat(null),
+ /// Synthetic registry subdirectories (dispatcher-owned), by slot.
+ synths: [max_synth_dirs]?*SynthDir = @splat(null),
fn deinit(mg: *Mntgen) void {
// Wake every worker, then join: at this point the child is gone (or
@@ -1031,6 +1067,12 @@ const Mntgen = struct {
slot.* = null;
}
}
+ for (&mg.synths) |*slot| {
+ if (slot.*) |sd| {
+ mg.gpa.destroy(sd);
+ slot.* = null;
+ }
+ }
if (mg.req_buf.len != 0) mg.gpa.free(mg.req_buf);
if (mg.data_buf.len != 0) mg.gpa.free(mg.data_buf);
if (mg.stage.len != 0) mg.gpa.free(mg.stage);
@@ -1079,12 +1121,23 @@ const Mntgen = struct {
mg.routeInterrupt(in.unique);
return true;
},
+ // A BATCH_FORGET carries entries for many owners at once and
+ // cannot be routed by its header nodeid (the kernel sends 0);
+ // see `distributeForgets`.
+ .batch_forget => {
+ mg.distributeForgets(req);
+ return true;
+ },
else => {},
}
if (h.nodeid == fuse.root_id) {
try mg.handleRoot(req);
return true;
}
+ if (mountIndex(h.nodeid) == synth_index) {
+ try mg.handleSynthDir(req);
+ return true;
+ }
const wants_reply = switch (op) {
.forget, .batch_forget => false,
else => true,
@@ -1133,6 +1186,68 @@ const Mntgen = struct {
}
}
+ /// FUSE_BATCH_FORGET carries (nodeid, nlookup) entries for many owners
+ /// at once — synthetic-root, synthetic-subdirectory and per-mount
+ /// nodeids can all appear in one batch, and the kernel puts 0 in the
+ /// header nodeid. Routing such a batch like an ordinary request would
+ /// hand every entry to one wrong owner and drop the rest: a mount's
+ /// bridge would leak the fid behind each dropped entry forever, and
+ /// a dropped synthetic-subdirectory entry would leak its slot until
+ /// every subdirectory lookup answers EIO. So the dispatcher keeps
+ /// what it owns — the root holds nothing, a synth slot is freed here —
+ /// and forwards each mount-owned entry to its worker as a plain
+ /// single FUSE_FORGET.
+ fn distributeForgets(mg: *Mntgen, req: fuse.Request) void {
+ const in = fuse.body(fuse.BatchForgetIn, req) catch return;
+ const rest = req.body[@sizeOf(fuse.BatchForgetIn)..];
+ const count: usize = in.count;
+ if (rest.len < count * @sizeOf(fuse.ForgetOne)) return;
+ for (0..count) |i| {
+ const one = std.mem.bytesToValue(fuse.ForgetOne, rest[i * @sizeOf(fuse.ForgetOne) ..][0..@sizeOf(fuse.ForgetOne)]);
+ mg.forgetOne(one.nodeid, one.nlookup, req.header.unique);
+ }
+ }
+
+ /// Forgets one node by id, whatever owns it: the root holds nothing,
+ /// a synthetic subdirectory's slot goes back to the pool, and a
+ /// mount-owned node is forwarded to its worker (which clunks the fid
+ /// behind it). Unknown or dead mounts drop the entry, like a single
+ /// FORGET routed by `route`.
+ fn forgetOne(mg: *Mntgen, nodeid: u64, nlookup: u64, unique: u64) void {
+ if (nodeid == fuse.root_id) return;
+ const idx = mountIndex(nodeid);
+ if (idx == synth_index) {
+ const local = nodeid & mount_node_mask;
+ if (local == 0 or local > max_synth_dirs) return;
+ const slot: usize = @intCast(local - 1);
+ if (mg.synths[slot]) |sd| {
+ mg.trace(" synthetic directory slot {d} forgotten", .{slot});
+ mg.gpa.destroy(sd);
+ mg.synths[slot] = null;
+ }
+ return;
+ }
+ const m = if (idx < max_mounts) mg.mounts[idx] else null;
+ if (m == null or m.?.dead.load(.seq_cst)) return;
+ // A synthesized single FORGET (no reply is expected for one, so
+ // the unique is only bookkeeping).
+ var buf: [@sizeOf(fuse.InHeader) + @sizeOf(fuse.ForgetIn)]u8 = undefined;
+ const hdr = fuse.InHeader{
+ .len = @sizeOf(fuse.InHeader) + @sizeOf(fuse.ForgetIn),
+ .opcode = @intFromEnum(fuse.Opcode.forget),
+ .unique = unique,
+ .nodeid = nodeid,
+ .uid = 0,
+ .gid = 0,
+ .pid = 0,
+ .total_extlen = 0,
+ .padding = 0,
+ };
+ @memcpy(buf[0..@sizeOf(fuse.InHeader)], std.mem.asBytes(&hdr));
+ @memcpy(buf[@sizeOf(fuse.InHeader)..], std.mem.asBytes(&fuse.ForgetIn{ .nlookup = nlookup }));
+ mg.enqueue(m.?, &buf);
+ }
+
/// 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
@@ -1257,11 +1372,34 @@ const Mntgen = struct {
const out = mg.rootEntryOut(m.root_node, m.root_attr);
return mg.reply(u, &.{std.mem.asBytes(&out)});
}
- const m = mg.dialMount(name) catch |e| switch (e) {
- error.NotPosted => {
- mg.trace(" lookup '{s}': nothing posted under that name", .{name});
- return mg.replyError(u, .NOENT);
+ // 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.
+ 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);
+ const st = std.Io.Dir.statFile(.cwd(), mg.mo.io, entry_path, .{}) catch {
+ mg.trace(" lookup '{s}': nothing posted under that name", .{name});
+ return mg.replyError(u, .NOENT);
+ };
+ switch (st.kind) {
+ .directory => {
+ const sd = mg.newSynth(entry_path, name, 1, fuse.root_id) catch |e| {
+ mg.trace(" lookup '{s}': no synthetic slot: {t}", .{ name, e });
+ return mg.replyError(u, .IO);
+ };
+ const node = sd.nodeid();
+ const out = mg.rootEntryOut(node, mg.synthAttr(node));
+ return mg.reply(u, &.{std.mem.asBytes(&out)});
},
+ .unix_domain_socket => {},
+ 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);
@@ -1297,6 +1435,204 @@ const Mntgen = struct {
return list;
}
+ // -- synthetic registry subdirectories ------------------------------------
+
+ /// A directory entry in the registry (or in one of its
+ /// subdirectories) is served like the root: a synthetic directory
+ /// listing the real one, whose sockets dial on walk and whose
+ /// directories recurse. Served on the dispatcher thread, like the
+ /// root.
+ fn handleSynthDir(mg: *Mntgen, req: fuse.Request) error{FuseIo}!void {
+ const u = req.header.unique;
+ const local = req.header.nodeid & mount_node_mask;
+ if (local == 0 or local > max_synth_dirs) return mg.replyError(u, .IO);
+ const slot: usize = @intCast(local - 1);
+ const sd = mg.synths[slot] orelse return mg.replyError(u, .IO);
+ switch (req.header.op()) {
+ .forget, .batch_forget => {
+ // The kernel dropped the dentry; the slot goes with it.
+ mg.trace(" synthetic directory slot {d} forgotten", .{slot});
+ mg.gpa.destroy(sd);
+ mg.synths[slot] = null;
+ },
+ .getattr => {
+ const out = fuse.AttrOut{ .attr = mg.synthAttr(sd.nodeid()) };
+ try mg.reply(u, &.{std.mem.asBytes(&out)});
+ },
+ .lookup => try mg.synthLookup(sd, req),
+ .opendir => {
+ var fh: ?usize = null;
+ for (&mg.root_dirs, 0..) |*dir_slot, i| {
+ if (dir_slot.* == null) {
+ fh = i;
+ break;
+ }
+ }
+ const dir_slot = fh orelse return mg.replyError(u, .MFILE);
+ const list = mg.synthListing(sd) catch {
+ return mg.replyError(u, .IO);
+ };
+ mg.root_dirs[dir_slot] = list;
+ const out = fuse.OpenOut{ .fh = dir_slot };
+ try mg.reply(u, &.{std.mem.asBytes(&out)});
+ },
+ .readdir => {
+ const in = fuse.body(fuse.ReadIn, req) catch return mg.replyError(u, .BADF);
+ if (in.fh >= max_root_dirs) return mg.replyError(u, .BADF);
+ const list = mg.root_dirs[@intCast(in.fh)] orelse return mg.replyError(u, .BADF);
+ const size: usize = @min(in.size, max_write);
+ const used = packDirents(list.entries.items, in.offset, mg.data_buf[0..size]);
+ try mg.reply(u, &.{mg.data_buf[0..used]});
+ },
+ .release, .releasedir => {
+ const in = fuse.body(fuse.ReleaseIn, req) catch return mg.replyError(u, .BADF);
+ if (in.fh < max_root_dirs) {
+ if (mg.root_dirs[@intCast(in.fh)]) |list| {
+ list.deinit(mg.gpa);
+ mg.gpa.destroy(list);
+ mg.root_dirs[@intCast(in.fh)] = null;
+ }
+ }
+ try mg.reply(u, &.{});
+ },
+ .statfs => {
+ const out = fuse.StatfsOut{ .st = .{ .bsize = 4096, .namelen = 255, .frsize = 4096 } };
+ try mg.reply(u, &.{std.mem.asBytes(&out)});
+ },
+ .flush, .fsync, .fsyncdir => try mg.reply(u, &.{}),
+ // Capability probes read as "not supported", like the root's.
+ .access, .setxattr, .getxattr, .listxattr, .removexattr, .statx => try mg.replyError(u, .NOSYS),
+ else => try mg.replyError(u, .PERM),
+ }
+ }
+
+ fn synthAttr(mg: *const Mntgen, nodeid: u64) fuse.Attr {
+ // Read-only like the root and like /srv.
+ return .{
+ .ino = nodeid,
+ .mode = fuse.S_IFDIR | 0o555,
+ .nlink = 2,
+ .uid = mg.opts.uid,
+ .gid = mg.opts.gid,
+ .blksize = 4096,
+ };
+ }
+
+ fn synthLookup(mg: *Mntgen, sd: *SynthDir, req: fuse.Request) error{FuseIo}!void {
+ 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()));
+ return mg.reply(u, &.{std.mem.asBytes(&out)});
+ }
+ // The kernel resolves ".." from its own dentry tree and a LOOKUP of
+ // it has never been observed, but it must not alias the directory
+ // 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));
+ return mg.reply(u, &.{std.mem.asBytes(&out)});
+ }
+ if (!post.legalName(name)) return mg.replyError(u, .NOENT);
+ // The registry-relative key this child dials under (a mount's
+ // name, for findMount).
+ var key_buf: [post.sun_path_len]u8 = undefined;
+ const key = std.fmt.bufPrint(&key_buf, "{s}/{s}", .{ sd.rel(), name }) catch
+ return mg.replyError(u, .NOTNAM);
+ var path_buf: [post.sun_path_len]u8 = undefined;
+ 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)});
+ }
+ const st = std.Io.Dir.statFile(.cwd(), mg.mo.io, child, .{}) catch |e| switch (e) {
+ error.FileNotFound => {
+ mg.trace(" lookup '{s}': no entry", .{key});
+ return mg.replyError(u, .NOENT);
+ },
+ else => {
+ mg.trace(" lookup '{s}': stat failed: {t}", .{ key, e });
+ return mg.replyError(u, .IO);
+ },
+ };
+ switch (st.kind) {
+ .directory => {
+ const child_sd = mg.newSynth(child, key, sd.depth + 1, sd.nodeid()) catch |e| {
+ mg.trace(" lookup '{s}': no synthetic slot: {t}", .{ key, e });
+ return mg.replyError(u, .IO);
+ };
+ const node = child_sd.nodeid();
+ const out = mg.rootEntryOut(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 });
+ 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)});
+ },
+ // Not a service and not a directory: the entry answers EIO on
+ // walk, like a plain file in the registry itself.
+ else => {
+ mg.trace(" lookup '{s}': entry is not a socket or directory", .{key});
+ return mg.replyError(u, .IO);
+ },
+ }
+ }
+
+ /// One OPENDIR of a synthetic subdirectory: `.` and `..` plus the
+ /// real directory's entries, snapshotted for the life of the handle
+ /// (a fresh OPENDIR sees fresh entries), like the root.
+ fn synthListing(mg: *Mntgen, sd: *SynthDir) !*DirList {
+ const list = try mg.gpa.create(DirList);
+ errdefer mg.gpa.destroy(list);
+ list.* = .{};
+ errdefer list.deinit(mg.gpa);
+ const node = sd.nodeid();
+ try list.entries.append(mg.gpa, .{ .name = try mg.gpa.dupe(u8, "."), .ino = node, .dtype = fuse.DT_DIR });
+ try list.entries.append(mg.gpa, .{ .name = try mg.gpa.dupe(u8, ".."), .ino = sd.parent, .dtype = fuse.DT_DIR });
+ var names = post.postedDir(mg.mo.io, sd.path(), mg.stage) catch |e| {
+ mg.trace(" directory listing failed: {t}", .{e});
+ return error.Registry;
+ };
+ while (names.next()) |name| {
+ if (!validDirentName(name)) continue;
+ try list.entries.append(mg.gpa, .{ .name = try mg.gpa.dupe(u8, name), .ino = nameIno(name), .dtype = fuse.DT_DIR });
+ }
+ return list;
+ }
+
+ /// Allocates a synthetic directory node mirroring `path`, keyed by
+ /// the registry-relative `key`, at `depth` under `parent`'s node id.
+ fn newSynth(mg: *Mntgen, path: [:0]const u8, key: []const u8, depth: u8, parent: u64) !*SynthDir {
+ if (depth > max_synth_depth) return error.TooDeep;
+ var slot: ?usize = null;
+ for (&mg.synths, 0..) |*s, i| {
+ if (s.* == null) {
+ slot = i;
+ break;
+ }
+ }
+ const i = slot orelse return error.TooMany;
+ const sd = try mg.gpa.create(SynthDir);
+ errdefer mg.gpa.destroy(sd);
+ sd.* = .{ .slot = i, .parent = parent, .depth = depth };
+ if (path.len + 1 > sd.path_buf.len) return error.NameTooLong;
+ @memcpy(sd.path_buf[0..path.len], path);
+ sd.path_buf[path.len] = 0;
+ sd.path_len = @intCast(path.len);
+ if (key.len > sd.rel_buf.len) return error.NameTooLong;
+ @memcpy(sd.rel_buf[0..key.len], key);
+ sd.rel_len = @intCast(key.len);
+ mg.synths[i] = sd;
+ return sd;
+ }
+
// -- dialing ---------------------------------------------------------------
/// Dials `name` out of the registry, attaches, stats the server root and
@@ -1304,10 +1640,22 @@ const Mntgen = struct {
/// 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);
+ }
+
+ /// The shared dial tail: session, attach, stat, bridge and worker.
+ fn mountStream(mg: *Mntgen, stream: std.Io.net.Stream, key: []const u8) !*Mount {
if (mg.next_index >= max_mounts) return error.TooManyMounts;
const index: u32 = mg.next_index;
-
- const stream = try post.dial(mg.mo.io, mg.mo.env, name);
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,
@@ -1369,7 +1717,7 @@ const Mntgen = struct {
.gpa = mg.gpa,
.io = mg.io,
.debug = mg.opts.debug,
- .name = try mg.gpa.dupe(u8, name),
+ .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),
@@ -1848,6 +2196,16 @@ test "mntgen node id layout: index in the top bits, local ids below" {
try testing.expect(mount_node_mask == (1 << 32) - 1);
}
+test "mntgen synthetic subdirectory node ids route apart from mounts" {
+ for (1..max_synth_dirs + 1) |i| {
+ const node = mountNode(synth_index, i);
+ try testing.expectEqual(synth_index, mountIndex(node));
+ try testing.expectEqual(i, node & mount_node_mask);
+ }
+ // The reserved index can never be a mount ordinal.
+ try testing.expect(synth_index >= max_mounts);
+}
+
test "mntgen synthetic-root inos are deterministic and name-derived" {
const a = nameIno("alpha");
try testing.expectEqual(a, nameIno("alpha"));
diff --git a/9ns/test/mntgen.sh b/9ns/test/mntgen.sh
index ea58e79..aade76c 100755
--- a/9ns/test/mntgen.sh
+++ b/9ns/test/mntgen.sh
@@ -2,6 +2,8 @@
# Integration tests for 9ns --mntgen: one FUSE mount whose synthetic root
# lists the posted-9P registry ($XDG_RUNTIME_DIR/9p), servers dialed lazily
# on the first walk into their name, one worker thread per server.
+# Registry entries that are directories are served like the root itself:
+# their sockets dial on walk and their directories recurse (depth-capped).
# Usage: bash 9ns/test/mntgen.sh <9ns> <9proc-demo> (zig build 9ns-itest)
# Exit 0 on success (or when the machine cannot run the tests), 1 on failure.
set -u
@@ -148,6 +150,78 @@ expect_eq "re-dial picks up the fresh post" "$(zig version)" "$(run_in "cat $M/a
echo "# mount defaults and environment"
expect_eq "default mountpoint is /mnt/9p" "$M" "$(timeout 60 "$NS" --mntgen -- sh -c 'echo $NINE_MOUNT')"
expect_contains "default mount is served" "alpha" "$(timeout 60 "$NS" --mntgen -- sh -c 'ls $NINE_MOUNT')"
+
+echo "# registry subdirectories are served like the root"
+mkdir -p "$REG/svc/deep"
+"$PROC" --unix "$REG/svc/delta" & PIDS+=($!)
+wait_socket "$REG/svc/delta"
+"$PROC" --unix "$REG/svc/deep/eps" & PIDS+=($!)
+wait_socket "$REG/svc/deep/eps"
+touch "$REG/svc/plainfile"
+expect_contains "root lists the subdirectory" "svc" "$(run_in "ls $M")"
+SVC_LS=$(run_in "ls $M/svc")
+expect_contains "subdirectory lists its sockets" "delta" "$SVC_LS"
+expect_contains "subdirectory lists nested directories" "deep" "$SVC_LS"
+expect_contains "subdirectory lists a plain file" "plainfile" "$SVC_LS"
+expect_eq "socket inside a directory dials" "$(zig version)" "$(run_in "cat $M/svc/delta/build/zig_version")"
+expect_eq "socket inside a nested directory dials" "$(zig version)" "$(run_in "cat $M/svc/deep/eps/build/zig_version")"
+expect_eq "walk into a plain file inside a directory is EIO, not a crash" "1" "$(run_in "cat $M/svc/plainfile 2>/dev/null; echo \$?")"
+expect_contains "the plain file stays listed" "plainfile" "$(run_in "ls $M/svc")"
+mkdir -p "$REG/a/b/c/d/e/f/g/h/i/j"
+expect_eq "eight levels of nesting serve" "ok" "$(run_in "[ -d $M/a/b/c/d/e/f/g/h ] && echo ok")"
+expect_eq "the ninth level answers EIO" "1" "$(run_in "[ -e $M/a/b/c/d/e/f/g/h/i/j ]; echo \$?")"
+expect_eq "\".\" inside a synthetic subdirectory is the subdirectory" "delta" \
+ "$(run_in "ls $M/svc/. | grep -x delta")"
+expect_eq "\"..\" from a synthetic subdirectory is the mount root" "svc" \
+ "$(run_in "ls $M/svc/.. | grep -x svc")"
+
+# A synthetic subdirectory holds one of 64 slots until the kernel forgets
+# its node. FUSE delivers those forgets in BATCH_FORGET messages whose
+# header nodeid is 0, so a dispatcher that routes a batch by that header
+# drops every entry in it and the slots never come back: subdirectories
+# served once then answer EIO forever. The forgets themselves arrive
+# lazily (the kernel frees dentries on its own schedule), so the check
+# retries rather than sampling once.
+if command -v python3 > /dev/null; then
+ echo "# synthetic subdirectory slots come back after the kernel forgets them"
+ for i in $(seq 1 65); do mkdir -p "$REG/s$(printf %02d "$i")"; done
+ cat > "$TMP/slots.py" <<'PYEOF'
+import os, sys, time
+M = sys.argv[1]
+n = 64
+def opendir(i):
+ return os.open("%s/s%02d" % (M, i), os.O_RDONLY | os.O_DIRECTORY)
+fds = [opendir(i) for i in range(1, n + 1)]
+try:
+ os.close(opendir(n + 1))
+ print("65th=opened") # the cap is not enforced
+except OSError:
+ print("65th=refused")
+for fd in fds:
+ os.close(fd)
+deadline, ok = time.time() + 10, 0
+while True:
+ ok = 0
+ for i in range(1, n + 1):
+ try:
+ os.close(opendir(i))
+ ok += 1
+ except OSError:
+ pass
+ if ok == n or time.time() > deadline:
+ break
+ time.sleep(0.2)
+print("recycled=%d" % ok)
+PYEOF
+ SLOTS=$(run_in "python3 $TMP/slots.py $M")
+ expect_contains "the 65th concurrent subdirectory is refused cleanly" "65th=refused" "$SLOTS"
+ expect_contains "all 64 slots come back after a close burst" "recycled=64" "$SLOTS"
+ rm -rf "$REG"/s[0-9][0-9]
+else
+ echo "SKIP: synthetic-slot recycling check needs python3"
+fi
+
+rm -rf "$REG/svc" "$REG/a"
expect_eq "mount is fuse" "yes" "$(run_in "grep -q \"^9ns $M fuse\" /proc/mounts && echo yes")"
echo "# usage errors"
diff --git a/9proc/src/linux/probe.zig b/9proc/src/linux/probe.zig
index c2f1026..00c256e 100644
--- a/9proc/src/linux/probe.zig
+++ b/9proc/src/linux/probe.zig
@@ -542,7 +542,12 @@ pub fn Probe(comptime Srv: type) type {
const st = unixStat(@ptrCast(&sa.path)) catch return error.Occupied;
if (st) |s| {
if (s.mode & linux.S.IFMT != linux.S.IFSOCK) return error.Occupied;
- if (cloud9.post.probe(@ptrCast(&sa.path)) != .stale) return error.AlreadyListening;
+ // The probe wants the path, not the whole `sun_path`
+ // array: a 108-byte slice is past the address budget, and
+ // `probe` owns that uncertainty as `.live` — which would
+ // make every stale socket look like a live server.
+ const probe_path: [:0]const u8 = sa.path[0..path.len :0];
+ if (cloud9.post.probe(probe_path) != .stale) return error.AlreadyListening;
_ = linux.unlink(@ptrCast(&sa.path));
}
try p.check(linux.bind(lfd, @ptrCast(&sa), @sizeOf(linux.sockaddr.un)));
diff --git a/9proc/test/adv_linux_probe.py b/9proc/test/adv_linux_probe.py
index 186f008..c921ebb 100755
--- a/9proc/test/adv_linux_probe.py
+++ b/9proc/test/adv_linux_probe.py
@@ -289,6 +289,40 @@ def attack_signals(server):
ok("SIGTERM unlinks the unix socket (clean stop path)", not os.path.exists(path))
s.stop()
+ # A server killed outright leaves its socket behind. The next server
+ # of the same name must recognise the corpse and take the name over:
+ # the listener probes the path, and a probe that cannot tell answers
+ # "live", so a caller that hands it the whole 108-byte sun_path array
+ # instead of the path turns every stale socket into AlreadyListening.
+ s = Srv(server)
+ client(s.path, timeout=5)
+ stale = s.path
+ s.proc.send_signal(signal.SIGKILL)
+ s.proc.wait(timeout=5)
+ ok("SIGKILL leaves the socket behind", os.path.exists(stale))
+ taker = subprocess.Popen([server, "--unix", stale], stderr=subprocess.PIPE)
+ try:
+ # The path exists throughout (the corpse, then the new socket), so
+ # the connect itself is the readiness signal.
+ err, c = None, None
+ for _ in range(250):
+ try:
+ c = client(stale, timeout=10)
+ break
+ except OSError as e:
+ err = e
+ time.sleep(0.02)
+ ok("a fresh server takes over a stale socket path",
+ c is not None and rd(c, [b"build", b"zig_version"])[0] == Rread,
+ err or taker.poll())
+ except Exception as e:
+ ok("a fresh server takes over a stale socket path", False, e)
+ finally:
+ taker.kill()
+ taker.wait(timeout=5)
+ if os.path.exists(stale):
+ os.unlink(stale)
+
def attack_probe(ns, server):
print("# probe loop and admission")
diff --git a/src/post.zig b/src/post.zig
index 47a7d3d..641970d 100644
--- a/src/post.zig
+++ b/src/post.zig
@@ -124,6 +124,14 @@ pub const PostedError = PathError || error{NoSpace} || Io.Dir.OpenError || Io.Di
pub fn posted(io: Io, env: Env, out: []u8) PostedError!Names {
var dir_buf: [std.fs.max_path_bytes]u8 = undefined;
const dir_path = try registryDir(env, &dir_buf);
+ return postedDir(io, dir_path, out);
+}
+
+/// Lists one directory's entry names into `out` and returns an iterator
+/// over them — the registry scan, generalized for the registry
+/// subdirectories that 9ns mntgen serves. Nothing is dialed; a missing
+/// directory lists as empty.
+pub fn postedDir(io: Io, dir_path: [:0]const u8, out: []u8) PostedError!Names {
const dir = Io.Dir.openDirAbsolute(io, dir_path, .{ .iterate = true }) catch |err| switch (err) {
error.FileNotFound, error.NotDir => return .{ .bytes = out[0..0] },
else => return err,
@@ -154,6 +162,10 @@ pub const Probe = enum { none, stale, live };
/// `live`, because uncertainty must be owned by the server, never
/// resolved by deleting what may be someone's socket.
pub fn probe(path: [:0]const u8) Probe {
+ // A path that cannot fit a `sun_path` cannot be asked about at all;
+ // like `transport.isListening`, uncertainty is `.live` (occupied) so
+ // a claim never deletes a name it cannot inspect.
+ if (path.len >= sun_path_len) return .live;
const fd = socketNonblocking() catch return .live;
defer _ = linux.close(fd);
var addr: linux.sockaddr.un = .{ .path = @splat(0) };
@@ -217,14 +229,37 @@ pub fn dial(io: Io, env: Env, name: []const u8) DialError!Io.net.Stream {
error.AccessDenied => return error.AccessDenied,
error.Loop => return error.SymLinkLoop,
error.NotDir => return error.NotDir,
+ error.NameTooLong => return error.NameTooLong,
+ };
+ _ = io;
+ return .{ .socket = .{ .handle = fd, .address = .{ .ip4 = .loopback(0) } } };
+}
+
+/// Dials the socket at `path` — any registry path, including a
+/// subdirectory entry — and returns its stream. The same blocking,
+/// close-on-exec and refused-detection semantics as `dial`, with no
+/// name validation: the caller composed the path.
+pub fn dialPath(io: Io, path: [:0]const u8) DialError!Io.net.Stream {
+ const fd = connectBlocking(path) catch |err| switch (err) {
+ error.Socket => return error.SystemResources,
+ error.Noent => return error.NotPosted,
+ error.Refused => return error.Stale,
+ error.AccessDenied => return error.AccessDenied,
+ error.Loop => return error.SymLinkLoop,
+ error.NotDir => return error.NotDir,
+ error.NameTooLong => return error.NameTooLong,
};
_ = io;
return .{ .socket = .{ .handle = fd, .address = .{ .ip4 = .loopback(0) } } };
}
-const ConnectError = error{ Socket, Noent, Refused, AccessDenied, Loop, NotDir };
+const ConnectError = error{ Socket, Noent, Refused, AccessDenied, Loop, NotDir, NameTooLong };
fn connectBlocking(path: [:0]const u8) ConnectError!i32 {
+ // The caller composed the path (dialPath runs no name validation), so
+ // it must still fit a `sun_path`; beyond it the copy into the kernel
+ // address would overrun the stack struct.
+ if (path.len >= sun_path_len) return error.NameTooLong;
const rc = linux.socket(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0);
if (rawErrno(@bitCast(rc)) != .SUCCESS) return error.Socket;
const fd: i32 = @intCast(rc);
@@ -1022,3 +1057,104 @@ test "post watch: queue overflow surfaces, a deleted registry is gone" {
try testing.expectEqualStrings("back", back.?.name);
try Io.Dir.deleteFileAbsolute(io, try std.fmt.bufPrint(&name_buf, "{s}/back", .{reg}));
}
+
+test "post probe/dialPath: a path past the sun_path budget is refused, never read" {
+ // A path longer than a sockaddr's `sun_path` cannot be handed to the
+ // kernel; the raw probe/dial copies must not overrun the stack
+ // address struct. `probe` owns the uncertainty as `.live` (never
+ // delete what cannot be inspected); `dialPath` names the failure.
+ var buf: [512]u8 = undefined;
+ const over = sun_path_len + 8;
+ @memset(buf[0..over], 'a');
+ buf[over] = 0;
+ const p: [:0]const u8 = buf[0..over :0];
+ try testing.expectEqual(Probe.live, probe(p));
+ try testing.expectError(error.NameTooLong, dialPath(testing.io, p));
+ // Exactly `sun_path_len` bytes leaves no room for a zero and is not a
+ // name either.
+ buf[sun_path_len] = 0;
+ const at: [:0]const u8 = buf[0..sun_path_len :0];
+ try testing.expectEqual(Probe.live, probe(at));
+ try testing.expectError(error.NameTooLong, dialPath(testing.io, at));
+}
+
+test "post dialPath: refused, missing, non-socket, loop and a CLOEXEC stream" {
+ if (builtin.os.tag != .linux) return error.SkipZigTest;
+ var s: Scratch = .{ .dir = undefined };
+ try s.start();
+ defer s.end();
+ const io = testing.io;
+ var real_buf: [std.fs.max_path_bytes]u8 = undefined;
+ const rlen = try s.dir.dir.realPath(io, &real_buf);
+ const root = real_buf[0..rlen];
+ var path_buf: [std.fs.max_path_bytes]u8 = undefined;
+
+ // A live socket: dialPath yields a blocking, close-on-exec stream
+ // (the mounted PROGRAM must not inherit a dialed descriptor).
+ const sock_path = try std.fmt.bufPrintZ(&path_buf, "{s}/live.sock", .{root});
+ var server = try (try Io.net.UnixAddress.init(sock_path)).listen(io, .{});
+ {
+ var st = try dialPath(io, sock_path);
+ const fd = st.socket.handle;
+ try testing.expect(linux.fcntl(fd, linux.F.GETFD, 0) & linux.FD_CLOEXEC != 0);
+ st.close(io);
+ }
+ // The server dies without unlinking: refused, and the caller can tell.
+ server.deinit(io);
+ try testing.expectError(error.Stale, dialPath(io, sock_path));
+
+ // No entry at all.
+ try testing.expectError(error.NotPosted, dialPath(io, try std.fmt.bufPrintZ(&path_buf, "{s}/missing", .{root})));
+
+ // An entry that is not a socket answers the same ECONNREFUSED the
+ // kernel gives a dead socket; the caller tells them apart by stat
+ // before dialing (post.post does exactly that).
+ try testing.expectError(error.Stale, dialPath(io, try std.fmt.bufPrintZ(&path_buf, "{s}", .{root})));
+ const file_path = try std.fmt.bufPrintZ(&path_buf, "{s}/plain", .{root});
+ {
+ var f = try Io.Dir.createFileAbsolute(io, file_path, .{});
+ f.close(io);
+ }
+ try testing.expectError(error.Stale, dialPath(io, file_path));
+
+ // A symlink loop is its own error, not a missing entry.
+ const loop_path = try std.fmt.bufPrintZ(&path_buf, "{s}/loop", .{root});
+ try Io.Dir.symLinkAbsolute(io, loop_path, loop_path, .{});
+ try testing.expectError(error.SymLinkLoop, dialPath(io, loop_path));
+}
+
+test "post postedDir: a 1000-entry directory lists fully or refuses cleanly" {
+ if (builtin.os.tag != .linux) return error.SkipZigTest;
+ var s: Scratch = .{ .dir = undefined };
+ try s.start();
+ defer s.end();
+ const io = testing.io;
+ var real_buf: [std.fs.max_path_bytes]u8 = undefined;
+ const rlen = try s.dir.dir.realPath(io, &real_buf);
+ var dir_buf: [std.fs.max_path_bytes]u8 = undefined;
+ const dir = try std.fmt.bufPrintZ(&dir_buf, "{s}", .{real_buf[0..rlen]});
+ for (0..1000) |i| {
+ var name_buf: [32]u8 = undefined;
+ var path_buf: [std.fs.max_path_bytes]u8 = undefined;
+ const name = try std.fmt.bufPrint(&name_buf, "e{d:0>4}", .{i});
+ var f = try Io.Dir.createFileAbsolute(io, try std.fmt.bufPrintZ(&path_buf, "{s}/{s}", .{ dir, name }), .{});
+ f.close(io);
+ }
+ // A staging buffer that cannot hold every record is an explicit
+ // NoSpace — never a silently truncated listing.
+ var small: [512]u8 = undefined;
+ try testing.expectError(error.NoSpace, postedDir(io, dir, &small));
+ // Big enough lists every entry, with the first and last intact.
+ var big: [16 * 1024]u8 = undefined;
+ var names = try postedDir(io, dir, &big);
+ var count: usize = 0;
+ var have_first = false;
+ var have_last = false;
+ while (names.next()) |name| {
+ count += 1;
+ have_first = have_first or std.mem.eql(u8, name, "e0000");
+ have_last = have_last or std.mem.eql(u8, name, "e0999");
+ }
+ try testing.expectEqual(@as(usize, 1000), count);
+ try testing.expect(have_first and have_last);
+}
diff --git a/src/serve.zig b/src/serve.zig
index 19a8ca2..10c8d01 100644
--- a/src/serve.zig
+++ b/src/serve.zig
@@ -441,6 +441,7 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits
const testing = std.testing;
const c9 = @import("root.zig");
const Client = c9.Client;
+const linux = std.os.linux;
/// A read-only tree in the style of the engine's StubFs: `/index` is a
/// file, `/event` parks its reads until `post()` and `/dir` is a directory.
@@ -936,6 +937,7 @@ test "serve: listenPosted serves the registry name and stop() unposts" {
while (names.next()) |n| seen = seen or std.mem.eql(u8, n, "posted");
try testing.expect(seen);
rig.runner.stop();
+ rig.runner.stop(); // idempotent: the second stop unposts nothing more
try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, path, .{}));
names = try post.posted(io, envp, &stage);
try testing.expect(names.next() == null);
@@ -1013,3 +1015,225 @@ test "serve: listenPosted beside listen(): stop unposts only the registry name"
try Io.Dir.deleteFileAbsolute(io, unix_path);
rig.dir.cleanup();
}
+
+// ---- runner capacity and process hygiene ----
+
+/// Counts the process's open descriptors through /proc, no allocation.
+fn openFdCount(io: Io) !usize {
+ var rb: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined;
+ const dir = try Io.Dir.openDirAbsolute(io, "/proc/self/fd", .{ .iterate = true });
+ defer Io.Dir.close(dir, io);
+ var it = Io.Dir.Reader.init(dir, &rb);
+ var n: usize = 0;
+ while (try it.next(io)) |_| n += 1;
+ return n;
+}
+
+/// Counts the process's OS threads through /proc, no allocation.
+fn taskCount(io: Io) !usize {
+ var rb: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined;
+ const dir = try Io.Dir.openDirAbsolute(io, "/proc/self/task", .{ .iterate = true });
+ defer Io.Dir.close(dir, io);
+ var it = Io.Dir.Reader.init(dir, &rb);
+ var n: usize = 0;
+ while (try it.next(io)) |_| n += 1;
+ return n;
+}
+
+fn waitCount(runner: anytype, n: usize) !void {
+ var tries: usize = 0;
+ while (runner.count() != n) : (tries += 1) {
+ if (tries == 5000) return error.Timeout;
+ try testing.io.sleep(.fromMilliseconds(1), .awake);
+ }
+}
+
+const BigRunner = Runner(Stub, .{ .fid_capacity = 8, .slot_capacity = 4 }, .{ .msize = 4096, .connections = 16, .listeners = 1 });
+
+fn serveNowBig(ctx: ?*anyopaque, conn: *BigRunner.Conn, req: fs.Req) void {
+ const a = Stub.of(ctx).handle(req);
+ conn.reply(&a.reply, a.bytes);
+}
+
+fn countClosedBig(ctx: ?*anyopaque, conn: *BigRunner.Conn) void {
+ _ = conn;
+ const st = Stub.of(ctx);
+ st.mutex.lockUncancelable(testing.io);
+ defer st.mutex.unlock(testing.io);
+ st.closed += 1;
+}
+
+test "serve: sixteen connections fill the table, the seventeenth is closed, stop ends all" {
+ const io = testing.io;
+ var dir = testing.tmpDir(.{});
+ defer dir.cleanup();
+ var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
+ var unix_buffer: [transport.sun_path_len]u8 = undefined;
+ const plen = try dir.dir.realPath(io, &path_buffer);
+ const unix = try std.fmt.bufPrintSentinel(&unix_buffer, "{s}/9p", .{path_buffer[0..plen]}, 0);
+ var stub: Stub = .{};
+ var runner: BigRunner = undefined;
+ runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &stub, .serve = serveNowBig, .closed = countClosedBig } });
+ _ = try runner.listen(.{ .unix = unix }, 16);
+ defer runner.stop();
+
+ var clients: [16]TestClient = undefined;
+ for (&clients) |*tc| {
+ tc.* = .{ .io = undefined, .stream = undefined };
+ try tc.open(io, .{ .unix = unix });
+ try tc.handshake();
+ }
+ defer for (&clients) |*tc| tc.close();
+ try waitCount(&runner, 16);
+
+ // The seventeenth connection is accepted and closed at once: it never
+ // gets a version reply, and the sixteen keep their slots.
+ var extra: TestClient = .{ .io = undefined, .stream = undefined };
+ try extra.open(io, .{ .unix = unix });
+ defer extra.close();
+ try expectHangup(extra.one(.{ .version = .{} }));
+ try waitCount(&runner, 16);
+
+ // stop() hangs every live connection up and frees every slot.
+ const closed0 = stub.closed;
+ runner.stop();
+ try testing.expectEqual(@as(usize, 0), runner.count());
+ try testing.expectEqual(@as(u32, 16), stub.closed - closed0);
+ for (&clients) |*tc| try expectHangup(tc.one(.{ .stat = .{ .fid = 0 } }));
+}
+
+test "serve: repeated connect/disconnect leaves descriptors and threads flat" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start(.{ .serve = serveNow }, 0);
+ defer rig.end();
+ const io = testing.io;
+ const Cycle = struct {
+ fn run(r: *Rig) !void {
+ var tc: TestClient = .{ .io = undefined, .stream = undefined };
+ try tc.open(testing.io, .{ .unix = r.unix });
+ try tc.handshake();
+ try tc.readIndex(1);
+ tc.close();
+ }
+ };
+ // Warm the task pool to its steady state (the first connections grow
+ // it; later ones must not).
+ for (0..400) |_| try Cycle.run(&rig);
+ try rig.settle(0);
+ const fds0 = try openFdCount(io);
+ const threads0 = try taskCount(io);
+ for (0..1200) |_| try Cycle.run(&rig);
+ try rig.settle(0);
+ // Descriptors are exactly flat: no stream, buffer or address is leaked
+ // per connection.
+ try testing.expectEqual(fds0, try openFdCount(io));
+ // Threads come from the shared Io worker pool and may lazily add a
+ // worker; they must not scale with the connection count (which would be
+ // a per-connection thread leak).
+ try testing.expect(try taskCount(io) <= threads0 + 2);
+}
+
+test "serve: a connection storm is absorbed and the runner keeps serving" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start(.{ .serve = serveNow }, 0);
+ defer rig.end();
+ const io = testing.io;
+ const fds0 = try openFdCount(io);
+ const Storm = struct {
+ fn run(r: *Rig, refused: *std.atomic.Value(u32)) void {
+ for (0..32) |_| {
+ var tc: TestClient = .{ .io = undefined, .stream = undefined };
+ tc.open(testing.io, .{ .unix = r.unix }) catch {
+ _ = refused.fetchAdd(1, .monotonic);
+ continue;
+ };
+ // No handshake: the connection is torn down as soon as it
+ // is accepted (or refused a slot by the runner).
+ tc.close();
+ }
+ }
+ };
+ var group: Io.Group = .init;
+ var refused = std.atomic.Value(u32).init(0);
+ for (0..16) |_| try group.concurrent(io, Storm.run, .{ &rig, &refused });
+ try group.await(io);
+ try rig.settle(0);
+ // 512 connections, two slots: the extra ones are closed by the runner,
+ // not refused by the kernel backlog.
+ try testing.expectEqual(@as(u32, 0), refused.load(.monotonic));
+ try testing.expectEqual(fds0, try openFdCount(io));
+ // The runner still accepts and serves a full session.
+ var tc: TestClient = .{ .io = undefined, .stream = undefined };
+ try tc.open(io, .{ .unix = rig.unix });
+ defer tc.close();
+ try tc.handshake();
+ try tc.readIndex(1);
+}
+
+/// The open descriptors of this process, from /proc, for the CLOEXEC audit.
+fn collectFds(io: Io, out: []i32) ![]i32 {
+ var rb: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined;
+ const dir = try Io.Dir.openDirAbsolute(io, "/proc/self/fd", .{ .iterate = true });
+ defer Io.Dir.close(dir, io);
+ var it = Io.Dir.Reader.init(dir, &rb);
+ var n: usize = 0;
+ while (try it.next(io)) |e| {
+ const fd = std.fmt.parseInt(i32, e.name, 10) catch continue;
+ if (n < out.len) {
+ out[n] = fd;
+ n += 1;
+ }
+ }
+ return out[0..n];
+}
+
+fn hasFd(fds: []const i32, fd: i32) bool {
+ for (fds) |f| if (f == fd) return true;
+ return false;
+}
+
+test "serve: every descriptor opened by the post/serve path is close-on-exec" {
+ if (@import("builtin").os.tag != .linux) return error.SkipZigTest;
+ const io = testing.io;
+ var dir = testing.tmpDir(.{});
+ defer dir.cleanup();
+ var real_buf: [std.fs.max_path_bytes]u8 = undefined;
+ const rlen = try dir.dir.realPath(io, &real_buf);
+ var env_buf: [std.fs.max_path_bytes]u8 = undefined;
+ const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{real_buf[0..rlen]});
+ const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null };
+ const envp: post.Env = @ptrCast(&env);
+
+ var before_buf: [64]i32 = undefined;
+ const before = try collectFds(io, &before_buf);
+
+ // Post a runner, dial its name, scan the registry and accept a client:
+ // every descriptor this opens must not survive into an exec.
+ var stub: Stub = .{};
+ var runner: TestRunner = undefined;
+ runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &stub, .serve = serveNow } });
+ defer runner.stop();
+ try runner.listenPosted(envp, "cloexec", 4);
+ var path_buf: [transport.sun_path_len]u8 = undefined;
+ const posted_path = try post.registryPath(envp, "cloexec", &path_buf);
+ var dialed = try post.dialPath(io, posted_path);
+ defer dialed.close(io);
+ var stage: [1024]u8 = undefined;
+ var names = try post.posted(io, envp, &stage);
+ try testing.expect(names.next() != null);
+ var tc: TestClient = .{ .io = undefined, .stream = undefined };
+ try tc.open(io, .{ .unix = posted_path });
+ defer tc.close();
+ try tc.handshake();
+
+ var after_buf: [64]i32 = undefined;
+ const after = try collectFds(io, &after_buf);
+ for (after) |fd| {
+ if (hasFd(before, fd)) continue;
+ const flags = linux.fcntl(fd, linux.F.GETFD, 0);
+ try testing.expect(flags & linux.FD_CLOEXEC != 0);
+ }
+ // Sanity: the audit saw the new descriptors (the listener, the dialed
+ // socket, the accepted connection).
+ try testing.expect(after.len > before.len);
+}