//! 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})); }