summaryrefslogtreecommitdiff
path: root/9proc/src/linux/probe.zig
diff options
context:
space:
mode:
Diffstat (limited to '9proc/src/linux/probe.zig')
-rw-r--r--9proc/src/linux/probe.zig829
1 files changed, 829 insertions, 0 deletions
diff --git a/9proc/src/linux/probe.zig b/9proc/src/linux/probe.zig
new file mode 100644
index 0000000..559d981
--- /dev/null
+++ b/9proc/src/linux/probe.zig
@@ -0,0 +1,829 @@
+//! 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,
+ /// 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,
+ 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) {
+ _ = 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.server.in.len - conn.server.in_len;
+ }
+
+ 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);
+ // No libc, so no "is it still listening" probe: unlink a stale socket and bind.
+ _ = 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;
+ }
+
+ 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);
+}