diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-19 21:26:05 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-09-19 21:26:05 -0300 |
| commit | b05abcba3ea09ea106ad28364c6e40a3ec31b890 (patch) | |
| tree | 9170fac5e7e5d8bde108de34a182aaa9d6844117 /introspect/src/linux | |
| parent | ae310a207534b33b7321dd2b9f423a73b1969159 (diff) | |
| download | cloud9-b05abcba3ea09ea106ad28364c6e40a3ec31b890.tar.gz cloud9-b05abcba3ea09ea106ad28364c6e40a3ec31b890.zip | |
Add 9player and introspect as programs beside the library
9player/: FUSE mount CLI that mounts a 9P2000 tree into a fresh user+mount
namespace and runs a program in it (no root, no libfuse, no libc).
introspect/: the 9P debug/introspection library (freestanding core, value
renderers, Linux probe with threads/stacks/memory/breakpoints/panics) and
its demo server. Each has its own build fragment; the root build.zig wires
them behind -D9player/-Dintrospect with namespaced steps (9player-itest,
introspect-check-freestanding, programs-test, ...) and exports the
introspect module for dependents. This is the layout for related programs.
Co-Authored-By: Claude Fable 5.1 <[email protected]>
Diffstat (limited to 'introspect/src/linux')
| -rw-r--r-- | introspect/src/linux/debug.zig | 1458 | ||||
| -rw-r--r-- | introspect/src/linux/probe.zig | 829 | ||||
| -rw-r--r-- | introspect/src/linux/provider.zig | 604 | ||||
| -rw-r--r-- | introspect/src/linux/runtime.zig | 90 |
4 files changed, 2981 insertions, 0 deletions
diff --git a/introspect/src/linux/debug.zig b/introspect/src/linux/debug.zig new file mode 100644 index 0000000..4d43d5b --- /dev/null +++ b/introspect/src/linux/debug.zig @@ -0,0 +1,1458 @@ +//! Linux debug facilities for the introspect 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/introspect/src/linux/probe.zig b/introspect/src/linux/probe.zig new file mode 100644 index 0000000..f38c491 --- /dev/null +++ b/introspect/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("introspect"); + 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("introspect: 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/introspect-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/introspect/src/linux/provider.zig b/introspect/src/linux/provider.zig new file mode 100644 index 0000000..9952398 --- /dev/null +++ b/introspect/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/introspect/src/linux/runtime.zig b/introspect/src/linux/runtime.zig new file mode 100644 index 0000000..503d6c2 --- /dev/null +++ b/introspect/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); +} |
