summaryrefslogtreecommitdiff
path: root/src/serve.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-21 14:23:27 -0300
committerGabriel Schneider <[email protected]>2026-09-21 15:20:27 -0300
commit0d7e295efee1fca0935cf4a8bee9629c007dd2b6 (patch)
treed371eb028c63da5ab5b506bd25ac6579cf8a93fc /src/serve.zig
parent3a23f6a29e47ace901bd4d82b9db4055fcc12bb9 (diff)
downloadcloud9-0d7e295efee1fca0935cf4a8bee9629c007dd2b6.tar.gz
cloud9-0d7e295efee1fca0935cf4a8bee9629c007dd2b6.zip
9ns --mntgen: registry subdirectories are mount points too
A registry entry that is a directory is now served the way the root is: a synthetic directory listing the real one, dialing the sockets inside it on walk and recursing into further directories, to max_synth_depth (8) levels across max_synth_dirs (64) synthetic nodes. That is the plan9port mntgen shape and the layout zmx now posts under, so a live session reads at /mnt/9p/zmx/<name>. Before this a directory in the registry was dialed like a socket and answered EIO for good. post gains the two entry points the traversal needs: postedDir (the registry scan, against any directory) and dialPath (a dial by composed path, no name validation). Hardening, each from an attack that broke the code: - BATCH_FORGET carries entries for many owners and puts 0 in the header nodeid, so routing it by the header dropped all of them: 32 of 64 synthetic slots leaked in one close burst and the subdirectories that held them answered EIO forever. distributeForgets unpacks the body and hands each entry to its owner. - probe() and connectBlocking() copied a caller's path into the kernel address with no bound: a path past sun_path overran the 110-byte stack sockaddr (a panic in Debug, silent corruption in ReleaseFast). Both refuse it now, probe as `.live` so a claim never deletes what it could not inspect. - That bound then caught 9proc's own listener, which handed probe() the whole 108-byte sun_path array instead of the path inside it. The probe reads `.live` for anything it cannot ask about, so every stale socket became AlreadyListening and no server could ever take a dead predecessor's name back. It passes the path now. Suites: 87/87 root (+7 post/serve attack regressions), 48/48 9ns, 51+88 9ns integration (+4 traversal and slot-recycling checks), 213/0 9ns adversarial, 60/60 9proc plus its adversarial suites with a new stale-socket takeover check, freestanding green.
Diffstat (limited to 'src/serve.zig')
-rw-r--r--src/serve.zig224
1 files changed, 224 insertions, 0 deletions
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);
+}