summaryrefslogtreecommitdiff
path: root/src/post.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/post.zig')
-rw-r--r--src/post.zig1024
1 files changed, 1024 insertions, 0 deletions
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/<name>`. 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/<name>` must stay inside `sun_path_len`, and a
+/// conforming `$XDG_RUNTIME_DIR` (`/run/user/<uid>`) 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/<name>` 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}));
+}