summaryrefslogtreecommitdiff
path: root/9proc/src/linux
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-19 23:28:22 -0300
committerGabriel Schneider <[email protected]>2026-09-19 23:28:22 -0300
commitba996acfcad1698adbf4a1834fe50e73b1c6cab9 (patch)
tree282ba00ce5b10d7416aecb9f2f0f0a439340a57d /9proc/src/linux
parentb05abcba3ea09ea106ad28364c6e40a3ec31b890 (diff)
downloadcloud9-ba996acfcad1698adbf4a1834fe50e73b1c6cab9.tar.gz
cloud9-ba996acfcad1698adbf4a1834fe50e73b1c6cab9.zip
Rename programs: 9player -> 9ns, introspect -> 9proc, app -> web (9web)
Directories, binaries, build options (-D9ns, -D9proc), step names, module name (9proc), thread and fs names, env var NINEPLAYER_MOUNT -> NINE_MOUNT, docs and test scripts. Browser assets move to web/static. Co-Authored-By: Claude Fable 5.1 <[email protected]>
Diffstat (limited to '9proc/src/linux')
-rw-r--r--9proc/src/linux/debug.zig1458
-rw-r--r--9proc/src/linux/probe.zig829
-rw-r--r--9proc/src/linux/provider.zig604
-rw-r--r--9proc/src/linux/runtime.zig90
4 files changed, 2981 insertions, 0 deletions
diff --git a/9proc/src/linux/debug.zig b/9proc/src/linux/debug.zig
new file mode 100644
index 0000000..a344380
--- /dev/null
+++ b/9proc/src/linux/debug.zig
@@ -0,0 +1,1458 @@
+//! Linux debug facilities for the 9proc server: threads, stacks,
+//! registers, address → source, memory, breakpoints and panics.
+//!
+//! This file is a pure API; a later adapter turns it into a core `Provider`.
+//! All text is written to a `*std.Io.Writer`. Nothing here allocates after
+//! `init` except from the caller-provided `text_buf`, which is used as a fixed
+//! arena for symbol text and reset before every query.
+//!
+//! Only one `Debug` may exist per process: the signal handlers and the panic
+//! hook find their state through the global `current` pointer set by `init`.
+//!
+//! Mechanics
+//!
+//! * Capturing another thread's stack or registers: the calling (server)
+//! thread sends `capture_signal` with `tgkill`. The SA_SIGINFO handler copies
+//! the interrupted register state (`cpu_context.fromPosixSignalContext`) into
+//! the single capture slot and parks on a futex. The server unwinds the
+//! parked thread's stack from that context, releases the target, then
+//! symbolizes. The handler is async-signal-safe: no allocation, no
+//! `std.debug`, no locks other than the futex. A target that does not run
+//! the handler within `capture_timeout_ns` (signal masked, thread in D
+//! state, ...) yields `error.Timeout`; a late-arriving handler run cannot
+//! corrupt a reused slot because it must match the requested tid and win a
+//! compare-and-swap from `armed` on the slot state (that pair plays the role
+//! of a generation counter: a stale run finds the slot idle, armed for
+//! another tid, or armed for itself, in which case its capture is simply the
+//! valid answer to the new request).
+//! * Breakpoints: `@breakpoint()` raises SIGTRAP on the executing thread only.
+//! The handler claims a pause slot, saves the context and parks on a futex
+//! until `resumeThread`. On x86_64 the saved PC is already past `int3`; on
+//! aarch64 the handler advances PC by 4 in the ucontext before returning
+//! (only for a real `brk`, i.e. a kernel-generated si_code; a SIGTRAP sent
+//! with kill/tgkill parks the thread where it was). With no free slot the
+//! thread steps over the breakpoint and keeps running (`traps_skipped`
+//! counts them): the debug layer never kills the process. The server thread
+//! itself (`server_tid`) is never parked, a breakpoint there is stepped
+//! over, because nobody could resume it. Only a stale handler run after
+//! `deinit` (no `current`) falls back to the default disposition.
+//! * Panics: `panicHook` records the message and a stack capture, then, if
+//! `hold_on_panic` and a `Debug` exists, parks until `panicContinue`; then
+//! `std.debug.defaultPanic` runs. A nested or second panic, or a panic on
+//! the server thread itself (which could never be continued), goes
+//! straight to the default handler.
+//! * std.debug's `SelfInfo` guards its state with an `Io.RwLock`. A target
+//! parked while holding it (a thread inside a stack-trace dump, say) would
+//! deadlock the unwind, so after parking a thread the lock is probed with
+//! `tryLock`; a held lock yields `error.Busy` and the target is released.
+//! * Known-module guard: `std.debug.SelfInfo` (Zig 0.16) rebuilds its module
+//! list whenever it is asked about an address outside every known module,
+//! freeing the CIE lists its unwind cache still points into; later unwinds
+//! then read freed memory. `init` records the PT_LOAD ranges of the
+//! executable (the same source std uses) and every lookup or unwind is
+//! first checked against them; addresses outside (unmapped, vDSO, ...)
+//! render as "?" and are never handed to std.
+
+const std = @import("std");
+const builtin = @import("builtin");
+const linux = std.os.linux;
+const cpu_context = std.debug.cpu_context;
+const Writer = std.Io.Writer;
+const Native = cpu_context.Native;
+const arch = builtin.cpu.arch;
+
+pub const Options = struct {
+ /// Used for `std.debug` symbolization (reading debug info from disk).
+ io: std.Io,
+ /// Fixed arena for symbol text. A `FixedBufferAllocator` is placed over it
+ /// and reset before every query. 16 KiB is plenty; 4 KiB is a sane floor.
+ text_buf: []u8,
+ /// Real-time signal used to snapshot other threads. SIGRTMIN is 32 on
+ /// Linux without libc; the default is SIGRTMIN+3.
+ capture_signal: u8 = default_capture_signal,
+ /// How long to wait for a target thread to run the capture handler.
+ capture_timeout_ns: u64 = 250 * std.time.ns_per_ms,
+ /// How many threads may be parked in `@breakpoint()` at once (≤ 32).
+ max_paused: u8 = 16,
+};
+
+pub const default_capture_signal: u8 = 32 + 3;
+
+/// Hard upper bound of `Options.max_paused` (slot storage is static).
+pub const max_paused_cap = 32;
+/// Maximum number of frames written by any stack function.
+pub const max_frames = 64;
+/// Maximum number of tids enumerated from /proc/self/task.
+pub const max_threads = 512;
+/// Upper bound of the recorded panic message.
+pub const panic_msg_cap = 1024;
+/// Maximum number of PT_LOAD ranges recorded by the known-module guard.
+pub const max_ranges = 64;
+
+/// Consulted by `panicHook`: when true and a `Debug` is initialized, the
+/// panicking thread is held until `panicContinue`.
+pub var hold_on_panic: bool = true;
+
+/// The one live instance, set by `init`, cleared by `deinit`.
+pub var current: ?*Debug = null;
+
+/// The tid of the thread serving requests (0 = none). That thread is never
+/// parked by a breakpoint or held by a panic, since nobody could release it.
+pub var server_tid: std.atomic.Value(u32) = .init(0);
+
+/// Breakpoints stepped over because no pause slot was free, or because they
+/// were hit on the server thread.
+pub var traps_skipped: std.atomic.Value(u32) = .init(0);
+
+pub const Error = error{
+ /// The target thread did not run the capture handler in time.
+ Timeout,
+ /// No thread with that tid exists in this process.
+ NoThread,
+ /// The address is not mapped (EFAULT from process_vm_readv/writev).
+ Unmapped,
+ /// The thread is not parked in a breakpoint.
+ NotPaused,
+ /// No panic has been recorded / is being held.
+ NoPanic,
+ /// Another `Debug` already exists in this process.
+ AlreadyInitialized,
+ /// The operation is not available on this architecture / kernel.
+ Unsupported,
+ /// The target thread is parked inside std.debug (holding its lock); its
+ /// stack cannot be unwound without deadlocking. Retry later.
+ Busy,
+ /// Invalid option value.
+ InvalidOptions,
+ /// A syscall or /proc read failed unexpectedly.
+ Unexpected,
+ /// The writer failed.
+ WriteFailed,
+};
+
+// Capture slot states.
+const cap_idle: u32 = 0;
+const cap_armed: u32 = 1;
+const cap_capturing: u32 = 2;
+const cap_captured: u32 = 3;
+const cap_failed: u32 = 4;
+
+// Pause slot states.
+const pause_free: u32 = 0;
+const pause_claimed: u32 = 1;
+const pause_paused: u32 = 2;
+const pause_resuming: u32 = 3;
+
+const CaptureSlot = struct {
+ state: std.atomic.Value(u32) = .init(cap_idle),
+ target_tid: std.atomic.Value(u32) = .init(0),
+ ctx: Native = undefined,
+};
+
+const PauseSlot = struct {
+ state: std.atomic.Value(u32) = .init(pause_free),
+ tid: std.atomic.Value(u32) = .init(0),
+ ctx: Native = undefined,
+};
+
+pub const Debug = struct {
+ io: std.Io,
+ text_buf: []u8,
+ capture_signal: linux.SIG,
+ capture_timeout_ns: u64,
+ max_paused: u8,
+
+ capture: CaptureSlot = .{},
+ paused: [max_paused_cap]PauseSlot = [_]PauseSlot{.{}} ** max_paused_cap,
+
+ old_capture_action: linux.Sigaction = undefined,
+ old_trap_action: linux.Sigaction = undefined,
+ breakpoints_enabled: bool = false,
+
+ tids: [max_threads]u32 = undefined,
+ tid_count: usize = 0,
+
+ ranges: [max_ranges]Range = undefined,
+ range_count: usize = 0,
+
+ const Range = struct { start: usize, len: usize };
+
+ /// Installs the capture handler (not the SIGTRAP handler) and publishes
+ /// `d` as `current`.
+ pub fn init(d: *Debug, opts: Options) Error!void {
+ if (current != null) return error.AlreadyInitialized;
+ if (opts.capture_signal < 32 or opts.capture_signal >= linux.NSIG) return error.InvalidOptions;
+ if (opts.max_paused == 0 or opts.max_paused > max_paused_cap) return error.InvalidOptions;
+ if (Native == noreturn) return error.Unsupported;
+ d.* = .{
+ .io = opts.io,
+ .text_buf = opts.text_buf,
+ .capture_signal = @enumFromInt(opts.capture_signal),
+ .capture_timeout_ns = opts.capture_timeout_ns,
+ .max_paused = opts.max_paused,
+ };
+ d.scanModules();
+ const act: linux.Sigaction = .{
+ .handler = .{ .sigaction = captureHandler },
+ .mask = linux.sigemptyset(),
+ .flags = linux.SA.SIGINFO | linux.SA.RESTART,
+ };
+ current = d;
+ if (linux.errno(linux.sigaction(d.capture_signal, &act, &d.old_capture_action)) != .SUCCESS) {
+ current = null;
+ return error.Unexpected;
+ }
+ }
+
+ /// Restores the signal dispositions and clears `current`. Threads parked
+ /// in a breakpoint are resumed first.
+ pub fn deinit(d: *Debug) void {
+ d.disableBreakpoints();
+ _ = linux.sigaction(d.capture_signal, &d.old_capture_action, null);
+ if (current == d) current = null;
+ }
+
+ /// Installs the SIGTRAP handler so that `@breakpoint()` parks the thread.
+ pub fn enableBreakpoints(d: *Debug) Error!void {
+ if (d.breakpoints_enabled) return;
+ if (arch != .x86_64 and !arch.isAARCH64()) return error.Unsupported;
+ const act: linux.Sigaction = .{
+ .handler = .{ .sigaction = trapHandler },
+ .mask = linux.sigemptyset(),
+ .flags = linux.SA.SIGINFO | linux.SA.RESTART,
+ };
+ if (linux.errno(linux.sigaction(.TRAP, &act, &d.old_trap_action)) != .SUCCESS) return error.Unexpected;
+ d.breakpoints_enabled = true;
+ }
+
+ /// Restores the previous SIGTRAP disposition and resumes every parked thread.
+ pub fn disableBreakpoints(d: *Debug) void {
+ if (!d.breakpoints_enabled) return;
+ _ = linux.sigaction(.TRAP, &d.old_trap_action, null);
+ d.breakpoints_enabled = false;
+ for (&d.paused) |*slot| {
+ if (slot.state.cmpxchgStrong(pause_paused, pause_resuming, .acq_rel, .acquire) == null)
+ futexWake(&slot.state);
+ }
+ }
+
+ // ---------------------------------------------------------------- threads
+
+ /// The nth tid of this process, numerically sorted; null past the end.
+ /// Index 0 rescans /proc/self/task; higher indices reuse that scan.
+ pub fn threadAt(d: *Debug, index: usize) ?u32 {
+ if (index == 0 or d.tid_count == 0) d.scanThreads();
+ if (index >= d.tid_count) return null;
+ return d.tids[index];
+ }
+
+ pub fn threadExists(d: *Debug, tid: u32) bool {
+ _ = d;
+ var path_buf: [64]u8 = undefined;
+ const path = std.fmt.bufPrintZ(&path_buf, "/proc/self/task/{d}/comm", .{tid}) catch return false;
+ var buf: [32]u8 = undefined;
+ _ = readFile(path, &buf) catch return false;
+ return true;
+ }
+
+ /// The thread's comm (without the trailing newline).
+ pub fn threadName(d: *Debug, tid: u32, w: *Writer) Error!void {
+ _ = d;
+ var path_buf: [64]u8 = undefined;
+ const path = std.fmt.bufPrintZ(&path_buf, "/proc/self/task/{d}/comm", .{tid}) catch return error.Unexpected;
+ var buf: [64]u8 = undefined;
+ const text = readFile(path, &buf) catch |err| switch (err) {
+ error.NotFound => return error.NoThread,
+ else => return error.Unexpected,
+ };
+ w.writeAll(std.mem.trimEnd(u8, text, "\n")) catch return error.WriteFailed;
+ }
+
+ /// A few fields of /proc/self/task/<tid>/stat, one "name value" per line:
+ /// state, utime, stime, minflt, majflt, priority, nice, processor.
+ pub fn threadStat(d: *Debug, tid: u32, w: *Writer) Error!void {
+ _ = d;
+ var path_buf: [64]u8 = undefined;
+ const path = std.fmt.bufPrintZ(&path_buf, "/proc/self/task/{d}/stat", .{tid}) catch return error.Unexpected;
+ var buf: [1024]u8 = undefined;
+ const text = readFile(path, &buf) catch |err| switch (err) {
+ error.NotFound => return error.NoThread,
+ else => return error.Unexpected,
+ };
+ // "<pid> (<comm>) S <fields...>"; comm may contain spaces and parens.
+ const close = std.mem.lastIndexOfScalar(u8, text, ')') orelse return error.Unexpected;
+ var it = std.mem.tokenizeScalar(u8, text[close + 1 ..], ' ');
+ // Field numbers below are 0-based from `state`.
+ const wanted = [_]struct { idx: usize, name: []const u8 }{
+ .{ .idx = 0, .name = "state" },
+ .{ .idx = 11, .name = "utime" },
+ .{ .idx = 12, .name = "stime" },
+ .{ .idx = 7, .name = "minflt" },
+ .{ .idx = 9, .name = "majflt" },
+ .{ .idx = 15, .name = "priority" },
+ .{ .idx = 16, .name = "nice" },
+ .{ .idx = 36, .name = "processor" },
+ };
+ var fields: [40][]const u8 = undefined;
+ var n: usize = 0;
+ while (it.next()) |f| : (n += 1) {
+ if (n == fields.len) break;
+ fields[n] = f;
+ }
+ for (wanted) |want| {
+ const value = if (want.idx < n) fields[want.idx] else "?";
+ w.print("{s} {s}\n", .{ want.name, value }) catch return error.WriteFailed;
+ }
+ }
+
+ /// "#n 0x<addr> in <fn> (<file>:<line>:<col>)" per frame. The calling
+ /// thread unwinds itself directly; any other thread is captured with the
+ /// capture signal.
+ pub fn threadStack(d: *Debug, tid: u32, w: *Writer) Error!void {
+ var addrs: [max_frames]usize = undefined;
+ var trace: std.debug.StackTrace = undefined;
+ if (tid == selfTid()) {
+ trace = std.debug.captureCurrentStackTrace(.{}, &addrs);
+ } else {
+ try d.captureThread(tid);
+ if (!d.selfInfoFree()) {
+ d.releaseCapture();
+ return error.Busy;
+ }
+ trace = d.unwindContext(&d.capture.ctx, &addrs);
+ d.releaseCapture();
+ }
+ try d.writeFrames(trace.return_addresses, w);
+ }
+
+ /// "<reg> 0x<hex>" per general register, plus pc/sp/fp aliases.
+ pub fn threadRegs(d: *Debug, tid: u32, w: *Writer) Error!void {
+ if (tid == selfTid()) {
+ const ctx = Native.current();
+ return writeRegs(&ctx, w);
+ }
+ try d.captureThread(tid);
+ const ctx = d.capture.ctx;
+ d.releaseCapture();
+ return writeRegs(&ctx, w);
+ }
+
+ // ------------------------------------------------------ addresses & memory
+
+ /// "fn\nfile:line:col\nmodule\n", unknown parts as "?".
+ pub fn resolveAddr(d: *Debug, addr: usize, w: *Writer) Error!void {
+ if (!d.knownCode(addr)) return w.writeAll("?\n?\n?\n") catch error.WriteFailed;
+ var fba = std.heap.FixedBufferAllocator.init(d.text_buf);
+ const alloc = fba.allocator();
+ const di = std.debug.getSelfDebugInfo() catch return error.Unsupported;
+ var sym = std.debug.Symbol.unknown;
+ var symbols: std.ArrayList(std.debug.Symbol) = .empty;
+ if (di.getSymbols(d.io, alloc, alloc, addr, true, &symbols)) {
+ if (symbols.items.len > 0) sym = symbols.items[0];
+ } else |_| {}
+ w.print("{s}\n", .{sym.name orelse "?"}) catch return error.WriteFailed;
+ if (sym.source_location) |sl| {
+ w.print("{s}:{d}:{d}\n", .{ sl.file_name, sl.line, sl.column }) catch return error.WriteFailed;
+ } else {
+ w.writeAll("?\n") catch return error.WriteFailed;
+ }
+ const module = di.getModuleName(d.io, addr) catch "?";
+ w.print("{s}\n", .{module}) catch return error.WriteFailed;
+ }
+
+ /// Reads `buf.len` bytes at `addr` via process_vm_readv on the own
+ /// process. Never faults. Returns the number of bytes read (short when the
+ /// range crosses into an unmapped page); `error.Unmapped` when nothing
+ /// could be read.
+ pub fn readMem(d: *Debug, addr: usize, buf: []u8) Error!usize {
+ _ = d;
+ if (buf.len == 0) return 0;
+ // Page 0 is never mapped (mmap_min_addr) and a null `iovec.base` is a
+ // safety-checked cast; the same answer without the trap.
+ if (addr == 0) return error.Unmapped;
+ const local = [_]std.posix.iovec{.{ .base = buf.ptr, .len = buf.len }};
+ const remote = [_]std.posix.iovec_const{.{ .base = @ptrFromInt(addr), .len = buf.len }};
+ const rc = linux.process_vm_readv(linux.getpid(), &local, &remote, 0);
+ switch (linux.errno(rc)) {
+ .SUCCESS => return rc,
+ .FAULT => return error.Unmapped,
+ .NOSYS, .PERM => return error.Unsupported,
+ else => return error.Unexpected,
+ }
+ }
+
+ /// Writes `data` at `addr` via process_vm_writev. Read-only mappings also
+ /// report `error.Unmapped` (the kernel says EFAULT for both).
+ pub fn writeMem(d: *Debug, addr: usize, data: []const u8) Error!usize {
+ _ = d;
+ if (data.len == 0) return 0;
+ if (addr == 0) return error.Unmapped;
+ const local = [_]std.posix.iovec_const{.{ .base = data.ptr, .len = data.len }};
+ const remote = [_]std.posix.iovec_const{.{ .base = @ptrFromInt(addr), .len = data.len }};
+ const rc = linux.process_vm_writev(linux.getpid(), &local, &remote, 0);
+ switch (linux.errno(rc)) {
+ .SUCCESS => return rc,
+ .FAULT => return error.Unmapped,
+ .NOSYS, .PERM => return error.Unsupported,
+ else => return error.Unexpected,
+ }
+ }
+
+ /// Hexdump of `len` bytes at `addr` in the shape of `std.debug.dumpHex`
+ /// (16 bytes per line, address column, bytes in two groups, ASCII column).
+ /// Stops early at the first unmapped byte; `error.Unmapped` only when the
+ /// very first chunk is unreadable.
+ pub fn hexdump(d: *Debug, addr: usize, len: usize, w: *Writer) Error!void {
+ var chunk: [256]u8 = undefined;
+ var done: usize = 0;
+ while (done < len) {
+ const want = @min(chunk.len, len - done);
+ const got = d.readMem(addr +% done, chunk[0..want]) catch |err| switch (err) {
+ error.Unmapped => if (done == 0) return error.Unmapped else break,
+ else => return err,
+ };
+ if (got == 0) break;
+ try writeHexLines(addr +% done, chunk[0..got], w);
+ done += got;
+ if (got < want) break;
+ }
+ }
+
+ /// Copies /proc/self/maps to `w`.
+ pub fn maps(d: *Debug, w: *Writer) Error!void {
+ _ = d;
+ return streamFile("/proc/self/maps", w);
+ }
+
+ /// Reads `buf.len` bytes of /proc/self/maps at `offset` (0 at the end).
+ /// Not a consistent snapshot across reads; a map appearing between two
+ /// reads shifts the text, like `cat` on /proc itself.
+ pub fn readMaps(d: *Debug, offset: u64, buf: []u8) Error!usize {
+ _ = d;
+ if (offset > std.math.maxInt(i64)) return 0;
+ return preadFile("/proc/self/maps", offset, buf);
+ }
+
+ // ------------------------------------------------------------ breakpoints
+
+ /// The nth tid currently parked in `@breakpoint()`.
+ pub fn pausedAt(d: *Debug, index: usize) ?u32 {
+ var n: usize = 0;
+ for (d.paused[0..d.max_paused]) |*slot| {
+ if (slot.state.load(.acquire) != pause_paused) continue;
+ if (n == index) return slot.tid.load(.acquire);
+ n += 1;
+ }
+ return null;
+ }
+
+ pub fn isPaused(d: *Debug, tid: u32) bool {
+ return d.pausedSlot(tid) != null;
+ }
+
+ pub fn pausedStack(d: *Debug, tid: u32, w: *Writer) Error!void {
+ const slot = d.pausedSlot(tid) orelse return error.NotPaused;
+ if (!d.selfInfoFree()) return error.Busy;
+ var addrs: [max_frames]usize = undefined;
+ const trace = d.unwindContext(&slot.ctx, &addrs);
+ try d.writeFrames(trace.return_addresses, w);
+ }
+
+ pub fn pausedRegs(d: *Debug, tid: u32, w: *Writer) Error!void {
+ const slot = d.pausedSlot(tid) orelse return error.NotPaused;
+ return writeRegs(&slot.ctx, w);
+ }
+
+ /// Lets a parked thread continue past its breakpoint.
+ pub fn resumeThread(d: *Debug, tid: u32) Error!void {
+ const slot = d.pausedSlot(tid) orelse return error.NotPaused;
+ if (slot.state.cmpxchgStrong(pause_paused, pause_resuming, .acq_rel, .acquire) != null) return error.NotPaused;
+ futexWake(&slot.state);
+ }
+
+ fn pausedSlot(d: *Debug, tid: u32) ?*PauseSlot {
+ for (d.paused[0..d.max_paused]) |*slot| {
+ if (slot.state.load(.acquire) == pause_paused and slot.tid.load(.acquire) == tid) return slot;
+ }
+ return null;
+ }
+
+ // ------------------------------------------------------------------ panic
+
+ /// The recorded panic message; nothing before any panic.
+ pub fn panicMessage(d: *Debug, w: *Writer) Error!void {
+ _ = d;
+ if (panic_state.load(.acquire) == panic_none) return;
+ w.writeAll(panic_msg[0..panic_msg_len]) catch return error.WriteFailed;
+ }
+
+ /// Frames of the panicking thread, symbolized lazily.
+ pub fn panicStack(d: *Debug, w: *Writer) Error!void {
+ if (panic_state.load(.acquire) == panic_none) return;
+ try d.writeFrames(panic_addrs[0..panic_addr_count], w);
+ }
+
+ /// True while a panicking thread is parked waiting for `panicContinue`.
+ pub fn panicHeld(d: *Debug) bool {
+ _ = d;
+ return panic_state.load(.acquire) == panic_held;
+ }
+
+ /// Releases the held panicking thread into `std.debug.defaultPanic`.
+ pub fn panicContinue(d: *Debug) Error!void {
+ _ = d;
+ if (panic_state.cmpxchgStrong(panic_held, panic_continued, .acq_rel, .acquire) != null) return error.NoPanic;
+ futexWake(&panic_state);
+ }
+
+ // -------------------------------------------------------------- internals
+
+ fn scanThreads(d: *Debug) void {
+ d.tid_count = 0;
+ const fd_rc = linux.open("/proc/self/task", .{ .ACCMODE = .RDONLY, .DIRECTORY = true, .CLOEXEC = true }, 0);
+ if (linux.errno(fd_rc) != .SUCCESS) return;
+ const fd: i32 = @intCast(fd_rc);
+ defer _ = linux.close(fd);
+ var buf: [4096]u8 align(@alignOf(linux.dirent64)) = undefined;
+ while (true) {
+ const rc = linux.getdents64(fd, &buf, buf.len);
+ if (linux.errno(rc) != .SUCCESS or rc == 0) break;
+ var off: usize = 0;
+ while (off < rc) {
+ const ent: *align(1) const linux.dirent64 = @ptrCast(&buf[off]);
+ const name_ptr: [*:0]const u8 = @ptrCast(&buf[off + @offsetOf(linux.dirent64, "name")]);
+ const name = std.mem.span(name_ptr);
+ if (std.fmt.parseInt(u32, name, 10)) |tid| {
+ if (d.tid_count < max_threads) {
+ d.tids[d.tid_count] = tid;
+ d.tid_count += 1;
+ }
+ } else |_| {}
+ off += ent.reclen;
+ }
+ }
+ std.mem.sort(u32, d.tids[0..d.tid_count], {}, std.sort.asc(u32));
+ }
+
+ /// Arms the capture slot for `tid`, signals it and waits until the handler
+ /// has parked with its context copied. On success the caller owns the
+ /// slot until `releaseCapture`.
+ fn captureThread(d: *Debug, tid: u32) Error!void {
+ const slot = &d.capture;
+ slot.target_tid.store(tid, .release);
+ slot.state.store(cap_armed, .release);
+ const rc = linux.tgkill(linux.getpid(), @intCast(tid), d.capture_signal);
+ switch (linux.errno(rc)) {
+ .SUCCESS => {},
+ .SRCH => {
+ slot.state.store(cap_idle, .release);
+ return error.NoThread;
+ },
+ else => {
+ slot.state.store(cap_idle, .release);
+ return error.Unexpected;
+ },
+ }
+ const deadline = monotonicNs() + d.capture_timeout_ns;
+ while (true) {
+ const s = slot.state.load(.acquire);
+ switch (s) {
+ cap_captured => return,
+ cap_failed => {
+ slot.state.store(cap_idle, .release);
+ return error.Unsupported;
+ },
+ cap_armed => {
+ const now = monotonicNs();
+ if (now >= deadline) {
+ // Disarm; if the handler raced us it has moved on to
+ // `capturing` and we simply keep waiting for it.
+ if (slot.state.cmpxchgStrong(cap_armed, cap_idle, .acq_rel, .acquire) == null) return error.Timeout;
+ continue;
+ }
+ futexWaitNs(&slot.state, cap_armed, deadline - now);
+ },
+ // The handler is copying registers; it finishes promptly.
+ cap_capturing => futexWaitNs(&slot.state, cap_capturing, 1 * std.time.ns_per_ms),
+ else => unreachable,
+ }
+ }
+ }
+
+ /// Records the PT_LOAD ranges of every module `dl_iterate_phdr` reports
+ /// (for a static executable: the executable itself, not the vDSO).
+ fn scanModules(d: *Debug) void {
+ d.range_count = 0;
+ std.posix.dl_iterate_phdr(d, error{}, struct {
+ fn cb(info: *std.posix.dl_phdr_info, _: usize, ctx: *Debug) error{}!void {
+ for (info.phdr[0..info.phnum]) |phdr| {
+ if (phdr.type != .LOAD) continue;
+ if (ctx.range_count == max_ranges) return;
+ ctx.ranges[ctx.range_count] = .{ .start = info.addr +% phdr.vaddr, .len = phdr.memsz };
+ ctx.range_count += 1;
+ }
+ }
+ }.cb) catch {};
+ }
+
+ /// True when `addr` lies in a module `std.debug` already knows about, so
+ /// that asking it about `addr` cannot trigger a module rescan.
+ fn knownCode(d: *const Debug, addr: usize) bool {
+ for (d.ranges[0..d.range_count]) |r| {
+ if (addr >= r.start and addr - r.start < r.len) return true;
+ }
+ return false;
+ }
+
+ /// Unwinds from a saved context. A pc outside every known module (e.g. a
+ /// thread inside the vDSO) is reported as a single frame and not unwound,
+ /// because std would otherwise rescan its module list (see the header).
+ fn unwindContext(d: *const Debug, ctx: *const Native, addrs: *[max_frames]usize) std.debug.StackTrace {
+ if (!d.knownCode(ctx.getPc())) {
+ addrs[0] = ctx.getPc() +| 1;
+ return .{ .return_addresses = addrs[0..1], .skipped = .unknown };
+ }
+ return std.debug.captureCurrentStackTrace(.{ .context = ctx }, addrs);
+ }
+
+ /// True when nobody holds std.debug's `SelfInfo` lock right now. Called
+ /// with the target parked, so a held lock means the *target* (or another
+ /// live thread, which will let go) holds it; only the former deadlocks,
+ /// and the caller cannot tell them apart, so both yield `error.Busy`.
+ fn selfInfoFree(d: *const Debug) bool {
+ if (comptime !@hasField(std.debug.SelfInfo, "rwlock")) return true;
+ const di = std.debug.getSelfDebugInfo() catch return true;
+ if (!di.rwlock.tryLock(d.io)) return false;
+ di.rwlock.unlock(d.io);
+ return true;
+ }
+
+ fn releaseCapture(d: *Debug) void {
+ d.capture.state.store(cap_idle, .release);
+ futexWake(&d.capture.state);
+ }
+
+ fn writeFrames(d: *Debug, addrs: []const usize, w: *Writer) Error!void {
+ var fba = std.heap.FixedBufferAllocator.init(d.text_buf);
+ const alloc = fba.allocator();
+ const di = std.debug.getSelfDebugInfo() catch return error.Unsupported;
+ for (addrs, 0..) |ret_addr, i| {
+ // Return addresses point after the call; the first frame of a
+ // context capture is stored as pc+1 by std for the same reason.
+ const addr = ret_addr -| 1;
+ fba.reset();
+ var symbols: std.ArrayList(std.debug.Symbol) = .empty;
+ var sym = std.debug.Symbol.unknown;
+ if (d.knownCode(addr)) {
+ if (di.getSymbols(d.io, alloc, alloc, addr, true, &symbols)) {
+ if (symbols.items.len > 0) sym = symbols.items[0];
+ } else |_| {}
+ }
+ w.print("#{d} 0x{x} in {s} (", .{ i, addr, sym.name orelse "?" }) catch return error.WriteFailed;
+ if (sym.source_location) |sl| {
+ w.print("{s}:{d}:{d})\n", .{ sl.file_name, sl.line, sl.column }) catch return error.WriteFailed;
+ } else {
+ w.writeAll("?)\n") catch return error.WriteFailed;
+ }
+ }
+ }
+};
+
+// ------------------------------------------------------------------ handlers
+
+fn selfTid() u32 {
+ return @intCast(linux.gettid());
+}
+
+fn captureHandler(_: linux.SIG, _: *const linux.siginfo_t, ctx_ptr: ?*anyopaque) callconv(.c) void {
+ const d = current orelse return;
+ const slot = &d.capture;
+ const me = selfTid();
+ if (slot.target_tid.load(.acquire) != me) return;
+ if (slot.state.cmpxchgStrong(cap_armed, cap_capturing, .acq_rel, .acquire) != null) return;
+ // The tid check and the swap are not one atomic step: a stale run (a
+ // signal that stayed pending while its request timed out) may have read
+ // the old tid and then won the swap of a request re-armed for another
+ // thread. `target_tid` is fixed while the slot is armed, so re-checking
+ // after the swap closes the window; hand the slot back untouched.
+ if (slot.target_tid.load(.acquire) != me) {
+ slot.state.store(cap_armed, .release);
+ futexWake(&slot.state);
+ return;
+ }
+ if (cpu_context.fromPosixSignalContext(ctx_ptr)) |ctx| {
+ slot.ctx = ctx;
+ slot.state.store(cap_captured, .release);
+ futexWake(&slot.state);
+ while (slot.state.load(.acquire) == cap_captured) futexWaitNs(&slot.state, cap_captured, null);
+ } else {
+ slot.state.store(cap_failed, .release);
+ futexWake(&slot.state);
+ }
+}
+
+/// aarch64 Linux ucontext_t, only as far as `mcontext.pc` (see
+/// std.debug.cpu_context's signal_ucontext_t).
+const UcontextAarch64 = extern struct {
+ flags: usize,
+ link: ?*UcontextAarch64,
+ stack: linux.stack_t,
+ sigmask: linux.sigset_t,
+ unused: [120]u8,
+ mcontext: extern struct {
+ fault_address: u64 align(16),
+ x: [30]u64,
+ lr: u64,
+ sp: u64,
+ pc: u64,
+ },
+};
+
+fn trapHandler(_: linux.SIG, info: *const linux.siginfo_t, ctx_ptr: ?*anyopaque) callconv(.c) void {
+ const d = current orelse return trapFallback();
+ const ctx = cpu_context.fromPosixSignalContext(ctx_ptr) orelse return trapFallback();
+ // si_code > 0 is kernel-generated (TRAP_BRKPT for int3/brk); <= 0 is
+ // kill/tgkill/sigqueue from user space, where PC points at the
+ // interrupted instruction and must not be touched.
+ const from_instruction = info.code > 0;
+ if (comptime arch.isAARCH64()) {
+ // `brk #imm` does not advance PC; step over it so returning from the
+ // handler does not re-trap.
+ if (from_instruction) {
+ const uc: *UcontextAarch64 = @ptrCast(@alignCast(ctx_ptr.?));
+ uc.mcontext.pc += 4;
+ }
+ } else if (comptime arch != .x86_64) {
+ return trapFallback();
+ }
+ const tid = selfTid();
+ if (tid == server_tid.load(.acquire)) {
+ // Nobody could resume the thread that serves /breakpoints: step over.
+ _ = traps_skipped.fetchAdd(1, .acq_rel);
+ return;
+ }
+ const slot: *PauseSlot = for (d.paused[0..d.max_paused]) |*slot| {
+ if (slot.state.cmpxchgStrong(pause_free, pause_claimed, .acq_rel, .acquire) == null) break slot;
+ } else {
+ _ = traps_skipped.fetchAdd(1, .acq_rel);
+ return;
+ };
+ slot.ctx = ctx;
+ slot.tid.store(tid, .release);
+ slot.state.store(pause_paused, .release);
+ while (slot.state.load(.acquire) == pause_paused) futexWaitNs(&slot.state, pause_paused, null);
+ slot.state.store(pause_free, .release);
+}
+
+/// Restores the default SIGTRAP disposition and re-raises it: the signal is
+/// blocked while the handler runs, so it is delivered (fatally) on return.
+/// Only for a handler run with no `Debug` (a trap in flight during `deinit`)
+/// or on an architecture whose context cannot be read.
+fn trapFallback() void {
+ const act: linux.Sigaction = .{
+ .handler = .{ .handler = linux.SIG.DFL },
+ .mask = linux.sigemptyset(),
+ .flags = 0,
+ };
+ _ = linux.sigaction(.TRAP, &act, null);
+ _ = linux.tkill(linux.gettid(), .TRAP);
+}
+
+// --------------------------------------------------------------------- panic
+
+const panic_none: u32 = 0;
+const panic_recording: u32 = 1;
+const panic_recorded: u32 = 2;
+const panic_held: u32 = 3;
+const panic_continued: u32 = 4;
+
+var panic_state: std.atomic.Value(u32) = .init(panic_none);
+var panic_msg: [panic_msg_cap]u8 = undefined;
+var panic_msg_len: usize = 0;
+var panic_addrs: [max_frames]usize = undefined;
+var panic_addr_count: usize = 0;
+/// The tid of the panicking thread (0 before any panic).
+pub var panic_tid: u32 = 0;
+
+/// Records the first panic: message (bounded copy) and stack addresses.
+/// Returns false if a panic was already recorded (nested or second panic).
+pub fn recordPanic(msg: []const u8, first_trace_addr: ?usize) bool {
+ if (panic_state.cmpxchgStrong(panic_none, panic_recording, .acq_rel, .acquire) != null) return false;
+ panic_tid = selfTid();
+ panic_msg_len = @min(msg.len, panic_msg.len);
+ @memcpy(panic_msg[0..panic_msg_len], msg[0..panic_msg_len]);
+ const trace = std.debug.captureCurrentStackTrace(.{ .first_address = first_trace_addr }, &panic_addrs);
+ panic_addr_count = trace.return_addresses.len;
+ panic_state.store(panic_recorded, .release);
+ return true;
+}
+
+/// Parks the panicking thread until `Debug.panicContinue` when holding is
+/// enabled and a `Debug` exists; then hands over to `std.debug.defaultPanic`.
+pub fn panicHook(msg: []const u8, first_trace_addr: ?usize) noreturn {
+ @branchHint(.cold);
+ if (recordPanic(msg, first_trace_addr)) {
+ // The server thread cannot be held: it is the one that would have to
+ // serve /panic/ctl.
+ if (hold_on_panic and current != null and panic_tid != server_tid.load(.acquire)) {
+ if (panic_state.cmpxchgStrong(panic_recorded, panic_held, .acq_rel, .acquire) == null) {
+ while (panic_state.load(.acquire) == panic_held) futexWaitNs(&panic_state, panic_held, null);
+ }
+ }
+ }
+ std.debug.defaultPanic(msg, first_trace_addr);
+}
+
+/// Clears the recorded panic. Only meaningful in tests of the record path.
+pub fn resetPanicRecord() void {
+ panic_msg_len = 0;
+ panic_addr_count = 0;
+ panic_tid = 0;
+ panic_state.store(panic_none, .release);
+}
+
+// ------------------------------------------------------------------- helpers
+
+fn futexWake(word: *std.atomic.Value(u32)) void {
+ _ = linux.futex_3arg(&word.raw, .{ .cmd = .WAKE, .private = true }, std.math.maxInt(u32));
+}
+
+/// Waits while `*word == expect`, at most `timeout_ns` (forever when null).
+/// Returns on wake, timeout, value change or EINTR; callers loop.
+fn futexWaitNs(word: *std.atomic.Value(u32), expect: u32, timeout_ns: ?u64) void {
+ var ts: linux.timespec = undefined;
+ const ts_ptr: ?*const linux.timespec = if (timeout_ns) |ns| blk: {
+ ts = .{ .sec = @intCast(ns / std.time.ns_per_s), .nsec = @intCast(ns % std.time.ns_per_s) };
+ break :blk &ts;
+ } else null;
+ _ = linux.futex_4arg(&word.raw, .{ .cmd = .WAIT, .private = true }, expect, ts_ptr);
+}
+
+fn monotonicNs() u64 {
+ var ts: linux.timespec = undefined;
+ _ = linux.clock_gettime(.MONOTONIC, &ts);
+ return @as(u64, @intCast(ts.sec)) * std.time.ns_per_s + @as(u64, @intCast(ts.nsec));
+}
+
+const FileError = error{ NotFound, Unexpected, TooBig };
+
+/// Reads a whole (small) file with raw syscalls.
+fn readFile(path: [*:0]const u8, buf: []u8) FileError![]u8 {
+ const fd_rc = linux.open(path, .{ .ACCMODE = .RDONLY, .CLOEXEC = true }, 0);
+ switch (linux.errno(fd_rc)) {
+ .SUCCESS => {},
+ .NOENT, .SRCH => return error.NotFound,
+ else => return error.Unexpected,
+ }
+ const fd: i32 = @intCast(fd_rc);
+ defer _ = linux.close(fd);
+ var len: usize = 0;
+ while (len < buf.len) {
+ const rc = linux.read(fd, buf[len..].ptr, buf.len - len);
+ switch (linux.errno(rc)) {
+ .SUCCESS => {},
+ .INTR => continue,
+ .SRCH, .NOENT => return error.NotFound,
+ else => return error.Unexpected,
+ }
+ if (rc == 0) return buf[0..len];
+ len += rc;
+ }
+ return error.TooBig;
+}
+
+/// One pread of `buf.len` bytes at `offset`; 0 at the end of the file.
+fn preadFile(path: [*:0]const u8, offset: u64, buf: []u8) Error!usize {
+ const fd_rc = linux.open(path, .{ .ACCMODE = .RDONLY, .CLOEXEC = true }, 0);
+ if (linux.errno(fd_rc) != .SUCCESS) return error.Unexpected;
+ const fd: i32 = @intCast(fd_rc);
+ defer _ = linux.close(fd);
+ var len: usize = 0;
+ while (len < buf.len) {
+ const rc = linux.pread(fd, buf[len..].ptr, buf.len - len, @intCast(offset + len));
+ switch (linux.errno(rc)) {
+ .SUCCESS => {},
+ .INTR => continue,
+ else => return error.Unexpected,
+ }
+ if (rc == 0) break;
+ len += rc;
+ }
+ return len;
+}
+
+/// Streams a file of any size to `w`.
+fn streamFile(path: [*:0]const u8, w: *Writer) Error!void {
+ const fd_rc = linux.open(path, .{ .ACCMODE = .RDONLY, .CLOEXEC = true }, 0);
+ if (linux.errno(fd_rc) != .SUCCESS) return error.Unexpected;
+ const fd: i32 = @intCast(fd_rc);
+ defer _ = linux.close(fd);
+ var buf: [4096]u8 = undefined;
+ while (true) {
+ const rc = linux.read(fd, &buf, buf.len);
+ switch (linux.errno(rc)) {
+ .SUCCESS => {},
+ .INTR => continue,
+ else => return error.Unexpected,
+ }
+ if (rc == 0) return;
+ w.writeAll(buf[0..rc]) catch return error.WriteFailed;
+ }
+}
+
+fn writeHexLines(base: usize, bytes: []const u8, w: *Writer) Error!void {
+ var offset: usize = 0;
+ while (offset < bytes.len) : (offset += 16) {
+ const line = bytes[offset..@min(offset + 16, bytes.len)];
+ w.print("{x:0>[1]} ", .{ base +% offset, @sizeOf(usize) * 2 }) catch return error.WriteFailed;
+ for (line, 0..) |byte, i| {
+ w.print("{X:0>2} ", .{byte}) catch return error.WriteFailed;
+ if (i == 7) w.writeByte(' ') catch return error.WriteFailed;
+ }
+ w.writeByte(' ') catch return error.WriteFailed;
+ if (line.len < 16) {
+ var missing = (16 - line.len) * 3;
+ if (line.len < 8) missing += 1;
+ w.splatByteAll(' ', missing) catch return error.WriteFailed;
+ }
+ for (line) |byte| {
+ w.writeByte(if (std.ascii.isPrint(byte)) byte else '.') catch return error.WriteFailed;
+ }
+ w.writeByte('\n') catch return error.WriteFailed;
+ }
+}
+
+fn writeRegs(ctx: *const Native, w: *Writer) Error!void {
+ if (comptime arch == .x86_64) {
+ inline for (@typeInfo(Native.Gpr).@"enum".fields) |f| {
+ w.print("{s} 0x{x}\n", .{ f.name, ctx.gprs.get(@field(Native.Gpr, f.name)) }) catch return error.WriteFailed;
+ }
+ w.print("pc 0x{x}\nsp 0x{x}\nfp 0x{x}\n", .{
+ ctx.gprs.get(.rip), ctx.gprs.get(.rsp), ctx.gprs.get(.rbp),
+ }) catch return error.WriteFailed;
+ } else if (comptime arch.isAARCH64()) {
+ for (ctx.x, 0..) |x, i| w.print("x{d} 0x{x}\n", .{ i, x }) catch return error.WriteFailed;
+ w.print("sp 0x{x}\npc 0x{x}\nfp 0x{x}\nlr 0x{x}\n", .{
+ ctx.sp, ctx.pc, ctx.x[29], ctx.x[30],
+ }) catch return error.WriteFailed;
+ } else {
+ w.print("pc 0x{x}\nfp 0x{x}\n", .{ ctx.getPc(), ctx.getFp() }) catch return error.WriteFailed;
+ }
+}
+
+// --------------------------------------------------------------------- tests
+
+const testing = std.testing;
+
+fn testOptions(text_buf: []u8) Options {
+ return .{ .io = testing.io, .text_buf = text_buf };
+}
+
+noinline fn sleepMs(ms: u64) void {
+ var ts: linux.timespec = .{ .sec = @intCast(ms / 1000), .nsec = @intCast((ms % 1000) * std.time.ns_per_ms) };
+ _ = linux.nanosleep(&ts, null);
+}
+
+// The test threads use atomic builtins rather than `std.atomic.Value` methods
+// so that, in release modes, their pc is never inside an inlined callee: the
+// DWARF symbolizer names the innermost inlined function at an address (see
+// the notes on `writeFrames`).
+const SpinState = struct {
+ tid: std.atomic.Value(u32) = .init(0),
+ stop: bool = false,
+ counter: u32 = 0,
+ done: bool = false,
+};
+
+noinline fn spinHere(st: *SpinState) void {
+ while (!@atomicLoad(bool, &st.stop, .acquire)) {
+ _ = @atomicRmw(u32, &st.counter, .Add, 1, .monotonic);
+ }
+}
+
+fn spinThreadMain(st: *SpinState) void {
+ st.tid.store(selfTid(), .release);
+ spinHere(st);
+ @atomicStore(bool, &st.done, true, .release); // keeps the call above from becoming a tail call
+}
+
+fn waitForTid(st: *SpinState) u32 {
+ var tries: usize = 0;
+ while (st.tid.load(.acquire) == 0) : (tries += 1) {
+ if (tries > 2000) return 0;
+ sleepMs(1);
+ }
+ return st.tid.load(.acquire);
+}
+
+test "capture own stack" {
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: Debug = undefined;
+ try d.init(testOptions(&text_buf));
+ defer d.deinit();
+ try testing.expect(current == &d);
+
+ var out: Writer.Allocating = .init(testing.allocator);
+ defer out.deinit();
+ try d.threadStack(selfTid(), &out.writer);
+ const text = out.written();
+ try testing.expect(std.mem.indexOf(u8, text, "#0 0x") != null);
+ try testing.expect(std.mem.indexOf(u8, text, "debug.zig:") != null);
+ try testing.expect(std.mem.indexOf(u8, text, "test.capture own stack") != null);
+
+ out.clearRetainingCapacity();
+ try d.threadRegs(selfTid(), &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "pc 0x") != null);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "pc 0x0\n") == null);
+}
+
+test "capture another thread: stack, regs, name, stat" {
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: Debug = undefined;
+ try d.init(testOptions(&text_buf));
+ defer d.deinit();
+
+ var st: SpinState = .{};
+ const th = try std.Thread.spawn(.{}, spinThreadMain, .{&st});
+ const tid = waitForTid(&st);
+ try testing.expect(tid != 0);
+
+ var out: Writer.Allocating = .init(testing.allocator);
+ defer out.deinit();
+ try d.threadStack(tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "spinHere") != null);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "spinThreadMain") != null);
+
+ out.clearRetainingCapacity();
+ try d.threadRegs(tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "pc 0x") != null);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "pc 0x0\n") == null);
+
+ out.clearRetainingCapacity();
+ try d.threadName(tid, &out.writer);
+ try testing.expect(out.written().len > 0);
+ try testing.expect(std.mem.indexOfScalar(u8, out.written(), '\n') == null);
+
+ out.clearRetainingCapacity();
+ try d.threadStat(tid, &out.writer);
+ try testing.expect(std.mem.startsWith(u8, out.written(), "state "));
+ try testing.expect(std.mem.indexOf(u8, out.written(), "\nutime ") != null);
+
+ // Enumeration lists both threads and nothing bogus.
+ try testing.expect(d.threadExists(tid));
+ try testing.expect(d.threadExists(selfTid()));
+ var found_self = false;
+ var found_other = false;
+ var i: usize = 0;
+ var prev: u32 = 0;
+ while (d.threadAt(i)) |t| : (i += 1) {
+ try testing.expect(t > prev);
+ prev = t;
+ if (t == tid) found_other = true;
+ if (t == selfTid()) found_self = true;
+ }
+ try testing.expect(found_self and found_other);
+
+ // Repeated captures of the same thread keep working.
+ var k: usize = 0;
+ while (k < 5) : (k += 1) {
+ out.clearRetainingCapacity();
+ try d.threadStack(tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "spinHere") != null);
+ }
+ const before = @atomicLoad(u32, &st.counter, .acquire);
+ sleepMs(2);
+ try testing.expect(@atomicLoad(u32, &st.counter, .acquire) != before); // the thread is running again
+
+ @atomicStore(bool, &st.stop, true, .release);
+ th.join();
+ try testing.expect(!d.threadExists(tid));
+ try testing.expectError(error.NoThread, d.threadStack(tid, &out.writer));
+ try testing.expectError(error.NoThread, d.threadName(tid, &out.writer));
+}
+
+/// The address of the call site in the caller, i.e. inside this file's test.
+noinline fn callerAddress() usize {
+ return @returnAddress() - 1;
+}
+
+test "resolveAddr names this file" {
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: Debug = undefined;
+ try d.init(testOptions(&text_buf));
+ defer d.deinit();
+ var out: Writer.Allocating = .init(testing.allocator);
+ defer out.deinit();
+ try d.resolveAddr(callerAddress(), &out.writer);
+ const text = out.written();
+ var lines = std.mem.splitScalar(u8, text, '\n');
+ const fn_name = lines.next().?;
+ const loc = lines.next().?;
+ const module = lines.next().?;
+ try testing.expect(fn_name.len > 0 and !std.mem.eql(u8, fn_name, "?"));
+ try testing.expect(std.mem.indexOf(u8, loc, "debug.zig:") != null);
+ try testing.expect(module.len > 0);
+
+ out.clearRetainingCapacity();
+ try d.resolveAddr(8, &out.writer);
+ try testing.expectEqualStrings("?\n?\n?\n", out.written());
+
+ // Regression: an unmapped lookup must not poison std's unwind cache (see
+ // the header); unwinding afterwards still works.
+ out.clearRetainingCapacity();
+ try d.threadStack(selfTid(), &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "test.resolveAddr names this file") != null);
+}
+
+test "readMem, writeMem, hexdump" {
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: Debug = undefined;
+ try d.init(testOptions(&text_buf));
+ defer d.deinit();
+
+ var value: [8]u8 = .{ 1, 2, 3, 4, 5, 6, 7, 8 };
+ var got: [8]u8 = undefined;
+ try testing.expectEqual(@as(usize, 8), try d.readMem(@intFromPtr(&value), &got));
+ try testing.expectEqualSlices(u8, &value, &got);
+ try testing.expectError(error.Unmapped, d.readMem(8, &got));
+
+ const new = [_]u8{ 0xaa, 0xbb, 0xcc };
+ try testing.expectEqual(@as(usize, 3), try d.writeMem(@intFromPtr(&value) + 2, &new));
+ try testing.expectEqualSlices(u8, &.{ 1, 2, 0xaa, 0xbb, 0xcc, 6, 7, 8 }, &value);
+ try testing.expectError(error.Unmapped, d.writeMem(8, &new));
+
+ var bytes: [19]u8 = .{ 0x00, 0x11, 0x22, 0x33, 0x44, 0x55, 0x66, 0x77, 0x88, 0x99, 0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff, 0x01, 0x12, 0x13 };
+ var out: Writer.Allocating = .init(testing.allocator);
+ defer out.deinit();
+ try d.hexdump(@intFromPtr(&bytes), bytes.len, &out.writer);
+ const expected = try std.fmt.allocPrint(testing.allocator,
+ \\{x:0>[2]} 00 11 22 33 44 55 66 77 88 99 AA BB CC DD EE FF .."3DUfw........
+ \\{x:0>[2]} 01 12 13 ...
+ \\
+ , .{ @intFromPtr(&bytes), @intFromPtr(&bytes) + 16, @sizeOf(usize) * 2 });
+ defer testing.allocator.free(expected);
+ try testing.expectEqualStrings(expected, out.written());
+ try testing.expectError(error.Unmapped, d.hexdump(8, 16, &out.writer));
+
+ // Address 0 (also reached by an offset that wraps) must be an error, not a
+ // safety-checked null pointer cast on the server thread.
+ try testing.expectError(error.Unmapped, d.readMem(0, &got));
+ try testing.expectError(error.Unmapped, d.writeMem(0, &new));
+ try testing.expectError(error.Unmapped, d.hexdump(0, 16, &out.writer));
+ try testing.expectError(error.Unmapped, d.readMem(std.math.maxInt(usize) - 3, &got));
+ try testing.expectError(error.Unmapped, d.hexdump(std.math.maxInt(usize) - 3, 16, &out.writer));
+
+ out.clearRetainingCapacity();
+ try d.maps(&out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "[stack]") != null);
+
+ // readMaps serves the file piecewise at any offset and ends with 0.
+ var piece: [4096]u8 = undefined;
+ var total: usize = 0;
+ while (true) {
+ const n = try d.readMaps(total, &piece);
+ if (n == 0) break;
+ total += n;
+ }
+ try testing.expect(total >= out.written().len / 2);
+ try testing.expectEqual(@as(usize, 0), try d.readMaps(std.math.maxInt(u64), &piece));
+}
+
+test "breakpoint on the server thread and past the slot table steps over; tgkill SIGTRAP parks" {
+ if (arch != .x86_64 and !arch.isAARCH64()) return error.SkipZigTest;
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: Debug = undefined;
+ var opts = testOptions(&text_buf);
+ opts.max_paused = 1;
+ try d.init(opts);
+ defer d.deinit();
+ try d.enableBreakpoints();
+ defer d.disableBreakpoints();
+ var out: Writer.Allocating = .init(testing.allocator);
+ defer out.deinit();
+
+ // The "server" thread (this one, for the test) hits a breakpoint: it keeps running.
+ const skipped0 = traps_skipped.load(.acquire);
+ server_tid.store(selfTid(), .release);
+ defer server_tid.store(0, .release);
+ @breakpoint();
+ try testing.expectEqual(skipped0 + 1, traps_skipped.load(.acquire));
+ try testing.expect(!d.isPaused(selfTid()));
+
+ // One slot: the first trapping thread parks, the second steps over.
+ var a: TrapState = .{};
+ const ta = try std.Thread.spawn(.{}, trapThreadMain, .{&a});
+ var tries: usize = 0;
+ while (a.tid.load(.acquire) == 0 or !d.isPaused(a.tid.load(.acquire))) : (tries += 1) {
+ try testing.expect(tries < 5000);
+ sleepMs(1);
+ }
+ var b: TrapState = .{};
+ const tb = try std.Thread.spawn(.{}, trapThreadMain, .{&b});
+ tb.join();
+ try testing.expectEqual(@as(u32, 1), @atomicLoad(u32, &b.counter, .acquire));
+ try testing.expectEqual(skipped0 + 2, traps_skipped.load(.acquire));
+ try testing.expectEqual(@as(u32, 0), @atomicLoad(u32, &a.counter, .acquire));
+ try d.resumeThread(a.tid.load(.acquire));
+ ta.join();
+ try testing.expectEqual(@as(u32, 1), @atomicLoad(u32, &a.counter, .acquire));
+
+ // A SIGTRAP sent with tgkill (not an int3/brk) parks the thread where it
+ // was; resuming it must not skip an instruction: the spinner keeps counting.
+ var st: SpinState = .{};
+ const th = try std.Thread.spawn(.{}, spinThreadMain, .{&st});
+ const tid = waitForTid(&st);
+ try testing.expect(tid != 0);
+ try testing.expectEqual(linux.E.SUCCESS, linux.errno(linux.tgkill(linux.getpid(), @intCast(tid), .TRAP)));
+ tries = 0;
+ while (!d.isPaused(tid)) : (tries += 1) {
+ try testing.expect(tries < 5000);
+ sleepMs(1);
+ }
+ const frozen = @atomicLoad(u32, &st.counter, .acquire);
+ sleepMs(5);
+ try testing.expectEqual(frozen, @atomicLoad(u32, &st.counter, .acquire));
+ out.clearRetainingCapacity();
+ try d.pausedStack(tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "spinHere") != null);
+ try d.resumeThread(tid);
+ sleepMs(5);
+ try testing.expect(@atomicLoad(u32, &st.counter, .acquire) != frozen);
+ @atomicStore(bool, &st.stop, true, .release);
+ th.join();
+}
+
+const LockState = struct {
+ tid: std.atomic.Value(u32) = .init(0),
+ release: std.atomic.Value(bool) = .init(false),
+ unlocked: std.atomic.Value(bool) = .init(false),
+ stop: std.atomic.Value(bool) = .init(false),
+ io: std.Io,
+};
+
+fn lockHolderMain(st: *LockState) void {
+ const di = std.debug.getSelfDebugInfo() catch return;
+ di.rwlock.lockUncancelable(st.io);
+ st.tid.store(selfTid(), .release);
+ while (!st.release.load(.acquire)) sleepMs(1);
+ di.rwlock.unlock(st.io);
+ st.unlocked.store(true, .release);
+ while (!st.stop.load(.acquire)) sleepMs(1);
+}
+
+test "a target parked while holding std.debug's lock is Busy, not a deadlock" {
+ if (comptime !@hasField(std.debug.SelfInfo, "rwlock")) return error.SkipZigTest;
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: Debug = undefined;
+ try d.init(testOptions(&text_buf));
+ defer d.deinit();
+ var st: LockState = .{ .io = testing.io };
+ const th = try std.Thread.spawn(.{}, lockHolderMain, .{&st});
+ var tries: usize = 0;
+ while (st.tid.load(.acquire) == 0) : (tries += 1) {
+ try testing.expect(tries < 2000);
+ sleepMs(1);
+ }
+ const tid = st.tid.load(.acquire);
+ // No allocation while the holder has the lock: `testing.allocator`
+ // records a stack trace per allocation, which needs that same lock.
+ var buf: [16 * 1024]u8 = undefined;
+ var w: Writer = .fixed(&buf);
+ try testing.expectError(error.Busy, d.threadStack(tid, &w));
+ try testing.expectEqual(cap_idle, d.capture.state.load(.acquire));
+ // Registers need no unwind and are still available.
+ try d.threadRegs(tid, &w);
+ try testing.expect(std.mem.indexOf(u8, w.buffered(), "pc 0x") != null);
+ // Handshake, not a sleep: a slow holder would otherwise still hold the
+ // lock and the next capture would legitimately be Busy again.
+ st.release.store(true, .release);
+ tries = 0;
+ while (!st.unlocked.load(.acquire)) : (tries += 1) {
+ try testing.expect(tries < 5000);
+ sleepMs(1);
+ }
+ w = .fixed(&buf);
+ try d.threadStack(tid, &w);
+ try testing.expect(std.mem.indexOf(u8, w.buffered(), "lockHolderMain") != null);
+ st.stop.store(true, .release);
+ th.join();
+}
+
+const TrapState = struct {
+ tid: std.atomic.Value(u32) = .init(0),
+ counter: u32 = 0,
+};
+
+noinline fn trapThreadMain(st: *TrapState) void {
+ st.tid.store(selfTid(), .release);
+ @breakpoint();
+ _ = @atomicRmw(u32, &st.counter, .Add, 1, .acq_rel);
+}
+
+test "breakpoint: pause, inspect, resume" {
+ if (arch != .x86_64 and !arch.isAARCH64()) return error.SkipZigTest;
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: Debug = undefined;
+ try d.init(testOptions(&text_buf));
+ defer d.deinit();
+ try d.enableBreakpoints();
+
+ var st: TrapState = .{};
+ const th = try std.Thread.spawn(.{}, trapThreadMain, .{&st});
+ var tries: usize = 0;
+ while (st.tid.load(.acquire) == 0 or !d.isPaused(st.tid.load(.acquire))) : (tries += 1) {
+ try testing.expect(tries < 5000);
+ sleepMs(1);
+ }
+ const tid = st.tid.load(.acquire);
+ try testing.expectEqual(@as(?u32, tid), d.pausedAt(0));
+ try testing.expectEqual(@as(?u32, null), d.pausedAt(1));
+ try testing.expectEqual(@as(u32, 0), @atomicLoad(u32, &st.counter, .acquire));
+
+ var out: Writer.Allocating = .init(testing.allocator);
+ defer out.deinit();
+ try d.pausedStack(tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "trapThreadMain") != null);
+ out.clearRetainingCapacity();
+ try d.pausedRegs(tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "pc 0x") != null);
+
+ // A paused thread can also be captured through the signal path.
+ out.clearRetainingCapacity();
+ try d.threadStack(tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "#0 0x") != null);
+
+ sleepMs(5);
+ try testing.expectEqual(@as(u32, 0), @atomicLoad(u32, &st.counter, .acquire));
+ try d.resumeThread(tid);
+ th.join();
+ try testing.expectEqual(@as(u32, 1), @atomicLoad(u32, &st.counter, .acquire));
+ try testing.expect(!d.isPaused(tid));
+ try testing.expectEqual(@as(?u32, null), d.pausedAt(0));
+ try testing.expectError(error.NotPaused, d.resumeThread(tid));
+ try testing.expectError(error.NotPaused, d.pausedStack(tid, &out.writer));
+ d.disableBreakpoints();
+}
+
+/// Stands in for `FullPanic`'s call: the first trace address is the return
+/// address into the panicking function.
+noinline fn panicLike(msg: []const u8) bool {
+ return recordPanic(msg, @returnAddress());
+}
+
+test "panic record path" {
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: Debug = undefined;
+ try d.init(testOptions(&text_buf));
+ defer d.deinit();
+ defer resetPanicRecord();
+
+ var out: Writer.Allocating = .init(testing.allocator);
+ defer out.deinit();
+ try d.panicMessage(&out.writer);
+ try testing.expectEqualStrings("", out.written());
+ try testing.expect(!d.panicHeld());
+ try testing.expectError(error.NoPanic, d.panicContinue());
+
+ try testing.expect(panicLike("something broke"));
+ try testing.expect(!recordPanic("nested", null));
+ try testing.expectEqual(selfTid(), panic_tid);
+
+ try d.panicMessage(&out.writer);
+ try testing.expectEqualStrings("something broke", out.written());
+ out.clearRetainingCapacity();
+ try d.panicStack(&out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "#0 0x") != null);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "test.panic record path") != null);
+ try testing.expect(!d.panicHeld());
+ try testing.expectError(error.NoPanic, d.panicContinue());
+
+ // A long message is truncated, not overflowed.
+ resetPanicRecord();
+ const long = [_]u8{'x'} ** (panic_msg_cap + 100);
+ try testing.expect(recordPanic(&long, null));
+ out.clearRetainingCapacity();
+ try d.panicMessage(&out.writer);
+ try testing.expectEqual(@as(usize, panic_msg_cap), out.written().len);
+}
+
+const MaskState = struct {
+ tid: std.atomic.Value(u32) = .init(0),
+ unblock: std.atomic.Value(bool) = .init(false),
+ stop: std.atomic.Value(bool) = .init(false),
+ signal: linux.SIG,
+};
+
+fn maskedThreadMain(st: *MaskState) void {
+ var set = linux.sigemptyset();
+ linux.sigaddset(&set, st.signal);
+ _ = linux.sigprocmask(linux.SIG.BLOCK, &set, null);
+ st.tid.store(selfTid(), .release);
+ while (!st.unblock.load(.acquire)) sleepMs(1);
+ _ = linux.sigprocmask(linux.SIG.UNBLOCK, &set, null);
+ while (!st.stop.load(.acquire)) sleepMs(1);
+}
+
+test "capture timeout on a thread with the signal masked" {
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: Debug = undefined;
+ var opts = testOptions(&text_buf);
+ opts.capture_timeout_ns = 50 * std.time.ns_per_ms;
+ try d.init(opts);
+ defer d.deinit();
+
+ var st: MaskState = .{ .signal = d.capture_signal };
+ const th = try std.Thread.spawn(.{}, maskedThreadMain, .{&st});
+ var tries: usize = 0;
+ while (st.tid.load(.acquire) == 0) : (tries += 1) {
+ try testing.expect(tries < 2000);
+ sleepMs(1);
+ }
+ const masked_tid = st.tid.load(.acquire);
+
+ var out: Writer.Allocating = .init(testing.allocator);
+ defer out.deinit();
+ const t0 = monotonicNs();
+ try testing.expectError(error.Timeout, d.threadStack(masked_tid, &out.writer));
+ try testing.expect(monotonicNs() - t0 >= 50 * std.time.ns_per_ms);
+ try testing.expectEqual(cap_idle, d.capture.state.load(.acquire));
+
+ // The process is healthy: another thread can still be captured...
+ var spin: SpinState = .{};
+ const spinner = try std.Thread.spawn(.{}, spinThreadMain, .{&spin});
+ const spin_tid = waitForTid(&spin);
+ try testing.expect(spin_tid != 0);
+ out.clearRetainingCapacity();
+ try d.threadStack(spin_tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "spinHere") != null);
+
+ // ...and the late delivery of the pending signal is harmless.
+ st.unblock.store(true, .release);
+ sleepMs(20);
+ out.clearRetainingCapacity();
+ try d.threadStack(spin_tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "spinHere") != null);
+ out.clearRetainingCapacity();
+ try d.threadStack(masked_tid, &out.writer);
+ try testing.expect(std.mem.indexOf(u8, out.written(), "maskedThreadMain") != null);
+
+ @atomicStore(bool, &spin.stop, true, .release);
+ spinner.join();
+ st.stop.store(true, .release);
+ th.join();
+}
+
+test "options validation and single instance" {
+ var text_buf: [4096]u8 = undefined;
+ var d: Debug = undefined;
+ var opts = testOptions(&text_buf);
+ opts.capture_signal = 5;
+ try testing.expectError(error.InvalidOptions, d.init(opts));
+ opts = testOptions(&text_buf);
+ opts.max_paused = max_paused_cap + 1;
+ try testing.expectError(error.InvalidOptions, d.init(opts));
+ try d.init(testOptions(&text_buf));
+ defer d.deinit();
+ var d2: Debug = undefined;
+ try testing.expectError(error.AlreadyInitialized, d2.init(testOptions(&text_buf)));
+}
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);
+}
diff --git a/9proc/src/linux/provider.zig b/9proc/src/linux/provider.zig
new file mode 100644
index 0000000..9952398
--- /dev/null
+++ b/9proc/src/linux/provider.zig
@@ -0,0 +1,604 @@
+//! Adapts `debug.Debug` into core `Provider`s. The core mounts providers at
+//! top level only, so one `DebugProvider` registers six of them, all sharing
+//! the same `Debug` and the same snapshot pool:
+//!
+//! /threads/<tid>/{name,stat,stack,regs} (lists only live tids)
+//! /addr/<hex> "fn\nfile:line:col\nmodule\n"
+//! /mem/maps, /mem/<hex> /proc/self/maps; raw bytes at address+offset (writable)
+//! /hex/<hex> hexdump of 256 bytes at the address
+//! /breakpoints/<tid>/{stack,regs,ctl} ctl accepts "continue" (lists only paused tids)
+//! /panic/{message,stack,ctl} ctl accepts "continue"
+//!
+//! /addr, /mem, /hex and /panic carry a README; /threads and /breakpoints
+//! list nothing but tids so that a shell glob over them sees only threads.
+//!
+//! Handles encode `(kind, tid-or-address)` in 56 bits (the core keeps the
+//! low 56 bits of a handle for the qid path): kind in bits 48..55, value in
+//! bits 0..47. Handles carry no reference count, so `clunk` is a no-op.
+//!
+//! The core hands providers a buffer-based `read`, not a writer, so every
+//! text file is generated at `open` into one of `snapshot_slots` fixed slots
+//! (keyed by handle, reference counted across fids) and served from there; a
+//! read at offset 0 regenerates, like the core's own dynamic files. `/mem/<hex>`
+//! is read and written directly at address+offset and never snapshotted, and
+//! `/mem/maps` is read straight from /proc/self/maps at the requested offset
+//! (a big process has more mappings than a snapshot slot holds).
+const std = @import("std");
+const cloud9 = @import("cloud9");
+const core = @import("../core.zig");
+const debug = @import("debug.zig");
+const Writer = std.Io.Writer;
+const Provider = core.Provider;
+const Handle = Provider.Handle;
+const Error = Provider.Error;
+const NodeStat = core.NodeStat;
+
+pub const snapshot_slots = 8;
+pub const snapshot_bytes = 32 * 1024;
+/// Bytes shown by /hex/<hex>.
+pub const hex_bytes = 256;
+
+pub const Tree = enum(u8) { threads, addr, mem, hex, breakpoints, panic };
+pub const tree_names = [_][]const u8{ "threads", "addr", "mem", "hex", "breakpoints", "panic" };
+
+const Kind = enum(u8) {
+ root = 0,
+ readme,
+ thread_dir,
+ thread_name,
+ thread_stat,
+ thread_stack,
+ thread_regs,
+ addr_file,
+ maps,
+ mem_file,
+ hex_file,
+ bp_dir,
+ bp_stack,
+ bp_regs,
+ bp_ctl,
+ panic_message,
+ panic_stack,
+ panic_ctl,
+
+ fn isDir(k: Kind) bool {
+ return k == .root or k == .thread_dir or k == .bp_dir;
+ }
+
+ /// Text files generated into a snapshot slot at open.
+ fn isText(k: Kind) bool {
+ return switch (k) {
+ .readme, .thread_name, .thread_stat, .thread_stack, .thread_regs, .addr_file, .hex_file, .bp_stack, .bp_regs, .panic_message, .panic_stack => true,
+ else => false,
+ };
+ }
+
+ fn isCtl(k: Kind) bool {
+ return k == .bp_ctl or k == .panic_ctl;
+ }
+
+ fn mode(k: Kind) u32 {
+ if (k.isDir()) return cloud9.dmdir | 0o555;
+ if (k.isCtl()) return 0o222;
+ if (k == .mem_file) return 0o666;
+ return 0o444;
+ }
+
+ fn fixedName(k: Kind) ?[]const u8 {
+ return switch (k) {
+ .readme => "README",
+ .thread_name => "name",
+ .thread_stat => "stat",
+ .thread_stack, .bp_stack, .panic_stack => "stack",
+ .thread_regs, .bp_regs => "regs",
+ .maps => "maps",
+ .bp_ctl, .panic_ctl => "ctl",
+ .panic_message => "message",
+ else => null,
+ };
+ }
+};
+
+const value_bits = 48;
+const value_mask: u64 = (1 << value_bits) - 1;
+
+fn mk(kind: Kind, value: u64) Handle {
+ return (@as(u64, @intFromEnum(kind)) << value_bits) | (value & value_mask);
+}
+
+fn kindOf(h: Handle) Kind {
+ return @enumFromInt(@as(u8, @truncate(h >> value_bits)));
+}
+
+fn valueOf(h: Handle) u64 {
+ return h & value_mask;
+}
+
+const readme_threads =
+ \\One directory per thread of this process, named by tid:
+ \\ name the thread's comm
+ \\ stat state and a few fields of /proc/self/task/<tid>/stat
+ \\ stack "#n 0x<addr> in <fn> (<file>:<line>:<col>)" per frame
+ \\ regs general registers captured while the thread was stopped
+ \\
+;
+const readme_addr =
+ \\Walk any hex address: /addr/<hex> reads as "fn\nfile:line:col\nmodule\n".
+ \\
+;
+const readme_mem =
+ \\maps /proc/self/maps
+ \\<hex> raw process memory at that address (+ file offset); writable
+ \\
+;
+const readme_hex =
+ \\Walk any hex address: /hex/<hex> is a hexdump of the 256 bytes there.
+ \\
+;
+const readme_breakpoints =
+ \\One directory per thread stopped in @breakpoint(), named by tid:
+ \\ stack, regs as under /threads
+ \\ ctl write "continue" to resume the thread
+ \\
+;
+const readme_panic =
+ \\message the first panic's message (empty before any panic)
+ \\stack frames of the panicking thread
+ \\ctl write "continue" to let the default panic handler run
+ \\
+;
+
+fn readmeFor(tree: Tree) []const u8 {
+ return switch (tree) {
+ .threads => readme_threads,
+ .addr => readme_addr,
+ .mem => readme_mem,
+ .hex => readme_hex,
+ .breakpoints => readme_breakpoints,
+ .panic => readme_panic,
+ };
+}
+
+const Slot = struct {
+ handle: Handle = 0,
+ refs: u32 = 0,
+ len: u32 = 0,
+ buf: [snapshot_bytes]u8 = undefined,
+};
+
+pub const DebugProvider = struct {
+ d: *debug.Debug,
+ slots: [snapshot_slots]Slot = @splat(.{}),
+ /// Backs `NodeStat.name` until the next call.
+ name_buf: [32]u8 = undefined,
+
+ pub fn init(dp: *DebugProvider, d: *debug.Debug) void {
+ dp.* = .{ .d = d };
+ }
+
+ /// The provider for one tree, to pass to `Shared.addProvider`.
+ pub fn provider(dp: *DebugProvider, comptime tree: Tree) Provider {
+ return .{ .name = tree_names[@intFromEnum(tree)], .ctx = dp, .vtable = vtableFor(tree) };
+ }
+
+ /// Mounts all six trees; `shared` is a `Server(cfg).Shared`.
+ pub fn mountAll(dp: *DebugProvider, shared: anytype) error{Full}!void {
+ inline for (comptime std.meta.tags(Tree)) |tree| try shared.addProvider(dp.provider(tree));
+ }
+
+ fn self(ctx: *anyopaque) *DebugProvider {
+ return @ptrCast(@alignCast(ctx));
+ }
+
+ fn vtableFor(comptime tree: Tree) *const Provider.VTable {
+ return &struct {
+ const vt: Provider.VTable = .{
+ .walk = walkFn,
+ .stat = statFn,
+ .list = listFn,
+ .open = openFn,
+ .read = readFn,
+ .write = writeFn,
+ .close = closeFn,
+ .clunk = clunkFn,
+ };
+ fn walkFn(ctx: *anyopaque, parent: Handle, name: []const u8) Error!Handle {
+ return self(ctx).walk(tree, parent, name);
+ }
+ fn statFn(ctx: *anyopaque, h: Handle, out: *NodeStat) Error!void {
+ return self(ctx).stat(tree, h, out);
+ }
+ fn listFn(ctx: *anyopaque, dir: Handle, index: usize, out: *NodeStat) Error!bool {
+ return self(ctx).list(tree, dir, index, out);
+ }
+ fn openFn(ctx: *anyopaque, h: Handle, mode: u8) Error!void {
+ return self(ctx).open(tree, h, mode);
+ }
+ fn readFn(ctx: *anyopaque, h: Handle, offset: u64, buf: []u8) Error!usize {
+ return self(ctx).read(tree, h, offset, buf);
+ }
+ fn writeFn(ctx: *anyopaque, h: Handle, offset: u64, data: []const u8) Error!usize {
+ return self(ctx).write(tree, h, offset, data);
+ }
+ fn closeFn(ctx: *anyopaque, h: Handle) void {
+ self(ctx).close(h);
+ }
+ fn clunkFn(_: *anyopaque, _: Handle) void {}
+ }.vt;
+ }
+
+ // -- naming ------------------------------------------------------------
+
+ fn parseTid(name: []const u8) ?u32 {
+ if (name.len == 0 or name.len > 10) return null;
+ for (name) |ch| if (!std.ascii.isDigit(ch)) return null;
+ return std.fmt.parseInt(u32, name, 10) catch null;
+ }
+
+ fn parseHex(name: []const u8) ?u64 {
+ const digits = if (std.mem.startsWith(u8, name, "0x")) name[2..] else name;
+ if (digits.len == 0 or digits.len > 12) return null;
+ for (digits) |ch| if (!std.ascii.isHex(ch)) return null;
+ const v = std.fmt.parseInt(u64, digits, 16) catch return null;
+ if (v > value_mask) return null;
+ return v;
+ }
+
+ fn nodeName(dp: *DebugProvider, h: Handle) []const u8 {
+ const k = kindOf(h);
+ if (k.fixedName()) |n| return n;
+ return switch (k) {
+ .root => "",
+ .thread_dir, .bp_dir => std.fmt.bufPrint(&dp.name_buf, "{d}", .{valueOf(h)}) catch unreachable,
+ .addr_file, .mem_file, .hex_file => std.fmt.bufPrint(&dp.name_buf, "{x}", .{valueOf(h)}) catch unreachable,
+ else => unreachable,
+ };
+ }
+
+ // -- vtable ------------------------------------------------------------
+
+ fn walk(dp: *DebugProvider, tree: Tree, parent: Handle, name: []const u8) Error!Handle {
+ const k = kindOf(parent);
+ if (std.mem.eql(u8, name, ".")) return parent;
+ if (!k.isDir()) return error.NotDir;
+ if (std.mem.eql(u8, name, "..")) return Provider.root;
+ switch (k) {
+ .root => {
+ if (tree != .threads and tree != .breakpoints and std.mem.eql(u8, name, "README")) return mk(.readme, 0);
+ switch (tree) {
+ .threads => {
+ const tid = parseTid(name) orelse return error.NotFound;
+ if (!dp.d.threadExists(tid)) return error.NotFound;
+ return mk(.thread_dir, tid);
+ },
+ .addr => return mk(.addr_file, parseHex(name) orelse return error.NotFound),
+ .mem => {
+ if (std.mem.eql(u8, name, "maps")) return mk(.maps, 0);
+ return mk(.mem_file, parseHex(name) orelse return error.NotFound);
+ },
+ .hex => return mk(.hex_file, parseHex(name) orelse return error.NotFound),
+ .breakpoints => {
+ const tid = parseTid(name) orelse return error.NotFound;
+ if (!dp.d.isPaused(tid)) return error.NotFound;
+ return mk(.bp_dir, tid);
+ },
+ .panic => {
+ if (std.mem.eql(u8, name, "message")) return mk(.panic_message, 0);
+ if (std.mem.eql(u8, name, "stack")) return mk(.panic_stack, 0);
+ if (std.mem.eql(u8, name, "ctl")) return mk(.panic_ctl, 0);
+ return error.NotFound;
+ },
+ }
+ },
+ .thread_dir => {
+ const tid = valueOf(parent);
+ if (std.mem.eql(u8, name, "name")) return mk(.thread_name, tid);
+ if (std.mem.eql(u8, name, "stat")) return mk(.thread_stat, tid);
+ if (std.mem.eql(u8, name, "stack")) return mk(.thread_stack, tid);
+ if (std.mem.eql(u8, name, "regs")) return mk(.thread_regs, tid);
+ return error.NotFound;
+ },
+ .bp_dir => {
+ const tid = valueOf(parent);
+ if (std.mem.eql(u8, name, "stack")) return mk(.bp_stack, tid);
+ if (std.mem.eql(u8, name, "regs")) return mk(.bp_regs, tid);
+ if (std.mem.eql(u8, name, "ctl")) return mk(.bp_ctl, tid);
+ return error.NotFound;
+ },
+ else => unreachable,
+ }
+ }
+
+ fn fill(dp: *DebugProvider, tree: Tree, h: Handle, out: *NodeStat) void {
+ const k = kindOf(h);
+ out.* = .{
+ .mode = k.mode(),
+ .length = if (k == .readme) readmeFor(tree).len else 0,
+ .name = dp.nodeName(h),
+ .handle = h,
+ };
+ }
+
+ fn stat(dp: *DebugProvider, tree: Tree, h: Handle, out: *NodeStat) Error!void {
+ dp.fill(tree, h, out);
+ }
+
+ fn list(dp: *DebugProvider, tree: Tree, dir: Handle, index: usize, out: *NodeStat) Error!bool {
+ const k = kindOf(dir);
+ if (!k.isDir()) return error.NotDir;
+ const h: Handle = switch (k) {
+ .root => switch (tree) {
+ .threads => mk(.thread_dir, dp.d.threadAt(index) orelse return false),
+ .addr, .hex => if (index == 0) mk(.readme, 0) else return false,
+ .mem => switch (index) {
+ 0 => mk(.readme, 0),
+ 1 => mk(.maps, 0),
+ else => return false,
+ },
+ .breakpoints => mk(.bp_dir, dp.d.pausedAt(index) orelse return false),
+ .panic => switch (index) {
+ 0 => mk(.readme, 0),
+ 1 => mk(.panic_message, 0),
+ 2 => mk(.panic_stack, 0),
+ 3 => mk(.panic_ctl, 0),
+ else => return false,
+ },
+ },
+ .thread_dir => switch (index) {
+ 0 => mk(.thread_name, valueOf(dir)),
+ 1 => mk(.thread_stat, valueOf(dir)),
+ 2 => mk(.thread_stack, valueOf(dir)),
+ 3 => mk(.thread_regs, valueOf(dir)),
+ else => return false,
+ },
+ .bp_dir => switch (index) {
+ 0 => mk(.bp_stack, valueOf(dir)),
+ 1 => mk(.bp_regs, valueOf(dir)),
+ 2 => mk(.bp_ctl, valueOf(dir)),
+ else => return false,
+ },
+ else => unreachable,
+ };
+ dp.fill(tree, h, out);
+ return true;
+ }
+
+ fn open(dp: *DebugProvider, tree: Tree, h: Handle, mode: u8) Error!void {
+ const k = kindOf(h);
+ const acc = mode & 3;
+ const wants_write = acc == cloud9.owrite or acc == cloud9.ordwr;
+ const wants_read = acc != cloud9.owrite;
+ if (k.isDir()) {
+ if (wants_write or mode & cloud9.otrunc != 0) return error.IsDir;
+ return;
+ }
+ if (k.isCtl()) {
+ if (wants_read) return error.Perm;
+ return;
+ }
+ if (k == .mem_file) return;
+ if (wants_write or mode & cloud9.otrunc != 0) return error.Perm;
+ if (k == .maps) return;
+ std.debug.assert(k.isText());
+ const slot = dp.takeSlot(h) orelse return error.NoSpace;
+ errdefer dp.releaseSlot(slot);
+ try dp.generate(tree, h, slot);
+ }
+
+ fn close(dp: *DebugProvider, h: Handle) void {
+ if (!kindOf(h).isText()) return;
+ if (dp.findSlot(h)) |s| dp.releaseSlot(s);
+ }
+
+ fn read(dp: *DebugProvider, tree: Tree, h: Handle, offset: u64, buf: []u8) Error!usize {
+ const k = kindOf(h);
+ if (k.isDir()) return error.IsDir;
+ if (k.isCtl()) return error.Perm;
+ if (k == .mem_file) {
+ const addr = valueOf(h) +% offset;
+ return dp.d.readMem(addr, buf) catch |e| mapErr(e);
+ }
+ if (k == .maps) return dp.d.readMaps(offset, buf) catch |e| mapErr(e);
+ const slot = dp.findSlot(h) orelse return error.Io;
+ if (offset == 0) try dp.generate(tree, h, slot);
+ if (offset >= slot.len) return 0;
+ const off: usize = @intCast(offset);
+ const n = @min(buf.len, slot.len - off);
+ @memcpy(buf[0..n], slot.buf[off..][0..n]);
+ return n;
+ }
+
+ fn write(dp: *DebugProvider, tree: Tree, h: Handle, offset: u64, data: []const u8) Error!usize {
+ _ = tree;
+ const k = kindOf(h);
+ if (k.isDir()) return error.IsDir;
+ switch (k) {
+ .mem_file => {
+ const addr = valueOf(h) +% offset;
+ return dp.d.writeMem(addr, data) catch |e| mapErr(e);
+ },
+ .bp_ctl, .panic_ctl => {
+ const cmd = std.mem.trim(u8, data, " \t\r\n\x00");
+ if (!std.mem.eql(u8, cmd, "continue")) return error.Unsupported;
+ if (k == .bp_ctl) {
+ dp.d.resumeThread(@intCast(valueOf(h))) catch |e| return mapErr(e);
+ } else {
+ dp.d.panicContinue() catch |e| return mapErr(e);
+ }
+ return data.len;
+ },
+ else => return error.Perm,
+ }
+ }
+
+ // -- snapshots ---------------------------------------------------------
+
+ fn findSlot(dp: *DebugProvider, h: Handle) ?*Slot {
+ for (&dp.slots) |*s| if (s.refs > 0 and s.handle == h) return s;
+ return null;
+ }
+
+ fn takeSlot(dp: *DebugProvider, h: Handle) ?*Slot {
+ if (dp.findSlot(h)) |s| {
+ s.refs += 1;
+ return s;
+ }
+ for (&dp.slots) |*s| if (s.refs == 0) {
+ s.* = .{ .handle = h, .refs = 1 };
+ return s;
+ };
+ return null;
+ }
+
+ fn releaseSlot(_: *DebugProvider, s: *Slot) void {
+ s.refs -= 1;
+ }
+
+ /// (Re)generates the text of `h` into `slot`. A text that does not fit is
+ /// truncated, not an error.
+ fn generate(dp: *DebugProvider, tree: Tree, h: Handle, slot: *Slot) Error!void {
+ var w: Writer = .fixed(&slot.buf);
+ slot.len = 0;
+ dp.render(tree, h, &w) catch |e| switch (e) {
+ error.WriteFailed => {},
+ else => return mapErr(e),
+ };
+ slot.len = @intCast(w.buffered().len);
+ }
+
+ fn render(dp: *DebugProvider, tree: Tree, h: Handle, w: *Writer) debug.Error!void {
+ const d = dp.d;
+ const v = valueOf(h);
+ switch (kindOf(h)) {
+ .readme => w.writeAll(readmeFor(tree)) catch return error.WriteFailed,
+ .thread_name => try d.threadName(@intCast(v), w),
+ .thread_stat => try d.threadStat(@intCast(v), w),
+ .thread_stack => try d.threadStack(@intCast(v), w),
+ .thread_regs => try d.threadRegs(@intCast(v), w),
+ .addr_file => try d.resolveAddr(@intCast(v), w),
+ .hex_file => try d.hexdump(@intCast(v), hex_bytes, w),
+ .bp_stack => try d.pausedStack(@intCast(v), w),
+ .bp_regs => try d.pausedRegs(@intCast(v), w),
+ .panic_message => try d.panicMessage(w),
+ .panic_stack => try d.panicStack(w),
+ else => unreachable,
+ }
+ }
+
+ fn mapErr(e: debug.Error) Error {
+ return switch (e) {
+ error.NoThread, error.NotPaused, error.NoPanic => error.NotFound,
+ error.Unsupported => error.Unsupported,
+ error.WriteFailed => error.NoSpace,
+ error.Timeout, error.Busy, error.Unmapped, error.Unexpected, error.AlreadyInitialized, error.InvalidOptions => error.Io,
+ };
+ }
+};
+
+// ---------------------------------------------------------------------------
+// Tests (through the core's in-memory harness)
+// ---------------------------------------------------------------------------
+
+const testing = std.testing;
+
+const TestCfg: core.Config = .{ .name = "dbgtest", .msize = 8192, .max_fids = 16, .max_providers = 6, .snapshot_slots = 2, .snapshot_bytes = 512 };
+const TS = core.Server(TestCfg);
+
+test "debug provider: threads, addr, mem, hex, breakpoints, panic through the core" {
+ var text_buf: [16 * 1024]u8 = undefined;
+ var d: debug.Debug = undefined;
+ try d.init(.{ .io = testing.io, .text_buf = &text_buf });
+ defer d.deinit();
+ var dp: DebugProvider = undefined;
+ dp.init(&d);
+
+ var dummy: u8 = 0;
+ var shared: TS.Shared = .init(&dummy);
+ try dp.mountAll(&shared);
+ var storage: TS.Storage = undefined;
+ var h: TS.Harness = undefined;
+ try h.init(&shared, &storage);
+ defer h.deinit();
+
+ // /threads lists tids only, among them this thread.
+ const names = try h.listPath(&.{"threads"});
+ defer TS.Harness.freeNames(names);
+ try testing.expect(names.len >= 1);
+ try testing.expect(!TS.Harness.hasName(names, "README"));
+ const no_readme = try h.ok(.{ .walk = .{ .fid = 0, .newfid = 7, .names = &.{ "threads", "README" } } });
+ try testing.expectEqual(@as(u16, 1), no_readme.walk.nwqid);
+ var tid_buf: [16]u8 = undefined;
+ const tid = try std.fmt.bufPrint(&tid_buf, "{d}", .{std.os.linux.gettid()});
+ try testing.expect(TS.Harness.hasName(names, tid));
+
+ // Own stack names this test function's file.
+ const stack = try h.readPath(&.{ "threads", tid, "stack" });
+ defer testing.allocator.free(stack);
+ try testing.expect(std.mem.indexOf(u8, stack, "#0 0x") != null);
+
+ // /addr/<hex> of a function here resolves to this file.
+ var addr_buf: [32]u8 = undefined;
+ const addr_name = try std.fmt.bufPrint(&addr_buf, "{x}", .{@intFromPtr(&DebugProvider.parseTid)});
+ const resolved = try h.readPath(&.{ "addr", addr_name });
+ defer testing.allocator.free(resolved);
+ try testing.expect(std.mem.indexOf(u8, resolved, "provider.zig") != null);
+ const partial = try h.ok(.{ .walk = .{ .fid = 0, .newfid = 5, .names = &.{ "addr", "zzz" } } });
+ try testing.expectEqual(@as(u16, 1), partial.walk.nwqid);
+ try h.walkTo(5, &.{"addr"});
+ try h.expectFail(.{ .walk = .{ .fid = 5, .newfid = 6, .names = &.{"zzz"} } }, "file does not exist");
+ _ = try h.ok(.{ .clunk = .{ .fid = 5 } });
+
+ // /mem/<hex> reads and writes live memory; unmapped is an error.
+ var cell: [8]u8 = "abcdefgh".*;
+ var mem_buf: [32]u8 = undefined;
+ const mem_name = try std.fmt.bufPrint(&mem_buf, "{x}", .{@intFromPtr(&cell)});
+ try h.walkTo(1, &.{ "mem", mem_name });
+ _ = try h.ok(.{ .open = .{ .fid = 1, .mode = cloud9.ordwr } });
+ const r = try h.ok(.{ .read = .{ .fid = 1, .offset = 2, .count = 4 } });
+ try testing.expectEqualStrings("cdef", r.read);
+ _ = try h.ok(.{ .write = .{ .fid = 1, .offset = 0, .data = "XY" } });
+ try testing.expectEqualStrings("XYcdefgh", &cell);
+ _ = try h.ok(.{ .clunk = .{ .fid = 1 } });
+ try h.walkTo(2, &.{ "mem", "8" });
+ _ = try h.ok(.{ .open = .{ .fid = 2, .mode = cloud9.oread } });
+ try h.expectFail(.{ .read = .{ .fid = 2, .offset = 0, .count = 4 } }, "i/o error");
+ _ = try h.ok(.{ .clunk = .{ .fid = 2 } });
+ const maps = try h.readPath(&.{ "mem", "maps" });
+ defer testing.allocator.free(maps);
+ try testing.expect(std.mem.indexOf(u8, maps, "r-xp") != null or std.mem.indexOf(u8, maps, "r--p") != null);
+
+ // /hex/<hex> is a hexdump.
+ const hex = try h.readPath(&.{ "hex", mem_name });
+ defer testing.allocator.free(hex);
+ try testing.expect(std.mem.indexOf(u8, hex, "XYcdefgh") != null);
+
+ // Nothing paused, no panic.
+ const bps = try h.listPath(&.{"breakpoints"});
+ defer TS.Harness.freeNames(bps);
+ try testing.expectEqual(@as(usize, 0), bps.len);
+ const msg = try h.readPath(&.{ "panic", "message" });
+ defer testing.allocator.free(msg);
+ try testing.expectEqualStrings("", msg);
+ try h.walkTo(3, &.{ "panic", "ctl" });
+ _ = try h.ok(.{ .open = .{ .fid = 3, .mode = cloud9.owrite } });
+ try h.expectFail(.{ .write = .{ .fid = 3, .offset = 0, .data = "continue" } }, "file does not exist");
+ try h.expectFail(.{ .write = .{ .fid = 3, .offset = 0, .data = "bogus" } }, "not supported");
+ _ = try h.ok(.{ .clunk = .{ .fid = 3 } });
+
+ // Snapshot slots are released on clunk: open more files than slots, sequentially.
+ for (0..4) |_| {
+ const t = try h.readPath(&.{ "threads", tid, "name" });
+ testing.allocator.free(t);
+ }
+ for (&dp.slots) |s| try testing.expectEqual(@as(u32, 0), s.refs);
+}
+
+test "handle encoding round-trips" {
+ const h = mk(.mem_file, 0x7fff_dead_beef);
+ try testing.expectEqual(Kind.mem_file, kindOf(h));
+ try testing.expectEqual(@as(u64, 0x7fff_dead_beef), valueOf(h));
+ try testing.expect(h < (1 << 56));
+ try testing.expectEqual(@as(?u64, null), DebugProvider.parseHex("1_0"));
+ try testing.expectEqual(@as(?u64, 0x10), DebugProvider.parseHex("0x10"));
+ try testing.expectEqual(@as(?u32, null), DebugProvider.parseTid("+5"));
+}
diff --git a/9proc/src/linux/runtime.zig b/9proc/src/linux/runtime.zig
new file mode 100644
index 0000000..503d6c2
--- /dev/null
+++ b/9proc/src/linux/runtime.zig
@@ -0,0 +1,90 @@
+//! Generators for `Config.runtime`: /runtime/{pid,ppid,uptime,argv,cwd,env,clients}.
+//! The core passes every generator `Shared.ctx`; `Fns(Ctx, field)` casts it
+//! to `*Ctx` and reads the `Info` stored in `@field(ctx, field)`.
+const std = @import("std");
+const linux = std.os.linux;
+const Writer = std.Io.Writer;
+
+/// What the generators report. Texts are borrowed for the server's lifetime.
+pub const Info = struct {
+ /// argv, one argument per line.
+ argv: []const u8 = "",
+ /// Environment, one KEY=VALUE per line.
+ env: []const u8 = "",
+ cwd: []const u8 = "",
+ /// Monotonic seconds at startup; /runtime/uptime is the difference.
+ start_mono: i64 = 0,
+ /// Live client count, published by the probe.
+ clients: ?*const std.atomic.Value(u32) = null,
+
+ pub fn now() Info {
+ return .{ .start_mono = monotonicSecs() };
+ }
+};
+
+pub fn monotonicSecs() i64 {
+ var ts: linux.timespec = undefined;
+ _ = linux.clock_gettime(.MONOTONIC, &ts);
+ return ts.sec;
+}
+
+/// Seconds since the epoch, clamped to u32 (for atime/mtime and `fn/now`).
+pub fn realtimeSecs() u32 {
+ var ts: linux.timespec = undefined;
+ _ = linux.clock_gettime(.REALTIME, &ts);
+ return @intCast(std.math.clamp(ts.sec, 0, std.math.maxInt(u32)));
+}
+
+/// The `Config.runtime` type: `Ctx` is the type behind `Shared.ctx`, `field`
+/// the name of its `Info` field.
+pub fn Fns(comptime Ctx: type, comptime field: []const u8) type {
+ return struct {
+ fn info(ctx: *anyopaque) *const Info {
+ const c: *Ctx = @ptrCast(@alignCast(ctx));
+ return &@field(c, field);
+ }
+ pub fn pid(_: *anyopaque, w: *Writer) anyerror!void {
+ try w.print("{d}", .{linux.getpid()});
+ }
+ pub fn ppid(_: *anyopaque, w: *Writer) anyerror!void {
+ try w.print("{d}", .{linux.getppid()});
+ }
+ pub fn uptime(ctx: *anyopaque, w: *Writer) anyerror!void {
+ try w.print("{d}", .{monotonicSecs() - info(ctx).start_mono});
+ }
+ pub fn argv(ctx: *anyopaque, w: *Writer) anyerror!void {
+ try w.writeAll(info(ctx).argv);
+ }
+ pub fn cwd(ctx: *anyopaque, w: *Writer) anyerror!void {
+ try w.writeAll(info(ctx).cwd);
+ }
+ pub fn env(ctx: *anyopaque, w: *Writer) anyerror!void {
+ try w.writeAll(info(ctx).env);
+ }
+ pub fn clients(ctx: *anyopaque, w: *Writer) anyerror!void {
+ const n: u32 = if (info(ctx).clients) |c| c.load(.acquire) else 0;
+ try w.print("{d}", .{n});
+ }
+ };
+}
+
+test "runtime generators read Info through the context" {
+ const Ctx = struct { x: u32, info: Info };
+ var count: std.atomic.Value(u32) = .init(3);
+ var ctx: Ctx = .{ .x = 0, .info = .{ .argv = "a\nb\n", .cwd = "/tmp", .env = "K=V\n", .start_mono = monotonicSecs(), .clients = &count } };
+ const F = Fns(Ctx, "info");
+ var buf: [64]u8 = undefined;
+ var w: Writer = .fixed(&buf);
+ try F.clients(&ctx, &w);
+ try std.testing.expectEqualStrings("3", w.buffered());
+ w = .fixed(&buf);
+ try F.argv(&ctx, &w);
+ try std.testing.expectEqualStrings("a\nb\n", w.buffered());
+ w = .fixed(&buf);
+ try F.uptime(&ctx, &w);
+ try std.testing.expect(w.buffered().len >= 1);
+ w = .fixed(&buf);
+ try F.pid(&ctx, &w);
+ try std.testing.expectEqual(linux.getpid(), try std.fmt.parseInt(i32, w.buffered(), 10));
+ try std.testing.expectEqual(@as(usize, 7), @typeInfo(F).@"struct".decls.len);
+}