diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-21 14:23:27 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-09-21 15:20:27 -0300 |
| commit | 0d7e295efee1fca0935cf4a8bee9629c007dd2b6 (patch) | |
| tree | d371eb028c63da5ab5b506bd25ac6579cf8a93fc /src/serve.zig | |
| parent | 3a23f6a29e47ace901bd4d82b9db4055fcc12bb9 (diff) | |
| download | cloud9-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.zig | 224 |
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); +} |
