//! The Linux platform layer: one background thread runs a `poll()` loop over //! a listener and every client connection, feeding each connection's core //! `Conn` with `push`/`step`/`output`/`wrote`. No per-connection threads, no //! allocation after `init`; every buffer lives in a caller-placed `Storage`. //! //! Also home of the debug facilities (`debug`, `provider`) and the /runtime //! generators (`runtime`). See docs/LIBRARY.md. //! //! Client admission: a new connection takes a free slot. When every slot is //! taken, the connection that has held a slot without any fid (never //! attached, or fully clunked) for longer than `evict_idle_ms` is dropped in //! its favour; if there is none, the new connection is closed ("refused"). //! Nothing that holds a fid is ever evicted. //! //! `sleepServing(ms)` lets a request handler (a ctl command, say) wait //! without stalling the other clients: called on the probe thread from inside //! a request it keeps running the poll loop for every client whose request //! is not in progress until the time is up. Requests served from inside such //! a wait may wait themselves, up to `max_nested_sleeps` deep (each level is a //! different client, so the depth is bounded by the client table anyway); the //! level past that, and any call off the probe thread, is a plain sleep. const std = @import("std"); const builtin = @import("builtin"); const linux = std.os.linux; const cloud9 = @import("cloud9"); const core = @import("../core.zig"); pub const debug = @import("debug.zig"); pub const provider = @import("provider.zig"); pub const runtime = @import("runtime.zig"); pub const DebugProvider = provider.DebugProvider; pub const Listen = union(enum) { /// A unix socket path (< 108 bytes); a stale socket file is unlinked first. unix: []const u8, /// An IPv4 literal "a.b.c.d:port". tcp: []const u8, /// An already listening socket, owned by the caller. fd: i32, /// One pre-connected client on these descriptors (stdio: 0 and 1). Nothing /// is accepted and the loop ends when the client hangs up. client: struct { in: i32, out: i32 }, }; pub const Options = struct { /// For `std.debug` symbolization. io: std.Io, listen: Listen, /// Largest msize offered to clients (clamped to the server's `cfg.msize`). msize: u32 = 64 * 1024, /// Hold a panicking thread until /panic/ctl says "continue". hold_on_panic: bool = true, /// Real-time signal used to snapshot other threads. capture_signal: u8 = debug.default_capture_signal, /// Install the SIGTRAP handler so `@breakpoint()` parks the thread. breakpoints: bool = true, /// Mount /threads, /addr, /mem, /hex, /breakpoints, /panic (six provider slots). mount_debug: bool = true, }; pub const Error = error{ /// Another `Debug` (another probe) exists in this process. AlreadyInitialized, /// No register capture on this architecture. Unsupported, /// `Shared` has fewer than six free provider slots. TooManyProviders, PathTooLong, BadAddress, /// The unix path holds a live server; nothing is deleted or taken. AlreadyListening, /// The unix path holds a foreign non-socket entry; never deleted. Occupied, /// A syscall failed; `last_errno` says which error. Syscall, }; /// Idle time without fids after which a slot holder may be evicted. pub const evict_idle_ms: i64 = 500; /// Largest number of connections accepted per poll wakeup. const accept_burst = 64; /// How deep `sleepServing` may nest (each level keeps a poll round on the stack). pub const max_nested_sleeps = 8; /// Static per-client storage: `max_clients` core `Storage`s and `Conn`s, the /// poll table, the debug text arena and the debug provider's snapshot pool. pub fn Storage(comptime max_clients: u8, comptime Srv: type) type { return Probe(Srv).Storage(max_clients); } pub fn Probe(comptime Srv: type) type { return struct { const Self = @This(); pub fn Storage(comptime max_clients: u8) type { comptime std.debug.assert(max_clients > 0); return struct { pub const capacity = max_clients; conns: [max_clients]Srv.Storage, clients: [max_clients]Client, /// [0] wake eventfd, [1] listener, [2..] one per client slot. pollfds: [max_clients + 2]linux.pollfd, text_buf: [16 * 1024]u8, dp: DebugProvider, }; } pub const Client = struct { conn: Srv.Conn, in: i32 = -1, out: i32 = -1, used: bool = false, /// Descriptors we opened (accepted) are closed on drop; borrowed ones are not. owned: bool = false, /// Send with MSG_NOSIGNAL; falls back to write(2) on ENOTSOCK. is_socket: bool = true, /// Monotonic ms of the last byte received. last_active: i64 = 0, }; shared: *Srv.Shared, clients: []Client, conns: []Srv.Storage, pollfds: []linux.pollfd, dbg: debug.Debug, dp: *DebugProvider, msize: u32, listen_fd: i32 = -1, own_listener: bool = false, is_tcp: bool = false, single: bool = false, wake_fd: i32 = -1, unix_path: [108]u8 = undefined, unix_len: usize = 0, /// The bound entry's inode; `stop` unlinks only its own socket. unix_ino: u64 = 0, thread: ?std.Thread = null, thread_tid: std.atomic.Value(u32) = .init(0), nclients: std.atomic.Value(u32) = .init(0), /// Connections closed because no slot was free. refused: u64 = 0, stopping: std.atomic.Value(bool) = .init(false), /// Slots whose request is being handled (excluded from nested servicing and eviction). serving: std.StaticBitSet(256) = .initEmpty(), /// Current `sleepServing` nesting depth. nested: u8 = 0, debug_ready: bool = false, last_errno: linux.E = .SUCCESS, // -- lifecycle ------------------------------------------------------- /// Installs the debug facilities, mounts the debug providers into /// `shared` and opens the listener. `storage` is a `*Storage(n)`. On /// failure `shared` may already hold the debug providers and must be /// discarded. pub fn init(p: *Self, shared: *Srv.Shared, storage: anytype, opts: Options) Error!void { p.* = .{ .shared = shared, .clients = &storage.clients, .conns = &storage.conns, .pollfds = &storage.pollfds, .dbg = undefined, .dp = &storage.dp, .msize = opts.msize, }; for (p.clients) |*c| c.used = false; shared.hash_seed = randomSeed(); debug.hold_on_panic = opts.hold_on_panic; p.dbg.init(.{ .io = opts.io, .text_buf = &storage.text_buf, .capture_signal = opts.capture_signal }) catch |e| return switch (e) { error.AlreadyInitialized => error.AlreadyInitialized, error.Unsupported => error.Unsupported, else => error.Syscall, }; p.debug_ready = true; errdefer { p.dbg.deinit(); p.debug_ready = false; } if (opts.breakpoints) p.dbg.enableBreakpoints() catch |e| switch (e) { // No breakpoint support on this architecture: everything else still works. error.Unsupported => {}, else => return error.Syscall, }; if (opts.mount_debug) { p.dp.init(&p.dbg); p.dp.mountAll(shared) catch return error.TooManyProviders; } const efd = linux.eventfd(0, linux.EFD.CLOEXEC | linux.EFD.NONBLOCK); try p.check(efd); p.wake_fd = @intCast(efd); errdefer { _ = linux.close(p.wake_fd); p.wake_fd = -1; } switch (opts.listen) { .unix => |path| try p.listenUnix(path), .tcp => |text| try p.listenTcp(text), .fd => |fd| { try p.setNonblock(fd); p.listen_fd = fd; }, .client => |c| { p.single = true; try p.setNonblock(c.in); if (c.out != c.in) try p.setNonblock(c.out); _ = p.addClient(c.in, c.out, false); }, } } /// Spawns the poll thread. pub fn start(p: *Self) std.Thread.SpawnError!void { std.debug.assert(p.thread == null); p.stopping.store(false, .release); p.thread = try std.Thread.spawn(.{}, run, .{p}); } /// Waits for the poll thread to end (only happens by itself in /// `.client` mode, when the client hangs up). pub fn wait(p: *Self) void { if (p.thread) |t| { t.join(); p.thread = null; } } /// Stops the poll thread, drops every client, closes what `init` /// opened and restores the signal dispositions. pub fn stop(p: *Self) void { // Joining the poll thread from itself would hang forever; a // request handler that wants the server gone uses `requestStop`. std.debug.assert(p.thread_tid.load(.acquire) != @as(u32, @intCast(linux.gettid()))); p.stopping.store(true, .release); p.wakeLoop(); p.wait(); for (p.clients, 0..) |*c, i| if (c.used) p.dropClient(i); if (p.listen_fd >= 0) { if (p.own_listener) _ = linux.close(p.listen_fd); p.listen_fd = -1; } if (p.unix_len > 0) { // Only our own entry: a late stop must never unlink a // name another server has since claimed. if (unixIno(@ptrCast(&p.unix_path))) |ino| { if (p.unix_ino != 0 and ino == p.unix_ino) _ = linux.unlink(@ptrCast(&p.unix_path)); } p.unix_len = 0; } if (p.wake_fd >= 0) { _ = linux.close(p.wake_fd); p.wake_fd = -1; } if (p.debug_ready) { p.dbg.deinit(); p.debug_ready = false; } } /// Live client count (for /runtime/clients). pub fn clientCount(p: *const Self) u32 { return p.nclients.load(.acquire); } pub fn clientCounter(p: *const Self) *const std.atomic.Value(u32) { return &p.nclients; } /// Waits `ms` while keeping the other clients served (see the file comment). pub fn sleepServing(p: *Self, ms: u64) void { const on_thread = p.thread_tid.load(.acquire) == @as(u32, @intCast(linux.gettid())); if (!on_thread or p.serving.count() == 0 or p.nested >= max_nested_sleeps) return sleepMs(ms); p.nested += 1; defer p.nested -= 1; const deadline = monotonicMs() + @as(i64, @intCast(@min(ms, std.math.maxInt(i32)))); while (!p.stopping.load(.acquire)) { const now = monotonicMs(); if (now >= deadline) break; p.pollOnce(@intCast(deadline - now)); } } /// Asks the poll thread to stop; safe to call from a signal handler /// (an atomic store and one write to the wake eventfd). `stop` (or /// `wait`) still has to run afterwards to release everything. pub fn requestStop(p: *Self) void { p.stopping.store(true, .release); p.wakeLoop(); } // -- the loop -------------------------------------------------------- fn run(p: *Self) void { const tid: u32 = @intCast(linux.gettid()); p.thread_tid.store(tid, .release); // A breakpoint or panic on this thread must never park it (see debug.zig). debug.server_tid.store(tid, .release); setThreadName("9proc"); while (!p.stopping.load(.acquire)) { if (p.single and p.clientCount() == 0) break; p.pollOnce(-1); } debug.server_tid.store(0, .release); p.thread_tid.store(0, .release); } fn wakeLoop(p: *Self) void { if (p.wake_fd < 0) return; const one: u64 = 1; _ = linux.write(p.wake_fd, @ptrCast(&one), 8); } /// One `poll()` round: accept, read, step, write. Slots whose request /// is in progress (`serving`, only inside `sleepServing`) are left untouched. fn pollOnce(p: *Self, timeout_ms: i32) void { p.pollfds[0] = .{ .fd = p.wake_fd, .events = linux.POLL.IN, .revents = 0 }; p.pollfds[1] = .{ .fd = p.listen_fd, .events = linux.POLL.IN, .revents = 0 }; for (p.clients, 0..) |*c, i| { var fd: i32 = -1; var events: i16 = 0; if (c.used and !p.serving.isSet(i)) { fd = c.in; if (c.conn.output().len > 0) { fd = c.out; events = linux.POLL.OUT; } else if (inputRoom(&c.conn) > 0) { events = linux.POLL.IN; } } p.pollfds[2 + i] = .{ .fd = fd, .events = events, .revents = 0 }; } const rc = linux.poll(p.pollfds.ptr, p.pollfds.len, timeout_ms); switch (linux.errno(rc)) { .SUCCESS => {}, .INTR => return, else => { sleepMs(10); return; }, } if (p.pollfds[0].revents != 0) { var v: u64 = 0; _ = linux.read(p.wake_fd, @ptrCast(&v), 8); } if (p.stopping.load(.acquire)) return; if (p.pollfds[1].revents != 0) p.acceptSome(); for (p.clients, 0..) |*c, i| { const re = p.pollfds[2 + i].revents; if (re == 0 or !c.used or p.serving.isSet(i)) continue; if (re & (linux.POLL.IN | linux.POLL.HUP | linux.POLL.ERR | linux.POLL.NVAL) != 0) { p.readClient(i, re & linux.POLL.HUP != 0); } else if (re & linux.POLL.OUT != 0) { p.service(i); } if (p.stopping.load(.acquire)) return; } } fn acceptSome(p: *Self) void { var n: usize = 0; while (n < accept_burst) : (n += 1) { const rc = linux.accept4(p.listen_fd, null, null, linux.SOCK.NONBLOCK | linux.SOCK.CLOEXEC); switch (linux.errno(rc)) { .SUCCESS => {}, .AGAIN => return, .INTR, .CONNABORTED => continue, // Descriptor/memory exhaustion is transient (clients hang up); // back off instead of spinning on a readable listener. .MFILE, .NFILE, .NOBUFS, .NOMEM, .PERM => { sleepMs(100); return; }, else => return, } const cfd: i32 = @intCast(rc); if (p.is_tcp) { const one: u32 = 1; _ = linux.setsockopt(cfd, linux.IPPROTO.TCP, linux.TCP.NODELAY, @ptrCast(&one), @sizeOf(u32)); } if (p.addClient(cfd, cfd, true) != null) continue; if (p.evictable()) |victim| { p.dropClient(victim); _ = p.addClient(cfd, cfd, true); continue; } p.refused += 1; // A flood must not flood stderr. if (p.refused == 1 or p.refused % 1000 == 0) std.debug.print("9proc: refused connection ({d} clients open, {d} refused so far)\n", .{ p.clients.len, p.refused }); _ = linux.close(cfd); } } /// The longest-idle slot holder without fids, if idle long enough. fn evictable(p: *Self) ?usize { const now = monotonicMs(); var best: ?usize = null; for (p.clients, 0..) |*c, i| { if (!c.used or !c.owned) continue; if (p.serving.isSet(i)) continue; if (c.conn.fidCount() != 0) continue; if (now - c.last_active < evict_idle_ms) continue; if (best == null or c.last_active < p.clients[best.?].last_active) best = i; } return best; } fn addClient(p: *Self, in: i32, out: i32, owned: bool) ?usize { for (p.clients, 0..) |*c, i| { if (c.used) continue; c.conn = Srv.Conn.init(p.shared, &p.conns[i], p.msize); c.in = in; c.out = out; c.used = true; c.owned = owned; c.is_socket = true; c.last_active = monotonicMs(); _ = p.nclients.fetchAdd(1, .acq_rel); return i; } return null; } fn dropClient(p: *Self, i: usize) void { const c = &p.clients[i]; if (!c.used) return; c.conn.hangup(); if (c.owned) { _ = linux.close(c.in); if (c.out != c.in) _ = linux.close(c.out); } c.used = false; c.in = -1; c.out = -1; _ = p.nclients.fetchSub(1, .acq_rel); } /// Free space in the connection's input buffer (cloud9 keeps one /// msize-sized frame; `push` copies at most this much). fn inputRoom(conn: *const Srv.Conn) usize { return conn.inputRoom(); } fn readClient(p: *Self, i: usize, hup: bool) void { const c = &p.clients[i]; var buf: [64 * 1024]u8 = undefined; const room = inputRoom(&c.conn); if (room == 0) return p.service(i); const want = @min(room, buf.len); while (true) { const rc = linux.read(c.in, &buf, want); switch (linux.errno(rc)) { .SUCCESS => { if (rc == 0) return p.dropClient(i); const taken = c.conn.push(buf[0..rc]); std.debug.assert(taken == rc); c.last_active = monotonicMs(); return p.service(i); }, .INTR => continue, .AGAIN => { if (hup) p.dropClient(i); return; }, else => return p.dropClient(i), } } } /// Runs requests and drains output until nothing moves. fn service(p: *Self, i: usize) void { const c = &p.clients[i]; std.debug.assert(!p.serving.isSet(i)); p.serving.set(i); defer p.serving.unset(i); while (c.used) { var moved = false; while (true) { const more = c.conn.step() catch return p.dropClient(i); if (!more) break; moved = true; } const before = c.conn.output().len; p.flush(c) catch return p.dropClient(i); if (c.conn.output().len != before) moved = true; if (!moved) return; } } fn flush(p: *Self, c: *Client) error{Closed}!void { _ = p; while (c.conn.output().len > 0) { const chunk = c.conn.output(); const rc = if (c.is_socket) linux.sendto(c.out, chunk.ptr, chunk.len, linux.MSG.NOSIGNAL, null, 0) else linux.write(c.out, chunk.ptr, chunk.len); switch (linux.errno(rc)) { .SUCCESS => { if (rc == 0) return; c.conn.wrote(rc); }, .INTR => continue, .AGAIN => return, .NOTSOCK => c.is_socket = false, else => return error.Closed, } } } // -- listeners ------------------------------------------------------- fn check(p: *Self, rc: usize) Error!void { const e = linux.errno(rc); if (e != .SUCCESS) { p.last_errno = e; return error.Syscall; } } fn setNonblock(p: *Self, fd: i32) Error!void { const rc = linux.fcntl(fd, linux.F.GETFL, 0); try p.check(rc); const nonblock: u32 = @bitCast(linux.O{ .NONBLOCK = true }); try p.check(linux.fcntl(fd, linux.F.SETFL, rc | nonblock)); } fn listenUnix(p: *Self, path: []const u8) Error!void { var sa: linux.sockaddr.un = .{ .path = @splat(0) }; if (path.len == 0 or path.len >= sa.path.len) return error.PathTooLong; @memcpy(sa.path[0..path.len], path); const rc = linux.socket(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK, 0); try p.check(rc); const lfd: i32 = @intCast(rc); errdefer _ = linux.close(lfd); // Nothing foreign is deleted: a non-socket entry at the path // is refused (Occupied), a live server is refused // (AlreadyListening), and only a socket that refuses a // connect — a corpse — is unlinked. (A hand racing the swap // between probe and unlink is the documented residual of this // cheap protocol; cloud9.post claims names atomically when // that matters.) const st = unixStat(@ptrCast(&sa.path)) catch return error.Occupied; if (st) |s| { if (s.mode & linux.S.IFMT != linux.S.IFSOCK) return error.Occupied; if (cloud9.post.probe(@ptrCast(&sa.path)) != .stale) return error.AlreadyListening; _ = linux.unlink(@ptrCast(&sa.path)); } try p.check(linux.bind(lfd, @ptrCast(&sa), @sizeOf(linux.sockaddr.un))); try p.check(linux.listen(lfd, 128)); p.listen_fd = lfd; p.own_listener = true; p.unix_path = sa.path; p.unix_len = path.len; // Our entry's inode, so `stop` never unlinks a name another // server has since claimed. p.unix_ino = if (unixStat(@ptrCast(&p.unix_path)) catch null) |s| s.ino else 0; } /// statx(2) of one path: null when it does not exist, and an /// error when the kernel cannot say — the caller refuses rather /// than guesses. fn unixStat(path: [*:0]const u8) !?linux.Statx { var stx: linux.Statx = undefined; const rc = linux.statx(linux.AT.FDCWD, path, 0, .{ .TYPE = true, .INO = true }, &stx); const s: isize = @bitCast(rc); if (s == 0) return stx; const noent: isize = @intCast(@intFromEnum(linux.E.NOENT)); if (s == -noent) return null; return error.StatFailed; } /// The socket at `path`, by inode; null when it is gone or /// unreadable. fn unixIno(path: [*:0]const u8) ?u64 { const stx = unixStat(path) catch return null; return if (stx) |s| s.ino else null; } fn listenTcp(p: *Self, text: []const u8) Error!void { const sa = parseIpv4(text) orelse return error.BadAddress; const rc = linux.socket(linux.AF.INET, linux.SOCK.STREAM | linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK, 0); try p.check(rc); const lfd: i32 = @intCast(rc); errdefer _ = linux.close(lfd); const one: u32 = 1; _ = linux.setsockopt(lfd, linux.SOL.SOCKET, linux.SO.REUSEADDR, @ptrCast(&one), @sizeOf(u32)); try p.check(linux.bind(lfd, @ptrCast(&sa), @sizeOf(linux.sockaddr.in))); try p.check(linux.listen(lfd, 128)); p.listen_fd = lfd; p.own_listener = true; p.is_tcp = true; } }; } /// Entropy for the core's fid hash (so fid numbers cannot be chosen to /// collide); falls back to the clock if getrandom fails. fn randomSeed() u32 { var b: [4]u8 = undefined; if (linux.errno(linux.getrandom(&b, b.len, 0)) == .SUCCESS) return std.mem.readInt(u32, &b, .little); var ts: linux.timespec = undefined; _ = linux.clock_gettime(.MONOTONIC, &ts); return @truncate(@as(u64, @bitCast(ts.nsec)) ^ (@as(u64, @bitCast(ts.sec)) << 20)); } /// "a.b.c.d:port" as a socket address, or null. pub fn parseIpv4(text: []const u8) ?linux.sockaddr.in { const colon = std.mem.lastIndexOfScalar(u8, text, ':') orelse return null; const port = std.fmt.parseInt(u16, text[colon + 1 ..], 10) catch return null; var octets: [4]u8 = undefined; var it = std.mem.splitScalar(u8, text[0..colon], '.'); for (&octets) |*o| o.* = std.fmt.parseInt(u8, it.next() orelse return null, 10) catch return null; if (it.next() != null) return null; return .{ .port = std.mem.nativeToBig(u16, port), .addr = @bitCast(octets) }; } /// Names the calling thread (comm, at most 15 bytes) via prctl. pub fn setThreadName(name: []const u8) void { var buf: [16]u8 = @splat(0); const n = @min(name.len, 15); @memcpy(buf[0..n], name[0..n]); _ = linux.prctl(@intFromEnum(linux.PR.SET_NAME), @intFromPtr(&buf), 0, 0, 0); } pub fn sleepMs(ms: u64) void { var req: linux.timespec = .{ .sec = @intCast(ms / 1000), .nsec = @intCast((ms % 1000) * 1_000_000) }; var rem: linux.timespec = undefined; while (linux.errno(linux.nanosleep(&req, &rem)) == .INTR) req = rem; } pub fn monotonicMs() i64 { var ts: linux.timespec = undefined; _ = linux.clock_gettime(.MONOTONIC, &ts); return ts.sec * 1000 + @divTrunc(ts.nsec, 1_000_000); } // --------------------------------------------------------------------------- // Tests: a real unix socket, a cloud9.Client on the other end. // --------------------------------------------------------------------------- const testing = std.testing; test { _ = debug; _ = provider; _ = runtime; } const TestBuild = struct { pub const zig_version: []const u8 = builtin.zig_version_string; pub const target: []const u8 = "test"; pub const optimize: []const u8 = "Debug"; pub const time: []const u8 = "2024-01-01T00:00:00Z"; pub const change: []const u8 = "none"; }; const test_cfg: core.Config = .{ .name = "probetest", .build = TestBuild, .msize = 8192, .max_fids = 16, .max_providers = 6, .snapshot_slots = 2, .snapshot_bytes = 1024, }; const TS = core.Server(test_cfg); const TP = Probe(TS); /// A blocking client over a connected socket. const SockClient = struct { fd: i32, client: cloud9.Client, cin: [8192]u8 = undefined, cout: [8192]u8 = undefined, fn connect(sc: *SockClient, path: []const u8) !void { var sa: linux.sockaddr.un = .{ .path = @splat(0) }; @memcpy(sa.path[0..path.len], path); const rc = linux.socket(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0); if (linux.errno(rc) != .SUCCESS) return error.Socket; sc.fd = @intCast(rc); if (linux.errno(linux.connect(sc.fd, &sa, @sizeOf(linux.sockaddr.un))) != .SUCCESS) return error.Connect; sc.client = .init(.{ .in = &sc.cin, .out = &sc.cout }); } fn close(sc: *SockClient) void { _ = linux.close(sc.fd); } /// One round trip; null when the server closed the connection. fn rpc(sc: *SockClient, req: cloud9.Client.Request) !?cloud9.Client.Result { _ = try sc.client.submit(req); while (sc.client.output().len > 0) { const out = sc.client.output(); const rc = linux.write(sc.fd, out.ptr, out.len); switch (linux.errno(rc)) { .SUCCESS => sc.client.wrote(rc), .PIPE, .CONNRESET => return null, else => return error.Write, } } var buf: [8192]u8 = undefined; while (true) { if (sc.client.take()) |done| return done.result; const rc = linux.read(sc.fd, &buf, buf.len); switch (linux.errno(rc)) { .SUCCESS => {}, .CONNRESET => return null, else => return error.Read, } if (rc == 0) return null; var rest: []const u8 = buf[0..rc]; while (rest.len > 0) rest = rest[sc.client.push(rest)..]; } } fn session(sc: *SockClient) !void { const v = (try sc.rpc(.{ .version = .{ .msize = 8192 } })) orelse return error.Closed; try testing.expectEqual(@as(u32, 8192), v.version.msize); const a = (try sc.rpc(.{ .attach = .{ .fid = 0, .uname = "t" } })) orelse return error.Closed; try testing.expect(a == .attach); } fn readFile(sc: *SockClient, names: []const []const u8, out: []u8) ![]u8 { const w = (try sc.rpc(.{ .walk = .{ .fid = 0, .newfid = 1, .names = names } })) orelse return error.Closed; try testing.expectEqual(@as(u16, @intCast(names.len)), w.walk.nwqid); _ = (try sc.rpc(.{ .open = .{ .fid = 1, .mode = cloud9.oread } })) orelse return error.Closed; const r = (try sc.rpc(.{ .read = .{ .fid = 1, .offset = 0, .count = @intCast(out.len) } })) orelse return error.Closed; const n = r.read.len; @memcpy(out[0..n], r.read); _ = (try sc.rpc(.{ .clunk = .{ .fid = 1 } })) orelse return error.Closed; return out[0..n]; } }; fn testSockPath(buf: []u8, tag: []const u8) ![]const u8 { return std.fmt.bufPrint(buf, "/tmp/9proc-probe-{d}-{s}.sock", .{ linux.getpid(), tag }); } const TestCtx = struct { info: runtime.Info }; test "probe: start, serve a client over a unix socket, stop" { var ctx: TestCtx = .{ .info = .now() }; var shared: TS.Shared = .init(&ctx); const storage = try testing.allocator.create(TP.Storage(2)); defer testing.allocator.destroy(storage); var probe: TP = undefined; var path_buf: [64]u8 = undefined; const path = try testSockPath(&path_buf, "basic"); try probe.init(&shared, storage, .{ .io = testing.io, .listen = .{ .unix = path } }); defer probe.stop(); try probe.start(); ctx.info.clients = probe.clientCounter(); var sc: SockClient = undefined; try sc.connect(path); defer sc.close(); try sc.session(); var buf: [1024]u8 = undefined; const zv = try sc.readFile(&.{ "build", "zig_version" }, &buf); try testing.expectEqualStrings(builtin.zig_version_string, zv); try testing.expectEqual(@as(u32, 1), probe.clientCount()); // The debug providers are mounted: /threads lists the probe thread by name. const names = (try sc.rpc(.{ .walk = .{ .fid = 0, .newfid = 2, .names = &.{"threads"} } })) orelse return error.Closed; try testing.expectEqual(@as(u16, 1), names.walk.nwqid); _ = (try sc.rpc(.{ .open = .{ .fid = 2, .mode = cloud9.oread } })) orelse return error.Closed; const dir = (try sc.rpc(.{ .read = .{ .fid = 2, .offset = 0, .count = 4096 } })) orelse return error.Closed; try testing.expect(dir.read.len > 0); _ = (try sc.rpc(.{ .clunk = .{ .fid = 2 } })) orelse return error.Closed; // A missing file is the Plan 9 error string. const bad = (try sc.rpc(.{ .walk = .{ .fid = 0, .newfid = 3, .names = &.{"nope"} } })) orelse return error.Closed; try testing.expect(bad == .fail); try testing.expectEqualStrings("file does not exist", bad.fail); probe.stop(); // stop() is idempotent and the socket file is gone. probe.stop(); var gone: SockClient = undefined; try testing.expectError(error.Connect, gone.connect(path)); // The debug facilities can be set up again after stop. var probe2: TP = undefined; var shared2: TS.Shared = .init(&ctx); try probe2.init(&shared2, storage, .{ .io = testing.io, .listen = .{ .unix = path } }); probe2.stop(); } test "probe: max_clients refusal and idle eviction" { var ctx: TestCtx = .{ .info = .now() }; var shared: TS.Shared = .init(&ctx); const storage = try testing.allocator.create(TP.Storage(2)); defer testing.allocator.destroy(storage); var probe: TP = undefined; var path_buf: [64]u8 = undefined; const path = try testSockPath(&path_buf, "limit"); try probe.init(&shared, storage, .{ .io = testing.io, .listen = .{ .unix = path } }); defer probe.stop(); try probe.start(); // Two attached clients fill the table; a third is accepted then closed. var a: SockClient = undefined; try a.connect(path); defer a.close(); try a.session(); var b: SockClient = undefined; try b.connect(path); defer b.close(); try b.session(); var c: SockClient = undefined; try c.connect(path); defer c.close(); try testing.expectEqual(@as(?cloud9.Client.Result, null), try c.rpc(.{ .version = .{ .msize = 8192 } })); try testing.expectEqual(@as(u64, 1), probe.refused); // Attached clients are never evicted, even when idle for long. sleepMs(evict_idle_ms + 100); var d: SockClient = undefined; try d.connect(path); defer d.close(); try testing.expectEqual(@as(?cloud9.Client.Result, null), try d.rpc(.{ .version = .{ .msize = 8192 } })); var buf: [256]u8 = undefined; _ = try a.readFile(&.{"README"}, &buf); // A client without fids that has been idle long enough gives way. _ = (try b.rpc(.{ .clunk = .{ .fid = 0 } })) orelse return error.Closed; sleepMs(evict_idle_ms + 100); var e: SockClient = undefined; try e.connect(path); defer e.close(); try e.session(); try testing.expectEqual(@as(?cloud9.Client.Result, null), try b.rpc(.{ .version = .{ .msize = 8192 } })); try testing.expectEqual(@as(u32, 2), probe.clientCount()); } test "probe: single pre-connected client mode ends when the client hangs up" { var ctx: TestCtx = .{ .info = .now() }; var shared: TS.Shared = .init(&ctx); const storage = try testing.allocator.create(TP.Storage(1)); defer testing.allocator.destroy(storage); var sv: [2]i32 = undefined; try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.socketpair(linux.AF.UNIX, linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0, &sv))); var probe: TP = undefined; try probe.init(&shared, storage, .{ .io = testing.io, .listen = .{ .client = .{ .in = sv[1], .out = sv[1] } } }); defer probe.stop(); try probe.start(); var sc: SockClient = .{ .fd = sv[0], .client = undefined }; sc.client = .init(.{ .in = &sc.cin, .out = &sc.cout }); try sc.session(); var buf: [256]u8 = undefined; try testing.expect((try sc.readFile(&.{"README"}, &buf)).len > 0); try testing.expectEqual(@as(u32, 1), probe.clientCount()); _ = linux.close(sv[0]); probe.wait(); try testing.expectEqual(@as(u32, 0), probe.clientCount()); _ = linux.close(sv[1]); } test "parseIpv4 and sleepServing off the probe thread" { const sa = parseIpv4("127.0.0.1:564").?; try testing.expectEqual(std.mem.nativeToBig(u16, 564), sa.port); try testing.expectEqual(@as(u32, @bitCast([4]u8{ 127, 0, 0, 1 })), sa.addr); try testing.expect(parseIpv4("localhost:1") == null); try testing.expect(parseIpv4("1.2.3:1") == null); try testing.expect(parseIpv4("1.2.3.4") == null); try testing.expect(parseIpv4("1.2.3.4:70000") == null); const t0 = monotonicMs(); var probe: TP = undefined; probe.thread_tid = .init(0); probe.serving = .initEmpty(); probe.nested = 0; probe.sleepServing(20); try testing.expect(monotonicMs() - t0 >= 20); }