From 3a23f6a29e47ace901bd4d82b9db4055fcc12bb9 Mon Sep 17 00:00:00 2001 From: Gabriel Schneider Date: Mon, 21 Sep 2026 14:13:43 -0300 Subject: post registry + 9ns --mntgen: the /srv translation cloud9.post: servers post their socket under a name in $XDG_RUNTIME_DIR/9p (post/unpost, posted, dial, Watch) and serve.Runner.listenPosted posts a server by name, unposting on stop. Names are budget-checked against the 108-byte socket path; a claim binds+listens at a private temp path and takes the name with atomic renames under flock (RENAME_NOREPLACE for free names, RENAME_EXCHANGE grab-verify-commit for stale ones): the registry path is never unlinked by a claim, live names refuse with AlreadyPosted, foreign files with NotSocket, and unpost removes only the caller's inode-matched entry. Watch surfaces inotify overflow and a replaced registry dir. 9ns --mntgen [--mount DIR] -- PROGRAM: one FUSE mount at /mnt/9p whose synthetic root lists the posted registry (no connection made); a walk into an unmounted name dials it and runs the existing bridge dispatch in a per-server worker thread, routed by mount index in the node id's top bits (ordinals never reused, cap 4096); a dead server answers EIO on its subtree and is re-dialed on the next walk. The dial watches stop_fd through Tversion (connectWatched). All existing 9ns forms are unchanged. 9proc's unix listener no longer blind-unlinks its path: a foreign non-socket is refused (Occupied), a live server is refused (AlreadyListening), only a refused socket is cleared, and stop() unlinks only the listener's own inode-matched socket. Hardened by adversarial review (GLM 5.3 x2 + DeepSeek V4.1 Flash, all high-thinking): double-bind races on one name (0 in 180k rounds), foreign-file TOCTOU deletions (0 in 4M flips), a 255-byte-name listing panic, inotify queue overflow silently dropped, listenPosted silently overwriting, dial-time Tversion hangs wedging the dispatcher, --debug silently ignored in mntgen, and xattr/statx probes answering EPERM on the synthetic root (broke `ls -l /mnt/9p`). Tests: root 80/80, 9ns 47/47, 9proc 60/60, integration 88/88 + mntgen 37/37, adversarial 213/0, freestanding riscv32 gate green. --- src/post.zig | 1024 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/root.zig | 4 + src/serve.zig | 175 +++++++++- 3 files changed, 1199 insertions(+), 4 deletions(-) create mode 100644 src/post.zig (limited to 'src') diff --git a/src/post.zig b/src/post.zig new file mode 100644 index 0000000..47a7d3d --- /dev/null +++ b/src/post.zig @@ -0,0 +1,1024 @@ +//! The `/srv` translation: a server posts itself under a name in one +//! per-user registry directory, and clients list the names and dial them +//! (`post9pservice` in plan9port, `devsrv.c` in Plan 9). The registry is +//! `$XDG_RUNTIME_DIR/9p/` — per-user tmpfs, the right lifetime — and a +//! name's socket lives at `$XDG_RUNTIME_DIR/9p/`. Mounting and +//! namespace policy stay with the client (9ns `--mntgen`); this module +//! only builds paths, posts, lists, dials and watches. +//! +//! House style as everywhere in cloud9: no assumed allocator, buffers +//! are the caller's, the listing is staged records, and the environment +//! is passed in as the raw block `main` receives (`post.Env`) — the +//! library never reads the process environment on its own. The Linux +//! syscalls (probe, dial, `Watch`) are raw and need no libc; the rest +//! goes through `std.Io`. See docs/design.md, "Post registry". +const std = @import("std"); +const builtin = @import("builtin"); +const Io = std.Io; +const linux = std.os.linux; +const transport = @import("transport.zig"); + +/// The raw environment block the C startup hands to `main`, the same +/// shape 9ns threads to `ns.getenv`. `null` terminates. +pub const Env = [*:null]const ?[*:0]const u8; + +/// The kernel's Unix socket path budget (`sun_path`). +pub const sun_path_len = transport.sun_path_len; + +/// The longest name a post may carry: the socket path +/// `$XDG_RUNTIME_DIR/9p/` must stay inside `sun_path_len`, and a +/// conforming `$XDG_RUNTIME_DIR` (`/run/user/`) is budgeted at 32 +/// bytes, leaving the rest of the budget for `/9p/`, the name and the +/// terminating zero. `registryPath` checks the real prefix again. +pub const max_name_len = sun_path_len - "/9p/".len - 1 - 32; + +comptime { + if (max_name_len < 8) @compileError("9P post name budget is too small to be useful"); +} + +/// One entry of a raw environment block. +pub fn getenv(env: Env, name: []const u8) ?[]const u8 { + var i: usize = 0; + while (env[i]) |entry| : (i += 1) { + const kv = std.mem.span(entry); + const eq = std.mem.indexOfScalar(u8, kv, '=') orelse continue; + if (std.mem.eql(u8, kv[0..eq], name)) return kv[eq + 1 ..]; + } + return null; +} +/// Is `name` legal for posting? It mirrors the engine's `legalName` +// (empty, "." and ".." refused, no '/' or NUL) with the post-specific +// budget: path traversal through a posted name is THE attack, and the +// name must also fit a Unix socket path. +pub fn legalName(name: []const u8) bool { + if (name.len == 0 or name.len > max_name_len) return false; + if (std.mem.eql(u8, name, ".") or std.mem.eql(u8, name, "..")) return false; + return std.mem.indexOfAny(u8, name, "/\x00") == null; +} + +pub const PathError = error{ + /// `$XDG_RUNTIME_DIR` is unset. There is no `/tmp` fallback. + NotRuntimeDir, + /// The name cannot be posted (see `legalName`). + IllegalName, + /// The caller's buffer is too small for the path and its zero. + NoSpace, + /// The full socket path would not fit `sun_path_len`. + NameTooLong, +}; + +const xdg_runtime_dir = "XDG_RUNTIME_DIR"; + +/// Writes `$XDG_RUNTIME_DIR/9p` into `out` and returns it zero-terminated. +pub fn registryDir(env: Env, out: []u8) PathError![:0]u8 { + const root = getenv(env, xdg_runtime_dir) orelse return error.NotRuntimeDir; + if (root.len == 0) return error.NotRuntimeDir; + const total = root.len + "/9p".len; + if (total + 1 > out.len) return error.NoSpace; + @memcpy(out[0..root.len], root); + @memcpy(out[root.len..total], "/9p"); + out[total] = 0; + return out[0..total :0]; +} + +/// Writes `$XDG_RUNTIME_DIR/9p/` into `out` and returns it +/// zero-terminated, ready for `probe`, `dial` or `Io.net.UnixAddress`. +pub fn registryPath(env: Env, name: []const u8, out: []u8) PathError![:0]u8 { + if (!legalName(name)) return error.IllegalName; + const root = getenv(env, xdg_runtime_dir) orelse return error.NotRuntimeDir; + if (root.len == 0) return error.NotRuntimeDir; + const total = root.len + "/9p/".len + name.len; + if (total >= sun_path_len) return error.NameTooLong; + if (total + 1 > out.len) return error.NoSpace; + @memcpy(out[0..root.len], root); + @memcpy(out[root.len..][0.."/9p/".len], "/9p/"); + @memcpy(out[total - name.len ..][0..name.len], name); + out[total] = 0; + return out[0..total :0]; +} + +/// The raw registry entries, staged by `posted` as `len:u8 name` +/// records back to back in the caller's buffer (the engine's staged- +/// record pattern for readdir). Nothing is dialed; staleness is the +/// caller's concern (`probe`). +pub const Names = struct { + bytes: []const u8, + i: usize = 0, + + pub fn next(n: *Names) ?[]const u8 { + if (n.i + 1 > n.bytes.len) return null; + const len = n.bytes[n.i]; + if (n.i + 1 + len > n.bytes.len) return null; + const name = n.bytes[n.i + 1 ..][0..len]; + n.i += 1 + @as(usize, len); + return name; + } +}; + +pub const PostedError = PathError || error{NoSpace} || Io.Dir.OpenError || Io.Dir.Reader.Error; + +/// Lists the registry into `out` and returns an iterator over the names. +/// Nothing is dialed; a missing registry lists as empty (no server has +/// posted yet). The iteration buffer is this stack frame's, so only the +/// staged names outlive the call. +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); + 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, + }; + defer Io.Dir.close(dir, io); + var read_buf: [Io.Dir.Iterator.reader_buffer_len]u8 align(@alignOf(usize)) = undefined; + var reader = Io.Dir.Reader.init(dir, &read_buf); + var n: usize = 0; + while (try reader.next(io)) |entry| { + if (n + 1 + entry.name.len > out.len) return error.NoSpace; + out[n] = @intCast(entry.name.len); + @memcpy(out[n + 1 ..][0..entry.name.len], entry.name); + n += 1 + entry.name.len; + } + return .{ .bytes = out[0..n] }; +} + +/// What `probe` found at a registry path. Nothing is modified: a stale +/// entry is reported, not removed (`post` owns replacement). +pub const Probe = enum { none, stale, live }; + +/// Probes a registry entry without touching it: connect to the socket +/// path without blocking. Refused means the path holds nothing to talk +/// to — a dead server's socket *or an entry that is not a socket at +/// all*, which the kernel answers with the same ECONNREFUSED; `post` +/// tells them apart by stat and never deletes the latter. No entry at +/// all is `none`; connected — or busy, or any unexpected problem — is +/// `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 { + const fd = socketNonblocking() catch return .live; + defer _ = linux.close(fd); + var addr: linux.sockaddr.un = .{ .path = @splat(0) }; + @memcpy(addr.path[0..path.len], path); + const rc: isize = @bitCast(linux.connect(fd, @ptrCast(&addr), @sizeOf(linux.sockaddr.un))); + return switch (rawErrno(rc)) { + .SUCCESS, .AGAIN, .INPROGRESS, .PERM, .ACCES => .live, + .NOENT, .NOTDIR => .none, + .CONNREFUSED => .stale, + else => .live, + }; +} + +fn socketNonblocking() !i32 { + const rc = linux.socket(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.NONBLOCK | linux.SOCK.CLOEXEC, 0); + if (rawErrno(@bitCast(rc)) != .SUCCESS) return error.Socket; + return @intCast(rc); +} + +/// The errno of a raw Linux syscall return, which encodes it as the +/// negated value (0 on success). The std.Io Unix connect does not +/// promise ECONNREFUSED, so `probe` and `dial` speak to the kernel +/// directly. +fn rawErrno(rc: isize) linux.E { + const n: usize = if (rc < 0 and -rc < 4096) @intCast(-rc) else 0; + return @enumFromInt(n); +} + +pub const DialError = error{ + NotRuntimeDir, + IllegalName, + NoSpace, + NameTooLong, + /// No registry entry under this name. + NotPosted, + /// The entry exists but the connection is refused: it is stale. + Stale, + SystemResources, + ProcessFdQuotaExceeded, + SystemFdQuotaExceeded, + AccessDenied, + PermissionDenied, + SymLinkLoop, + NotDir, + Io, + Unexpected, +}; + +/// Dials the server posted under `name` and returns its stream. The +/// descriptor is blocking and close-on-exec; it is not the registry's — +/// closing it changes nothing in the registry. On Linux the connect is +/// a raw syscall (so ECONNREFUSED is distinguishable, which `std.Io`'s +/// Unix connect does not promise). +pub fn dial(io: Io, env: Env, name: []const u8) DialError!Io.net.Stream { + var path_buf: [sun_path_len]u8 = undefined; + const path = try registryPath(env, name, &path_buf); + 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, + }; + _ = io; + return .{ .socket = .{ .handle = fd, .address = .{ .ip4 = .loopback(0) } } }; +} + +const ConnectError = error{ Socket, Noent, Refused, AccessDenied, Loop, NotDir }; + +fn connectBlocking(path: [:0]const u8) ConnectError!i32 { + 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); + errdefer _ = linux.close(fd); + var addr: linux.sockaddr.un = .{ .path = @splat(0) }; + @memcpy(addr.path[0..path.len], path); + const crc: isize = @bitCast(linux.connect(fd, @ptrCast(&addr), @sizeOf(linux.sockaddr.un))); + switch (rawErrno(crc)) { + .SUCCESS => return fd, + .NOENT, .NOTDIR => return error.Noent, + .CONNREFUSED => return error.Refused, + .ACCES, .PERM => return error.AccessDenied, + .LOOP => return error.Loop, + else => return error.Socket, + } +} + +/// What `post` bound, and what the caller owes `unpost`: the registry +/// path and the inode of the bound entry, so a late `unpost` can prove +/// the entry is still its own before unlinking it. +pub const Posted = struct { + server: Io.net.Server, + path: [:0]const u8, + inode: u64, +}; + +pub const PostError = error{ + NotRuntimeDir, + IllegalName, + NoSpace, + NameTooLong, + /// A live server owns the name. + AlreadyPosted, + /// The registry entry exists but is not a socket. It is never + /// deleted; whoever put it there must remove it themselves. + NotSocket, +} || Io.net.UnixAddress.ListenError || Io.Dir.CreateDirPathError || Io.Dir.StatFileError || + Io.Dir.DeleteFileError || Io.UnexpectedError || Io.Cancelable; + +/// How many rounds of the claim protocol `post` runs before giving the +/// name up as contested. +const post_attempts = 8; + +/// Distinct temp socket names for posts sharing one pid (threads). +var temp_serial: std.atomic.Value(u64) = .init(0); + +/// Posts a listening socket under `name`: the registry directory is +/// created (0o750; already present is fine), the name's path — written +/// into `path_buf` — is bound and returned for `unpost`. +/// +/// The protocol: a live server owning the name is `AlreadyPosted`; an +/// entry that is not a socket is `NotSocket` and is never deleted. The +/// claiming socket is bound and *listening* at a private temp path +/// first (so any probe of it answers live — no window in which the +/// claim could be mistaken for a stale entry), and the name is taken +/// under the stale-removal lock by atomic `renameat2` calls only +/// (`claimName`): of two posts racing on one stale name exactly one +/// ends up owning the entry, the loser re-runs into `AlreadyPosted`, +/// and the registry path is never unlinked by a claim — what a claim +/// displaces is parked under a private temp name and left alone +/// unless provably ours or provably the dead entry. +pub fn post(io: Io, env: Env, name: []const u8, backlog: u31, path_buf: []u8) PostError!Posted { + const path = try registryPath(env, name, path_buf); + var dir_buf: [std.fs.max_path_bytes]u8 = undefined; + const dir_path = try registryDir(env, &dir_buf); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, dir_path, .fromMode(0o750)); + const parent = dir_path[0 .. dir_path.len - "/9p".len]; + const serial = temp_serial.fetchAdd(1, .monotonic); + var tmp_buf: [std.fs.max_path_bytes]u8 = undefined; + const tmp = std.fmt.bufPrintZ(&tmp_buf, "{s}/.post.sock.{d}.{d}", .{ parent, linux.getpid(), serial }) catch return error.NoSpace; + var attempts: usize = 0; + var saw_not_socket = false; + claim: while (attempts < post_attempts) : (attempts += 1) { + // A live listener at a private temp path, invisible to registry + // listings, before the name is even looked at. + const addr = try Io.net.UnixAddress.init(tmp); + var server = addr.listen(io, .{ .kernel_backlog = backlog }) catch |err| switch (err) { + error.AddressInUse => { + // A crashed former self under a reused pid, or stale + // garbage: the dead temp is ours to clear. + _ = linux.unlink(tmp.ptr); + continue :claim; + }, + else => { + _ = linux.unlink(tmp.ptr); + return err; + }, + }; + switch (claimName(io, dir_path, path, tmp)) { + .won => { + // Ours, and live: nothing in the protocol displaces a + // live entry. Record the inode `unpost` checks — and + // only from a socket: whatever a foreign hand may have + // put in the entry's place between the claim and here + // is not ours to name. + const st = Io.Dir.statFile(.cwd(), io, path, .{}) catch |err| switch (err) { + error.FileNotFound => { + // Unlinked behind our back: not ours now. (The + // temp is already spent; whatever a foreign hand + // parked there was deliberately left alone.) + server.deinit(io); + continue :claim; + }, + else => { + server.deinit(io); + return err; + }, + }; + if (st.kind != .unix_domain_socket) { + server.deinit(io); + continue :claim; + } + return .{ .server = server, .path = path, .inode = st.inode }; + }, + .contended => { + _ = linux.unlink(tmp.ptr); + server.deinit(io); + return error.AlreadyPosted; + }, + .not_socket => { + saw_not_socket = true; + _ = linux.unlink(tmp.ptr); + server.deinit(io); + // Another post's stale-swap parks a dummy here for a + // moment; give the dust a beat to settle before the + // entry is declared bogus. + io.sleep(.fromMilliseconds(1), .awake) catch {}; + continue :claim; + }, + .retry => { + _ = linux.unlink(tmp.ptr); + server.deinit(io); + continue :claim; + }, + } + } + return if (saw_not_socket) error.NotSocket else error.AlreadyPosted; +} + +/// The outcome of one claim round. +const Claim = enum { + /// The live listener now sits at the registry path. + won, + /// A live server owns the name. + contended, + /// The registry entry is not a socket (never deleted). + not_socket, + /// The landscape raced; look again. + retry, +}; + +/// Takes the registry path for the listener bound at `sock_tmp`, under +/// the stale-removal lock: +/// +/// * No entry: an atomic `renameat2(RENAME_NOREPLACE)` claims it — two +/// posts race to exactly one winner. +/// * A live entry: `contended`. +/// * A non-socket entry: `not_socket` — it stays strictly alone. +/// * A stale entry: it is grabbed with `renameat2(RENAME_EXCHANGE)` +/// against a private dummy and verified by inode *and* a fresh probe +/// (a live server that came up since the probe is swapped back, +/// untouched), and only then is the live listener exchanged into the +/// entry's place. The registry path is never unlinked here — the +/// dummy and the dead entry die under private temp names, provably +/// by inode, and anything a foreign hand parked in their place is +/// left alone. +fn claimName(io: Io, registry_dir: [:0]const u8, path: [:0]const u8, sock_tmp: [:0]const u8) Claim { + const lock_fd = lockStaleRemoval(registry_dir); + defer { + if (lock_fd >= 0) _ = linux.close(lock_fd); + } + const parent = registry_dir[0 .. registry_dir.len - "/9p".len]; + + const st = Io.Dir.statFile(.cwd(), io, path, .{}) catch |err| switch (err) { + error.FileNotFound => { + // A fresh name: the rename is the arbiter. + const rc: isize = @bitCast(linux.renameat2(linux.AT.FDCWD, sock_tmp.ptr, linux.AT.FDCWD, path.ptr, .{ .NOREPLACE = true })); + return switch (rawErrno(rc)) { + .SUCCESS => .won, + .EXIST => .retry, + .INVAL, .NOSYS, .PERM, .OPNOTSUPP => legacyClaim(path, sock_tmp), + else => .retry, + }; + }, + else => return .retry, + }; + if (st.kind != .unix_domain_socket) return .not_socket; + switch (probe(path)) { + .live => return .contended, + .none, .stale => {}, + } + + // A dead entry. Grab it for a look: the dummy never enters the + // registry directory, so no listing ever sees it. + const dummy_serial = temp_serial.fetchAdd(1, .monotonic); + var dummy_buf: [std.fs.max_path_bytes]u8 = undefined; + const dummy = std.fmt.bufPrintZ(&dummy_buf, "{s}/.post.tmp.{d}.{d}", .{ parent, linux.getpid(), dummy_serial }) catch return .retry; + var dummy_file = Io.Dir.createFileAbsolute(io, dummy, .{ .exclusive = true, .truncate = false }) catch return .retry; + dummy_file.close(io); + const dummy_st = Io.Dir.statFile(.cwd(), io, dummy, .{}) catch { + _ = linux.unlink(dummy.ptr); + return .retry; + }; + + const rc: isize = @bitCast(linux.renameat2(linux.AT.FDCWD, path.ptr, linux.AT.FDCWD, dummy.ptr, .{ .EXCHANGE = true })); + switch (rawErrno(rc)) { + .NOENT => { + // The entry unposted meanwhile; look again. + unlinkIfOurs(io, dummy, dummy_st.inode); + return .retry; + }, + .INVAL, .NOSYS, .PERM, .OPNOTSUPP => { + // No RENAME_EXCHANGE on this filesystem: the grab cannot + // be made safe here; the documented fallback is an + // inode-checked in-place unlink (the only place a claim + // may delete a raced non-socket, on filesystems without + // the primitive). + unlinkIfOurs(io, dummy, dummy_st.inode); + const now = Io.Dir.statFile(.cwd(), io, path, .{}) catch return .retry; + if (now.inode != st.inode) return .retry; + Io.Dir.deleteFileAbsolute(io, path) catch return .retry; + return switch (claimFresh(path, sock_tmp)) { + .won => .won, + else => .retry, + }; + }, + .SUCCESS => {}, + else => { + unlinkIfOurs(io, dummy, dummy_st.inode); + return .retry; + }, + } + // The dummy sits at the entry's place; the grabbed entry is at + // `dummy`. Verify it is still the dead one we probed. + const grabbed = Io.Dir.statFile(.cwd(), io, dummy, .{}) catch { + _ = swapBack(path, dummy); + unlinkIfOurs(io, dummy, dummy_st.inode); + return .retry; + }; + if (grabbed.inode != st.inode or probe(dummy) == .live) { + // A live server took the name between the probe and the grab: + // restore it, untouched. + const back = swapBack(path, dummy); + if (rawErrno(back) == .SUCCESS) unlinkIfOurs(io, dummy, dummy_st.inode); + return if (grabbed.inode != st.inode) .retry else .contended; + } + // Verified dead: exchange the live listener into the entry's + // place. What the exchange parks at `sock_tmp` is our dummy — + // unless a foreign hand replaced it in the meantime, in which + // case it is left alone, intact. + const claim_rc: isize = @bitCast(linux.renameat2(linux.AT.FDCWD, sock_tmp.ptr, linux.AT.FDCWD, path.ptr, .{ .EXCHANGE = true })); + if (rawErrno(claim_rc) != .SUCCESS) { + _ = swapBack(path, dummy); + unlinkIfOurs(io, dummy, dummy_st.inode); + return .retry; + } + unlinkIfOurs(io, sock_tmp, dummy_st.inode); + // The dead entry dies under the private temp name, provably by + // inode; anything a foreign hand parked there survives. + unlinkIfOurs(io, dummy, grabbed.inode); + return .won; +} + +/// The NOREPLACE rename for a provably free path. +fn claimFresh(path: [:0]const u8, sock_tmp: [:0]const u8) Claim { + const rc: isize = @bitCast(linux.renameat2(linux.AT.FDCWD, sock_tmp.ptr, linux.AT.FDCWD, path.ptr, .{ .NOREPLACE = true })); + return switch (rawErrno(rc)) { + .SUCCESS => .won, + else => .retry, + }; +} + +/// The claim on a filesystem without `renameat2` flags at all: a plain +/// rename over a provably empty path (the replace window is the price +/// of such a filesystem). +fn legacyClaim(path: [:0]const u8, sock_tmp: [:0]const u8) Claim { + if (probe(path) != .none) return .retry; + const rc: isize = @bitCast(linux.rename(sock_tmp.ptr, path.ptr)); + return if (rawErrno(rc) == .SUCCESS) .won else .retry; +} + +/// Unlinks `p` only while it still holds inode `ino`: a file some +/// other hand has put in a private temp's place is left alone (a +/// leaked dotfile beats deleting what may be someone's file). +fn unlinkIfOurs(io: Io, p: [:0]const u8, ino: u64) void { + const st = Io.Dir.statFile(.cwd(), io, p, .{}) catch return; + if (st.inode != ino) return; + Io.Dir.deleteFileAbsolute(io, p) catch {}; +} + +/// Puts a swapped-out entry back: `path` and `tmp` re-exchange, so the +/// displaced entry returns to the registry and our dummy comes back to +/// the temp name. The result is ignored: the removal is serialized by +/// `lockStaleRemoval`, so the entry can only fail to return when the +/// path was cleared by hand in the meantime. +fn swapBack(path: [:0]const u8, tmp: [:0]const u8) isize { + return @bitCast(linux.renameat2(linux.AT.FDCWD, path.ptr, linux.AT.FDCWD, tmp.ptr, .{ .EXCHANGE = true })); +} + +/// `flock(LOCK_EX)` on the registry-side `.post.lock`, serializing the +/// stale-removal window between posts (including threads: the lock is +/// taken per open file description). The file is an empty dotfile; the +/// kernel drops the lock when the holder dies. A lock that cannot be +/// taken does not block posting — the window just loses its +/// serialization. +fn lockStaleRemoval(registry_dir: [:0]const u8) i32 { + var buf: [std.fs.max_path_bytes]u8 = undefined; + const lock_path = std.fmt.bufPrintZ(&buf, "{s}/.post.lock", .{registry_dir[0 .. registry_dir.len - "/9p".len]}) catch return -1; + const fd: isize = @bitCast(linux.open(lock_path.ptr, .{ .ACCMODE = .RDWR, .CREAT = true, .CLOEXEC = true }, 0o600)); + if (rawErrno(fd) != .SUCCESS) return -1; + const lfd: i32 = @intCast(fd); + while (rawErrno(@bitCast(linux.flock(lfd, lock_ex))) == .INTR) {} + return lfd; +} + +/// `LOCK_EX` for the raw `flock` call (the constant is not in std's +/// linux namespace). +const lock_ex: i32 = 2; + +/// Unposts: unlinks the registry path — but only the caller's own +/// entry. `inode` is the one `post` bound (in `Posted`); a path whose +/// entry has been replaced (the socket lost and the name re-posted by +/// another server) is left strictly alone, so one server's late stop +/// can never unpost another's live name. Idempotent; all errors are +/// swallowed — an unpost must never be the reason a server fails to +/// shut down. +pub fn unpost(io: Io, path: [:0]const u8, inode: u64) void { + const st = Io.Dir.statFile(.cwd(), io, path, .{}) catch return; + if (st.inode != inode) return; + Io.Dir.deleteFileAbsolute(io, path) catch {}; +} + +/// An inotify watcher on the registry directory, for hosts (9ns's +/// mntgen) that cache its listing: `init`, `add` the directory, `next` +/// yields names as they are added and removed. Hosted Linux only; the +/// rest of `post` is platform-neutral. +pub const Watch = if (builtin.os.tag == .linux) struct { + const Self = @This(); + + fd: i32, + /// Event staging, so `next` never truncates a record. + stage: [4096]u8 = undefined, + head: usize = 0, + tail: usize = 0, + + /// The kind of change. `overflow` (empty name) reports that the + /// kernel dropped events — the queue overflowed; a caching consumer + /// must rescan. `gone` (empty name) reports that the watch itself + /// ended (the directory was removed, or the kernel dropped the + /// watch): re-`add` it and rescan; until then `next` is null. + pub const Kind = enum { added, removed, overflow, gone }; + pub const Event = struct { name: []const u8, kind: Kind }; + + pub const Error = error{ SystemResources, ProcessFdQuotaExceeded, Io, Unexpected }; + + pub fn init() Error!Self { + const rc = linux.inotify_init1(linux.IN.CLOEXEC | linux.IN.NONBLOCK); + return switch (rawErrno(@bitCast(rc))) { + .SUCCESS => .{ .fd = @intCast(rc) }, + .NOMEM => error.SystemResources, + .MFILE, .NFILE => error.ProcessFdQuotaExceeded, + else => error.Io, + }; + } + + /// Watches the directory for names appearing and disappearing. + pub fn add(w: *Self, dir_path: [:0]const u8) Error!void { + const mask = linux.IN.CREATE | linux.IN.DELETE | linux.IN.MOVED_TO | linux.IN.MOVED_FROM; + const rc = linux.inotify_add_watch(w.fd, dir_path.ptr, mask); + switch (rawErrno(@bitCast(rc))) { + .SUCCESS => {}, + .NOMEM => return error.SystemResources, + else => return error.Io, + } + } + /// Yields the next registry change, or null when nothing is + /// pending. `Event.name` points into the watch's staging and is + /// valid until the following `next`. Events for anything but a name + /// appearing or disappearing (directories, metadata) are skipped — + /// except the two a caching consumer must not miss: a queue + /// overflow (`Kind.overflow`) and the watch ending + /// (`Kind.gone`). + pub fn next(w: *Self) Error!?Event { + while (true) { + if (w.head + @sizeOf(linux.inotify_event) <= w.tail) { + const ev: *const linux.inotify_event = @ptrCast(@alignCast(w.stage[w.head..].ptr)); + const total = @sizeOf(linux.inotify_event) + ev.len; + if (w.head + total <= w.tail) { + const kind: ?Kind = if (ev.mask & linux.IN.Q_OVERFLOW != 0) + // Events were dropped: the consumer's cache may + // be arbitrarily wrong and must rescan. + .overflow + else if (ev.mask & linux.IN.IGNORED != 0) + // The watch itself ended (the directory was + // removed); nothing further will be reported. + .gone + else if (ev.mask & linux.IN.ISDIR != 0) + null + else if (ev.mask & (linux.IN.CREATE | linux.IN.MOVED_TO) != 0) + .added + else if (ev.mask & (linux.IN.DELETE | linux.IN.MOVED_FROM) != 0) + .removed + else + null; + const name = if (ev.getName()) |n| n else ""; + w.head += total; + if (kind) |k| return .{ .name = name, .kind = k }; + continue; + } + } + // Not enough staged: keep the tail, refill. + if (w.head > 0) { + const left = w.tail - w.head; + std.mem.copyForwards(u8, w.stage[0..left], w.stage[w.head..w.tail]); + w.head = 0; + w.tail = left; + } + const rc: isize = @bitCast(linux.read(w.fd, @as([*]u8, &w.stage) + w.tail, w.stage.len - w.tail)); + switch (rawErrno(rc)) { + .SUCCESS => w.tail += @intCast(rc), + .AGAIN => return null, // nothing pending + .INTR => continue, + else => return error.Io, + } + } + } + + pub fn deinit(w: *Self) void { + _ = linux.close(w.fd); + } +} else @compileError("post.Watch requires Linux inotify"); + +// ---- tests ---- + +const testing = std.testing; + +/// A scratch registry: a fake environment block whose XDG_RUNTIME_DIR is +/// a per-test temporary directory. +const Scratch = struct { + dir: testing.TmpDir, + path_buf: [std.fs.max_path_bytes]u8 = undefined, + env_buf: [std.fs.max_path_bytes]u8 = undefined, + env: [2]?[*:0]const u8 = undefined, + + fn start(s: *Scratch) !void { + const io = testing.io; + s.dir = testing.tmpDir(.{}); + errdefer s.dir.cleanup(); + const len = try s.dir.dir.realPath(io, &s.path_buf); + const value = try std.fmt.bufPrintZ(&s.env_buf, "XDG_RUNTIME_DIR={s}", .{s.path_buf[0..len]}); + s.env[0] = @ptrCast(value.ptr); + s.env[1] = null; + } + + fn end(s: *Scratch) void { + s.dir.cleanup(); + } + + fn envp(s: *Scratch) Env { + return @ptrCast(&s.env); + } +}; + +test "post getenv: the block is scanned, keys match whole" { + const env = [_:null]?[*:0]const u8{ "A=1", "XDG_RUNTIME_DIR=/run/user/1000", "X=", "XDG=no" }; + try testing.expectEqualStrings("/run/user/1000", getenv(@ptrCast(&env), "XDG_RUNTIME_DIR").?); + try testing.expectEqualStrings("no", getenv(@ptrCast(&env), "XDG").?); + try testing.expect(getenv(@ptrCast(&env), "NOPE") == null); + try testing.expect(getenv(@ptrCast(&env), "PAT") == null); + try testing.expect(getenv(@ptrCast(&env), "A") != null); +} + +test "post legalName: traversal, empties and the socket budget" { + try testing.expect(legalName("demo")); + try testing.expect(legalName("...")); + try testing.expect(legalName("a" ** max_name_len)); + // Empty and dot names are not postable. + try testing.expect(!legalName("")); + try testing.expect(!legalName(".")); + try testing.expect(!legalName("..")); + // Path traversal is THE attack: no separators, no NUL. + try testing.expect(!legalName("a/b")); + try testing.expect(!legalName("/")); + try testing.expect(!legalName("..\x00..")); + try testing.expect(!legalName("a\x00b")); + // The name must fit a Unix socket path. + try testing.expect(!legalName("a" ** (max_name_len + 1))); + try testing.expect(!legalName("a" ** 255)); +} + +test "post registryPath: unset XDG is refused, no /tmp fallback" { + const empty = [_:null]?[*:0]const u8{"PATH=/bin"}; + var buf: [sun_path_len]u8 = undefined; + try testing.expectError(error.NotRuntimeDir, registryPath(@ptrCast(&empty), "demo", &buf)); + try testing.expectError(error.NotRuntimeDir, registryDir(@ptrCast(&empty), &buf)); + const empty_value = [_:null]?[*:0]const u8{"XDG_RUNTIME_DIR="}; + try testing.expectError(error.NotRuntimeDir, registryPath(@ptrCast(&empty_value), "demo", &buf)); +} + +test "post registryPath: shape, budget and caller buffer" { + const env = [_:null]?[*:0]const u8{"XDG_RUNTIME_DIR=/run/user/1000"}; + const e: Env = @ptrCast(&env); + var buf: [sun_path_len]u8 = undefined; + const p = try registryPath(e, "demo", &buf); + try testing.expectEqualStrings("/run/user/1000/9p/demo", p); + try testing.expectEqual(@as(u8, 0), buf[p.len]); + var dbuf: [32]u8 = undefined; + try testing.expectEqualStrings("/run/user/1000/9p", try registryDir(e, &dbuf)); + // Illegal names never build a path at all. + try testing.expectError(error.IllegalName, registryPath(e, "a/b", &buf)); + try testing.expectError(error.IllegalName, registryPath(e, "..", &buf)); + // The 108-byte sockaddr budget, checked against the real prefix. + const long_env = [_:null]?[*:0]const u8{"XDG_RUNTIME_DIR=/tmp/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}; + const name = "n" ** max_name_len; // legal against the 32-byte XDG budget + try testing.expectError(error.NameTooLong, registryPath(@ptrCast(&long_env), name, &buf)); + // Too small a caller buffer is NoSpace, not a truncation. + var tiny: [4]u8 = undefined; + try testing.expectError(error.NoSpace, registryPath(e, "demo", &tiny)); +} + +test "post stale protocol: fresh, stale, live and non-socket entries" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var pbuf: [sun_path_len]u8 = undefined; + + // A fresh post binds and creates the 0o750 registry directory. + var first = try post(io, env, "demo", 4, &pbuf); + // While the listener is open the name is live and owned. + try testing.expect(probe(first.path) == .live); + try testing.expectError(error.AlreadyPosted, post(io, env, "demo", 4, &pbuf)); + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(env, &dbuf); + const st = try Io.Dir.statFile(.cwd(), io, reg, .{}); + try testing.expect(st.kind == .directory); + try testing.expect(st.permissions.toMode() & 0o750 == 0o750); + + // The listener dies without unlinking: the entry is stale, and a + // new post unlinks and replaces it. + first.server.deinit(io); + try testing.expect(probe(first.path) == .stale); + var revived = try post(io, env, "demo", 4, &pbuf); + try testing.expect(probe(revived.path) == .live); + unpost(io, revived.path, revived.inode); + revived.server.deinit(io); + + // A non-socket entry is refused, never deleted. + const plain = try registryPath(env, "plain", &pbuf); + { + var f = try Io.Dir.createFileAbsolute(io, plain, .{}); + f.close(io); + } + try testing.expectError(error.NotSocket, post(io, env, "plain", 4, &pbuf)); + const after = try Io.Dir.statFile(.cwd(), io, plain, .{}); + try testing.expect(after.kind == .file); + + // dial tells the three apart. + try testing.expectError(error.NotPosted, dial(io, env, "missing")); + try testing.expectError(error.Stale, dial(io, env, "plain")); + var live = try post(io, env, "dials", 4, &pbuf); + var stream = try dial(io, env, "dials"); + stream.close(io); + // The listener dies without unposting: the entry stays, stale. + live.server.deinit(io); + try testing.expectError(error.Stale, dial(io, env, "dials")); + unpost(io, live.path, live.inode); +} + +test "post posted: listing stages the raw names, missing lists empty" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var stage: [512]u8 = undefined; + var pbuf: [sun_path_len]u8 = undefined; + + // Before anything posts (or even creates the registry): empty. + var names = try posted(io, env, &stage); + try testing.expect(names.next() == null); + + // Two posts and a decoy regular file: the listing is raw. + var a = try post(io, env, "alpha", 4, &pbuf); + var b = try post(io, env, "beta", 4, &pbuf); + const decoy = try registryPath(env, "decoy", &pbuf); + { + var f = try Io.Dir.createFileAbsolute(io, decoy, .{}); + f.close(io); + } + names = try posted(io, env, &stage); + var seen: usize = 0; + var has_alpha = false; + var has_beta = false; + var has_decoy = false; + while (names.next()) |name| { + seen += 1; + has_alpha = has_alpha or std.mem.eql(u8, name, "alpha"); + has_beta = has_beta or std.mem.eql(u8, name, "beta"); + has_decoy = has_decoy or std.mem.eql(u8, name, "decoy"); + } + try testing.expectEqual(@as(usize, 3), seen); + try testing.expect(has_alpha and has_beta and has_decoy); + + unpost(io, a.path, a.inode); + a.server.deinit(io); + unpost(io, b.path, b.inode); + b.server.deinit(io); + try Io.Dir.deleteFileAbsolute(io, decoy); +} + +test "post watch: names appear and disappear" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(env, &dbuf); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, reg, .fromMode(0o750)); + + var w = try Watch.init(); + defer w.deinit(); + try w.add(reg); + try testing.expect((try w.next()) == null); + + var pbuf: [sun_path_len]u8 = undefined; + var a = try post(io, env, "alpha", 4, &pbuf); + const seen = try w.next(); + try testing.expect(seen != null); + try testing.expectEqualStrings("alpha", seen.?.name); + try testing.expectEqual(Watch.Kind.added, seen.?.kind); + unpost(io, a.path, a.inode); + a.server.deinit(io); + const gone = try w.next(); + try testing.expect(gone != null); + try testing.expectEqualStrings("alpha", gone.?.name); + try testing.expectEqual(Watch.Kind.removed, gone.?.kind); + try testing.expect((try w.next()) == null); +} + +test "post unpost is ownership-checked: a re-posted name survives a late unpost" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var pbuf: [sun_path_len]u8 = undefined; + + // A posts; the socket file is lost behind its back; B re-posts. + var a = try post(io, env, "svc", 4, &pbuf); + try Io.Dir.deleteFileAbsolute(io, a.path); + var b = try post(io, env, "svc", 4, &pbuf); + try testing.expect(probe(b.path) == .live); + + // A's late unpost (a delayed stop, say) must leave B's entry alone. + unpost(io, a.path, a.inode); + const st = try Io.Dir.statFile(.cwd(), io, b.path, .{}); + try testing.expect(st.kind == .unix_domain_socket); + try testing.expect(probe(b.path) == .live); + + // B's own unpost still works — and is idempotent. + unpost(io, b.path, b.inode); + unpost(io, b.path, b.inode); + try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, b.path, .{})); + a.server.deinit(io); + b.server.deinit(io); +} + +test "post names: a 255-byte entry stages and iterates (NAME_MAX)" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(env, &dbuf); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, reg, .fromMode(0o750)); + // A 255-byte name can never be a socket (108-byte budget), but an + // attacker can drop one in the registry; the raw listing must carry + // it without overflowing the staged `len:u8` record. + const fat = "z" ** 255; + var fat_buf: [std.fs.max_path_bytes]u8 = undefined; + const fat_path = try std.fmt.bufPrintZ(&fat_buf, "{s}/{s}", .{ reg, fat }); + { + var f = try Io.Dir.createFileAbsolute(io, fat_path, .{}); + f.close(io); + } + var stage: [1024]u8 = undefined; + var names = try posted(io, env, &stage); + var seen: usize = 0; + var have_fat = false; + while (names.next()) |name| { + seen += 1; + have_fat = have_fat or std.mem.eql(u8, name, fat); + } + try testing.expectEqual(@as(usize, 1), seen); + try testing.expect(have_fat); + try Io.Dir.deleteFileAbsolute(io, fat_path); +} + +test "post long name: max_name_len posts, probes, dials and unposts" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + // The budget assumes a conforming (short) XDG_RUNTIME_DIR; the + // testing tmpdir's path is longer, so the boundary needs a short + // scratch directory of its own. + const io = testing.io; + var xdg_buf: [64]u8 = undefined; + const xdg = try std.fmt.bufPrintZ(&xdg_buf, "/tmp/.p9t{d}", .{linux.getpid()}); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, xdg, .fromMode(0o700)); + var env_buf: [96]u8 = undefined; + const value = try std.fmt.bufPrintZ(&env_buf, "XDG_RUNTIME_DIR={s}", .{xdg}); + const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; + const envp: Env = @ptrCast(&env); + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(envp, &dbuf); + defer { + _ = linux.rmdir(reg.ptr); + _ = linux.rmdir(xdg.ptr); + } + var pbuf: [sun_path_len]u8 = undefined; + const name = "l" ** max_name_len; + try testing.expect(legalName(name)); + var p = try post(io, envp, name, 4, &pbuf); + try testing.expect(probe(p.path) == .live); + var stream = try dial(io, envp, name); + stream.close(io); + unpost(io, p.path, p.inode); + p.server.deinit(io); + // One past the cap is not even a path. + const too_long = "l" ** (max_name_len + 1); + try testing.expect(!legalName(too_long)); + try testing.expectError(error.IllegalName, registryPath(envp, too_long, &pbuf)); +} + +test "post watch: queue overflow surfaces, a deleted registry is gone" { + if (builtin.os.tag != .linux) return error.SkipZigTest; + var s: Scratch = .{ .dir = undefined }; + try s.start(); + defer s.end(); + const env = s.envp(); + const io = testing.io; + var dbuf: [std.fs.max_path_bytes]u8 = undefined; + const reg = try registryDir(env, &dbuf); + _ = try Io.Dir.createDirPathStatus(.cwd(), io, reg, .fromMode(0o750)); + var w = try Watch.init(); + defer w.deinit(); + try w.add(reg); + + // Overwhelm the kernel queue (default 16384 events): the overflow + // must be surfaced, not silently skipped, or a caching consumer + // stays stale forever. + var name_buf: [std.fs.max_path_bytes]u8 = undefined; + const spam = 20000; + for (0..spam) |i| { + const p = try std.fmt.bufPrintZ(&name_buf, "{s}/q{d}", .{ reg, i }); + var f = try Io.Dir.createFileAbsolute(io, p, .{}); + f.close(io); + } + var overflow = false; + var events: usize = 0; + while (try w.next()) |ev| { + events += 1; + if (ev.kind == .overflow) overflow = true; + } + try testing.expect(overflow); + + // The registry directory itself is replaced: the watch must say so. + for (0..spam) |i| { + const p = try std.fmt.bufPrint(&name_buf, "{s}/q{d}", .{ reg, i }); + try Io.Dir.deleteFileAbsolute(io, p); + } + while (try w.next()) |_| {} + _ = linux.rmdir(reg.ptr); // empty now + try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, reg, .{})); + const gone = try w.next(); + try testing.expect(gone != null); + try testing.expectEqual(Watch.Kind.gone, gone.?.kind); + try testing.expect((try w.next()) == null); + // Re-added, the (recreated) registry is watched again. + _ = try Io.Dir.createDirPathStatus(.cwd(), io, reg, .fromMode(0o750)); + try w.add(reg); + { + var f = try Io.Dir.createFileAbsolute(io, try std.fmt.bufPrintZ(&name_buf, "{s}/back", .{reg}), .{}); + f.close(io); + } + const back = try w.next(); + try testing.expect(back != null); + try testing.expectEqualStrings("back", back.?.name); + try Io.Dir.deleteFileAbsolute(io, try std.fmt.bufPrint(&name_buf, "{s}/back", .{reg})); +} diff --git a/src/root.zig b/src/root.zig index 586a41c..fb0bf5c 100644 --- a/src/root.zig +++ b/src/root.zig @@ -50,6 +50,10 @@ test { pub const http = @import("http.zig"); pub const transport = @import("transport.zig"); pub const serve = @import("serve.zig"); +pub const post = @import("post.zig"); +test { + _ = @import("post.zig"); +} pub const Quic = @import("quic.zig").Quic; test { _ = @import("session_test.zig"); diff --git a/src/serve.zig b/src/serve.zig index 255e2ff..19a8ca2 100644 --- a/src/serve.zig +++ b/src/serve.zig @@ -12,6 +12,7 @@ const std = @import("std"); const Io = std.Io; const fs = @import("fs.zig"); const transport = @import("transport.zig"); +const post = @import("post.zig"); /// Comptime bounds of one runner. pub const Limits = struct { @@ -72,6 +73,7 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits }; pub const ListenError = error{TooManyListeners} || Io.net.IpAddress.ListenError || Io.net.UnixAddress.ListenError || Io.net.UnixAddress.InitError || Io.ConcurrentError; + pub const ListenPostedError = post.PostError || error{TooManyListeners} || Io.ConcurrentError; io: Io, root: u64, @@ -81,6 +83,14 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits conns: [limits.connections]Conn, listeners: [limits.listeners]Io.net.Server, nlisteners: usize, + /// The registry socket of a `listenPosted`, zero-terminated; + /// `stop()` unlinks it (unpost on stop) — but only while it is + /// still this runner's entry (`posted_ino`). + posted_path: [transport.sun_path_len + 1]u8 = @splat(0), + posted_len: usize = 0, + /// The bound registry entry's inode, the ownership proof for + /// the unpost in `stop()`. + posted_ino: u64 = 0, /// Accept tasks and connection tasks; `stop()` cancels it. group: Io.Group, /// Guards `Conn.used`. @@ -292,6 +302,7 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits r.seed = o.seed; r.handler = o.handler; r.greet_timeout_ms = o.greet_timeout_ms; + r.posted_len = 0; r.nlisteners = 0; r.group = .init; r.slots = .init; @@ -307,15 +318,46 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits pub fn listen(r: *Self, address: transport.Address, backlog: u31) ListenError!Io.net.IpAddress { if (r.nlisteners == limits.listeners) return error.TooManyListeners; if (r.stopping.load(.acquire)) return error.TooManyListeners; + var server = try transport.listen(r.io, address, backlog); + errdefer server.deinit(r.io); + try r.startListener(server); + return server.socket.address; + } + + /// Posts the runner on the registry socket + /// `$XDG_RUNTIME_DIR/9p/` and starts accepting on it: + /// `post.post` creates the 0o750 registry directory and runs the + /// stale protocol (a refused entry is replaced; a live server + /// owning the name is `AlreadyPosted`; a non-socket entry is + /// never deleted). One posted name per runner: a second + /// `listenPosted` is `AlreadyPosted` (its socket would otherwise + /// be orphaned in the registry — nothing would unpost it). + /// `stop()` unposts — the socket is unlinked when the runner + /// stops, as long as the entry is still the runner's own. + pub fn listenPosted(r: *Self, env: post.Env, name: []const u8, backlog: u31) ListenPostedError!void { + if (r.nlisteners == limits.listeners) return error.TooManyListeners; + if (r.stopping.load(.acquire)) return error.TooManyListeners; + if (r.posted_len != 0) return error.AlreadyPosted; + var p = try post.post(r.io, env, name, backlog, &r.posted_path); + errdefer { + post.unpost(r.io, p.path, p.inode); + p.server.deinit(r.io); + } + try r.startListener(p.server); + r.posted_len = p.path.len; + r.posted_ino = p.inode; + } + + /// Registers a bound listener and starts its accept task. + fn startListener(r: *Self, server: Io.net.Server) Io.ConcurrentError!void { const i = r.nlisteners; - r.listeners[i] = try transport.listen(r.io, address, backlog); + r.listeners[i] = server; errdefer r.listeners[i].deinit(r.io); r.nlisteners += 1; r.group.concurrent(r.io, acceptLoop, .{ r, i }) catch |err| { r.nlisteners -= 1; return err; }; - return r.listeners[i].socket.address; } /// Connections held right now. @@ -333,13 +375,21 @@ pub fn Runner(comptime Backend: type, comptime opts: fs.Options, comptime limits for (&r.conns) |*c| if (c.live()) c.close(); } - /// Stops accepting, hangs every connection up, waits for their tasks - /// and closes the listeners. Idempotent; the runner is spent after. + /// Stops accepting, hangs every connection up, waits for their + /// tasks and closes the listeners. A posted listener is unposted + /// — its registry socket is unlinked, but only while the entry + /// is still the runner's own: a name that was re-posted by + /// another server (this one's socket file having been lost) + /// survives the stop. Idempotent; the runner is spent after. pub fn stop(r: *Self) void { if (r.stopping.swap(true, .acq_rel)) return; r.group.cancel(r.io); for (r.listeners[0..r.nlisteners]) |*l| l.deinit(r.io); r.nlisteners = 0; + if (r.posted_len != 0) { + post.unpost(r.io, r.posted_path[0..r.posted_len :0], r.posted_ino); + r.posted_len = 0; + } } fn acceptLoop(r: *Self, i: usize) void { @@ -846,3 +896,120 @@ test "serve: a backend answered from another thread under lock()" { rig.runner.stop(); try expectHangup(tc.one(.{ .stat = .{ .fid = 0 } })); } + +test "serve: listenPosted serves the registry name and stop() unposts" { + if (@import("builtin").os.tag != .linux) return error.SkipZigTest; + var rig: Rig = .{ .dir = undefined }; + const io = testing.io; + // A scratch registry: XDG_RUNTIME_DIR is the rig's own temp dir. + rig.dir = testing.tmpDir(.{}); + errdefer rig.dir.cleanup(); + var real_buf: [std.fs.max_path_bytes]u8 = undefined; + const len = try rig.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..len]}); + const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; + const envp: post.Env = @ptrCast(&env); + + rig.stub = .{}; + rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &rig.stub, .serve = serveNow } }); + try rig.runner.listenPosted(envp, "posted", 4); + var path_buf: [transport.sun_path_len]u8 = undefined; + const path = try post.registryPath(envp, "posted", &path_buf); + + // The name is posted, live, and serves a full 9P session. + try testing.expect(post.probe(path) == .live); + var tc: TestClient = .{ .io = undefined, .stream = undefined }; + try tc.open(io, .{ .unix = path }); + defer tc.close(); + try tc.handshake(); + try tc.readIndex(1); + + // A second post of the same name is refused while the runner lives. + var pbuf: [transport.sun_path_len]u8 = undefined; + try testing.expectError(error.AlreadyPosted, post.post(io, envp, "posted", 4, &pbuf)); + + // The listing sees it; stop() unposts and the entry disappears. + var stage: [512]u8 = undefined; + var names = try post.posted(io, envp, &stage); + var seen = false; + while (names.next()) |n| seen = seen or std.mem.eql(u8, n, "posted"); + try testing.expect(seen); + rig.runner.stop(); + try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, path, .{})); + names = try post.posted(io, envp, &stage); + try testing.expect(names.next() == null); + rig.dir.cleanup(); +} + +test "serve: a second listenPosted is refused; stop() never unposts another's name" { + if (@import("builtin").os.tag != .linux) return error.SkipZigTest; + var rig: Rig = .{ .dir = undefined }; + const io = testing.io; + rig.dir = testing.tmpDir(.{}); + errdefer rig.dir.cleanup(); + var real_buf: [std.fs.max_path_bytes]u8 = undefined; + const len = try rig.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..len]}); + const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; + const envp: post.Env = @ptrCast(&env); + + rig.stub = .{}; + rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &rig.stub, .serve = serveNow } }); + try rig.runner.listenPosted(envp, "twice", 4); + // One posted name per runner: a second would overwrite the first's + // path and orphan its socket in the registry. + try testing.expectError(error.AlreadyPosted, rig.runner.listenPosted(envp, "twice", 4)); + try testing.expectError(error.AlreadyPosted, rig.runner.listenPosted(envp, "other", 4)); + + var path_buf: [transport.sun_path_len]u8 = undefined; + const path = try post.registryPath(envp, "twice", &path_buf); + // The runner's socket file is lost behind its back (rm, crash + // cleanup), and another server takes the now-free name. + try Io.Dir.deleteFileAbsolute(io, path); + var thief = try post.post(io, envp, "twice", 4, &path_buf); + try testing.expect(post.probe(thief.path) == .live); + // stop() unposts only what it still owns: the thief survives. + rig.runner.stop(); + const st = try Io.Dir.statFile(.cwd(), io, thief.path, .{}); + try testing.expect(st.kind == .unix_domain_socket); + try testing.expect(post.probe(thief.path) == .live); + post.unpost(io, thief.path, thief.inode); + thief.server.deinit(io); + rig.dir.cleanup(); +} + +test "serve: listenPosted beside listen(): stop unposts only the registry name" { + if (@import("builtin").os.tag != .linux) return error.SkipZigTest; + var rig: Rig = .{ .dir = undefined }; + const io = testing.io; + rig.dir = testing.tmpDir(.{}); + errdefer rig.dir.cleanup(); + var real_buf: [std.fs.max_path_bytes]u8 = undefined; + const len = try rig.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..len]}); + const env = [2]?[*:0]const u8{ @ptrCast(value.ptr), null }; + const envp: post.Env = @ptrCast(&env); + + rig.stub = .{}; + rig.runner.init(.{ .io = io, .root = Stub.root_node, .handler = .{ .ctx = &rig.stub, .serve = serveNow } }); + // A plain Unix listener (the application's path policy) beside the + // posted name: both listener slots fill. + var unix_buf: [std.fs.max_path_bytes]u8 = undefined; + const unix_path = try std.fmt.bufPrintZ(&unix_buf, "{s}/plain.sock", .{real_buf[0..len]}); + _ = try rig.runner.listen(.{ .unix = unix_path }, 4); + try rig.runner.listenPosted(envp, "mixed", 4); + + var path_buf: [transport.sun_path_len]u8 = undefined; + const posted_path = try post.registryPath(envp, "mixed", &path_buf); + try testing.expect(post.probe(posted_path) == .live); + rig.runner.stop(); + // The posted name is unposted; the plain path is the application's. + try testing.expectError(error.FileNotFound, Io.Dir.statFile(.cwd(), io, posted_path, .{})); + const plain_st = try Io.Dir.statFile(.cwd(), io, unix_path, .{}); + try testing.expect(plain_st.kind == .unix_domain_socket); + try Io.Dir.deleteFileAbsolute(io, unix_path); + rig.dir.cleanup(); +} -- cgit v1.3