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 | |
| 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')
| -rw-r--r-- | src/post.zig | 138 | ||||
| -rw-r--r-- | src/serve.zig | 224 |
2 files changed, 361 insertions, 1 deletions
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); +} |
