From 251894cf3375812564ecf28392179b395cdda9c7 Mon Sep 17 00:00:00 2001 From: Paul Buetow Date: Wed, 13 May 2026 20:04:48 +0300 Subject: refactor: break down functions exceeding 50 lines into smaller helpers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Split 22 production files across the codebase — event loop, TUI models, probe manager, dashboard, export, flag parsing, code generation, and ioworkload scenarios — so that no function body exceeds 50 lines. Each extracted helper carries its own comment explaining its role. Co-Authored-By: Claude Sonnet 4.6 --- internal/benchutil/eventmix.go | 139 +++++++++---- internal/eventloop_exit.go | 38 +++- internal/eventloop_runtime.go | 33 +++ internal/export/snapshot_csv.go | 114 ++++++----- internal/flags/flags.go | 57 ++++-- internal/generate/bpfhandler.go | 120 ++++++----- internal/generate/format.go | 84 ++++---- internal/ior.go | 101 ++++++---- internal/ior_parquet_sink.go | 107 ++++++---- internal/ior_profiling.go | 81 +++++--- internal/probemanager/manager.go | 154 +++++++++----- internal/tui/common/keys.go | 60 +++--- internal/tui/dashboard/bubbles.go | 121 ++++++----- internal/tui/dashboard/histogram.go | 63 +++--- internal/tui/dashboard/icicle.go | 88 +++++--- internal/tui/dashboard/treemap.go | 91 ++++----- internal/tui/eventstream/export.go | 129 +++++++----- internal/tui/eventstream/model.go | 390 ++++++++++++++++++++---------------- internal/tui/export/model.go | 69 ++++--- internal/tui/probes/model.go | 172 +++++++++------- internal/tui/tracefilter/model.go | 233 +++++++++++---------- 21 files changed, 1453 insertions(+), 991 deletions(-) (limited to 'internal') diff --git a/internal/benchutil/eventmix.go b/internal/benchutil/eventmix.go index 914bb93..7961f2e 100644 --- a/internal/benchutil/eventmix.go +++ b/internal/benchutil/eventmix.go @@ -184,63 +184,114 @@ func sumWeights(entries []MixEntry) int { return total } +// pair generates a matched enter/exit raw event pair for the mix event type. +// fd/seq provide per-call identity for fd-based and path-based events. func (e MixEvent) pair(gen EventGenerator, time uint64, pid, tid uint32, fd int32, seq int) ([]byte, []byte, error) { switch e { - case MixRead: - return gen.FdPair(time, pid, tid, fd, types.SYS_ENTER_READ, types.SYS_EXIT_READ, 128) - case MixWrite: - return gen.FdPair(time, pid, tid, fd, types.SYS_ENTER_WRITE, types.SYS_EXIT_WRITE, 256) + case MixRead, MixWrite, MixFsync: + return e.pairFd(gen, time, pid, tid, fd) case MixOpen: return gen.OpenPair(time, pid, tid) case MixClose: - _, enter, err := gen.EnterFdEvent(time, pid, tid, fd, types.SYS_ENTER_CLOSE) - if err != nil { - return nil, nil, err - } - _, exit, err := gen.ExitFdEvent(time+gen.pairDelta(), pid, tid, fd, types.SYS_EXIT_CLOSE) - if err != nil { - return nil, nil, err - } - return enter, exit, nil - case MixStat: - path := fmt.Sprintf("/tmp/ior-stat-%d-%d", tid, seq) - return gen.PathPair(time, pid, tid, path, types.SYS_ENTER_NEWSTAT, types.SYS_EXIT_NEWSTAT, 0) + return pairClose(gen, time, pid, tid, fd) + case MixStat, MixAccess, MixMkdir, MixUnlink: + return e.pairPath(gen, time, pid, tid, tid, seq) case MixSync: return gen.NullPair(time, pid, tid, types.SYS_ENTER_SYNC, types.SYS_EXIT_SYNC) - case MixFsync: - return gen.FdPair(time, pid, tid, fd, types.SYS_ENTER_FSYNC, types.SYS_EXIT_FSYNC, 0) - case MixAccess: - path := fmt.Sprintf("/tmp/ior-access-%d-%d", tid, seq) - return gen.PathPair(time, pid, tid, path, types.SYS_ENTER_ACCESS, types.SYS_EXIT_ACCESS, 0) - case MixMkdir: - path := fmt.Sprintf("/tmp/ior-mkdir-%d-%d", tid, seq) - return gen.PathPair(time, pid, tid, path, types.SYS_ENTER_MKDIR, types.SYS_EXIT_MKDIR, 0) - case MixUnlink: - path := fmt.Sprintf("/tmp/ior-unlink-%d-%d", tid, seq) - return gen.PathPair(time, pid, tid, path, types.SYS_ENTER_UNLINK, types.SYS_EXIT_UNLINK, 0) - case MixRename: - oldname := fmt.Sprintf("/tmp/ior-old-%d-%d", tid, seq) - newname := fmt.Sprintf("/tmp/ior-new-%d-%d", tid, seq) - return gen.NamePair(time, pid, tid, oldname, newname, types.SYS_ENTER_RENAME, types.SYS_EXIT_RENAME, 0) - case MixLink: - oldname := fmt.Sprintf("/tmp/ior-link-old-%d-%d", tid, seq) - newname := fmt.Sprintf("/tmp/ior-link-new-%d-%d", tid, seq) - return gen.NamePair(time, pid, tid, oldname, newname, types.SYS_ENTER_LINK, types.SYS_EXIT_LINK, 0) + case MixRename, MixLink: + return e.pairName(gen, time, pid, tid, tid, seq) case MixFcntl: return gen.FcntlPair(time, pid, tid, uint32(fd), syscall.F_SETFL, syscall.O_NONBLOCK, types.SYS_EXIT_FCNTL, 0) case MixDup3: return gen.Dup3Pair(time, pid, tid, fd, syscall.O_CLOEXEC, types.SYS_EXIT_DUP3, int64(fd+1)) case MixOpenByHandleAt: - _, enter, err := gen.EnterOpenByHandleAtEvent(time, pid, tid, syscall.O_RDWR) - if err != nil { - return nil, nil, err - } - _, exit, err := gen.ExitRetEvent(time+gen.pairDelta(), pid, tid, types.SYS_EXIT_OPEN_BY_HANDLE_AT, int64(fd)) - if err != nil { - return nil, nil, err - } - return enter, exit, nil + return pairOpenByHandleAt(gen, time, pid, tid, fd) default: return gen.NullPair(time, pid, tid, types.SYS_ENTER_SYNC, types.SYS_EXIT_SYNC) } } + +// pairFd generates fd-based pairs for read/write/fsync-like events. +func (e MixEvent) pairFd(gen EventGenerator, time uint64, pid, tid uint32, fd int32) ([]byte, []byte, error) { + switch e { + case MixRead: + return gen.FdPair(time, pid, tid, fd, types.SYS_ENTER_READ, types.SYS_EXIT_READ, 128) + case MixWrite: + return gen.FdPair(time, pid, tid, fd, types.SYS_ENTER_WRITE, types.SYS_EXIT_WRITE, 256) + default: // MixFsync + return gen.FdPair(time, pid, tid, fd, types.SYS_ENTER_FSYNC, types.SYS_EXIT_FSYNC, 0) + } +} + +// pairPath generates path-based pairs for stat/access/mkdir/unlink-like events. +func (e MixEvent) pairPath(gen EventGenerator, time uint64, pid, tid uint32, seqTid uint32, seq int) ([]byte, []byte, error) { + name := mixEventPathName(e, seqTid, seq) + enterEv, exitEv := mixEventPathTraceIDs(e) + return gen.PathPair(time, pid, tid, name, enterEv, exitEv, 0) +} + +// mixEventPathTraceIDs returns the enter/exit TraceId constants for a path-based mix event. +func mixEventPathTraceIDs(e MixEvent) (types.TraceId, types.TraceId) { + switch e { + case MixStat: + return types.SYS_ENTER_NEWSTAT, types.SYS_EXIT_NEWSTAT + case MixAccess: + return types.SYS_ENTER_ACCESS, types.SYS_EXIT_ACCESS + case MixMkdir: + return types.SYS_ENTER_MKDIR, types.SYS_EXIT_MKDIR + default: // MixUnlink + return types.SYS_ENTER_UNLINK, types.SYS_EXIT_UNLINK + } +} + +// mixEventPathName returns the synthetic /tmp path for path-based mix events. +func mixEventPathName(e MixEvent, tid uint32, seq int) string { + switch e { + case MixStat: + return fmt.Sprintf("/tmp/ior-stat-%d-%d", tid, seq) + case MixAccess: + return fmt.Sprintf("/tmp/ior-access-%d-%d", tid, seq) + case MixMkdir: + return fmt.Sprintf("/tmp/ior-mkdir-%d-%d", tid, seq) + default: // MixUnlink + return fmt.Sprintf("/tmp/ior-unlink-%d-%d", tid, seq) + } +} + +// pairName generates name-pair (two-path) events for rename/link-like events. +func (e MixEvent) pairName(gen EventGenerator, time uint64, pid, tid uint32, seqTid uint32, seq int) ([]byte, []byte, error) { + if e == MixRename { + oldname := fmt.Sprintf("/tmp/ior-old-%d-%d", seqTid, seq) + newname := fmt.Sprintf("/tmp/ior-new-%d-%d", seqTid, seq) + return gen.NamePair(time, pid, tid, oldname, newname, types.SYS_ENTER_RENAME, types.SYS_EXIT_RENAME, 0) + } + oldname := fmt.Sprintf("/tmp/ior-link-old-%d-%d", seqTid, seq) + newname := fmt.Sprintf("/tmp/ior-link-new-%d-%d", seqTid, seq) + return gen.NamePair(time, pid, tid, oldname, newname, types.SYS_ENTER_LINK, types.SYS_EXIT_LINK, 0) +} + +// pairClose generates a close syscall pair (not available as a single gen helper). +func pairClose(gen EventGenerator, time uint64, pid, tid uint32, fd int32) ([]byte, []byte, error) { + _, enter, err := gen.EnterFdEvent(time, pid, tid, fd, types.SYS_ENTER_CLOSE) + if err != nil { + return nil, nil, err + } + _, exit, err := gen.ExitFdEvent(time+gen.pairDelta(), pid, tid, fd, types.SYS_EXIT_CLOSE) + if err != nil { + return nil, nil, err + } + return enter, exit, nil +} + +// pairOpenByHandleAt generates an open_by_handle_at pair using enter and ret events. +func pairOpenByHandleAt(gen EventGenerator, time uint64, pid, tid uint32, fd int32) ([]byte, []byte, error) { + _, enter, err := gen.EnterOpenByHandleAtEvent(time, pid, tid, syscall.O_RDWR) + if err != nil { + return nil, nil, err + } + _, exit, err := gen.ExitRetEvent(time+gen.pairDelta(), pid, tid, types.SYS_EXIT_OPEN_BY_HANDLE_AT, int64(fd)) + if err != nil { + return nil, nil, err + } + return enter, exit, nil +} diff --git a/internal/eventloop_exit.go b/internal/eventloop_exit.go index e4ae6eb..79c1b5b 100644 --- a/internal/eventloop_exit.go +++ b/internal/eventloop_exit.go @@ -94,12 +94,33 @@ func (e *eventLoop) handlePathExit(ep *event.Pair, pathEv *types.PathEvent) bool return true } +// handleFdExit processes exit events for fd-based syscalls. It resolves the fd +// to a file, applies close/close_range state transitions, filters the pair, and +// handles dup/pidfd_getfd fd-transfer operations before finalising bytes. func (e *eventLoop) handleFdExit(ep *event.Pair, fdEv *types.FdEvent) bool { fd := fdEv.Fd ep.File = e.fdState().resolve(fd, fdEv.Pid) + e.applyFdCloseState(ep, fd, fdEv.Pid) + ep.Comm = e.comm(fdEv.GetTid()) + if !e.Filter().MatchPair(ep) { + ep.Recycle() + return false + } + if ok := e.applyFdTransferOp(ep, fdEv); !ok { + return false + } + if retEv, ok := ep.ExitEv.(*types.RetEvent); ok { + ep.Bytes = bytesFromRet(retEv) + } + return true +} + +// applyFdCloseState updates fd-tracking state for close and close_range syscalls. +func (e *eventLoop) applyFdCloseState(ep *event.Pair, fd int32, pid uint32) { if ep.Is(types.SYS_ENTER_CLOSE) { e.fdState().delete(fd) - e.fdState().deleteProcFdCache(fd, fdEv.Pid) + e.fdState().deleteProcFdCache(fd, pid) + return } if ep.Is(types.SYS_ENTER_CLOSE_RANGE) { // close_range provides (first, last), but fd_event only carries the first @@ -107,15 +128,14 @@ func (e *eventLoop) handleFdExit(ep *event.Pair, fdEv *types.FdEvent) bool { retEv, ok := ep.ExitEv.(*types.RetEvent) if ok && retEv.Ret == 0 { e.fdState().closeRangeFrom(fd) - e.fdState().deleteProcFdCacheFrom(fd, fdEv.Pid) + e.fdState().deleteProcFdCacheFrom(fd, pid) } } - ep.Comm = e.comm(fdEv.GetTid()) - if !e.Filter().MatchPair(ep) { - ep.Recycle() - return false - } +} +// applyFdTransferOp handles dup/dup2 and pidfd_getfd fd-transfer operations. +// Returns false if the pair should be dropped due to a malformed event. +func (e *eventLoop) applyFdTransferOp(ep *event.Pair, fdEv *types.FdEvent) bool { if ep.Is(types.SYS_ENTER_DUP) || ep.Is(types.SYS_ENTER_DUP2) { fdFile, ok := ep.File.(*file.FdFile) if !ok { @@ -127,7 +147,6 @@ func (e *eventLoop) handleFdExit(ep *event.Pair, fdEv *types.FdEvent) bool { e.recyclePair(ep, "Dropped malformed dup exit event") return false } - // Duplicating fd e.registerDup(fdFile, int32(retEvent.Ret), 0) } if ep.Is(types.SYS_ENTER_PIDFD_GETFD) { @@ -142,9 +161,6 @@ func (e *eventLoop) handleFdExit(ep *event.Pair, fdEv *types.FdEvent) bool { ep.File = transferredFile } } - if retEv, ok := ep.ExitEv.(*types.RetEvent); ok { - ep.Bytes = bytesFromRet(retEv) - } return true } diff --git a/internal/eventloop_runtime.go b/internal/eventloop_runtime.go index 5addd46..74571c8 100644 --- a/internal/eventloop_runtime.go +++ b/internal/eventloop_runtime.go @@ -129,6 +129,10 @@ func (e *eventLoop) processRawEvent(raw []byte, ch chan<- *event.Pair) { handler(raw, ch) } +// initRawHandlers registers all BPF event-type dispatch callbacks. It is +// idempotent: a second call after the map is populated is a no-op. Handlers +// are grouped by event class (open, fd, null, ret, name/path, misc) so that +// each helper stays under 30 lines. func (e *eventLoop) initRawHandlers() { if e.rawHandlers == nil { e.rawHandlers = make(map[types.EventType]rawEventHandler) @@ -136,7 +140,16 @@ func (e *eventLoop) initRawHandlers() { if len(e.rawHandlers) != 0 { return } + e.registerOpenHandlers() + e.registerFdHandlers() + e.registerNullHandlers() + e.registerRetHandlers() + e.registerNamePathHandlers() + e.registerMiscHandlers() +} +// registerOpenHandlers wires enter/exit handlers for open-family events. +func (e *eventLoop) registerOpenHandlers() { e.rawHandlers[types.ENTER_OPEN_EVENT] = func(raw []byte, _ chan<- *event.Pair) { openEv, ok := decodeRawEvent(e, types.ENTER_OPEN_EVENT, raw, types.NewOpenEventFast) if !ok { @@ -153,6 +166,10 @@ func (e *eventLoop) initRawHandlers() { } e.tracepointExited(retEv, ch) } +} + +// registerFdHandlers wires enter/exit handlers for fd-family events (read/write/close…). +func (e *eventLoop) registerFdHandlers() { e.rawHandlers[types.ENTER_FD_EVENT] = func(raw []byte, _ chan<- *event.Pair) { fdEv, ok := decodeRawEvent(e, types.ENTER_FD_EVENT, raw, types.NewFdEventFast) if !ok { @@ -167,6 +184,10 @@ func (e *eventLoop) initRawHandlers() { } e.tracepointExited(fdEv, ch) } +} + +// registerNullHandlers wires enter/exit handlers for syscalls with no interesting arguments. +func (e *eventLoop) registerNullHandlers() { e.rawHandlers[types.ENTER_NULL_EVENT] = func(raw []byte, _ chan<- *event.Pair) { nullEv, ok := decodeRawEvent(e, types.ENTER_NULL_EVENT, raw, types.NewNullEventFast) if !ok { @@ -181,6 +202,10 @@ func (e *eventLoop) initRawHandlers() { } e.tracepointExited(nullEv, ch) } +} + +// registerRetHandlers wires the exit handler for generic return-value events. +func (e *eventLoop) registerRetHandlers() { e.rawHandlers[types.EXIT_RET_EVENT] = func(raw []byte, ch chan<- *event.Pair) { retEv, ok := decodeRawEvent(e, types.EXIT_RET_EVENT, raw, types.NewRetEventFast) if !ok { @@ -188,6 +213,10 @@ func (e *eventLoop) initRawHandlers() { } e.tracepointExited(retEv, ch) } +} + +// registerNamePathHandlers wires enter handlers for name- and path-carrying events. +func (e *eventLoop) registerNamePathHandlers() { e.rawHandlers[types.ENTER_NAME_EVENT] = func(raw []byte, _ chan<- *event.Pair) { nameEv, ok := decodeRawEvent(e, types.ENTER_NAME_EVENT, raw, types.NewNameEventFast) if !ok { @@ -206,6 +235,10 @@ func (e *eventLoop) initRawHandlers() { e.tracepointEntered(pathEv) } } +} + +// registerMiscHandlers wires enter handlers for fcntl, open_by_handle_at, and dup3. +func (e *eventLoop) registerMiscHandlers() { e.rawHandlers[types.ENTER_FCNTL_EVENT] = func(raw []byte, _ chan<- *event.Pair) { fcntlEv, ok := decodeRawEvent(e, types.ENTER_FCNTL_EVENT, raw, types.NewFcntlEventFast) if !ok { diff --git a/internal/export/snapshot_csv.go b/internal/export/snapshot_csv.go index 591bd67..6f7312a 100644 --- a/internal/export/snapshot_csv.go +++ b/internal/export/snapshot_csv.go @@ -24,65 +24,89 @@ func SnapshotCSV(snap *statsengine.Snapshot) (filename string, retErr error) { }() w := csv.NewWriter(f) + if err := writeSnapshotRows(w, snap); err != nil { + return "", err + } + w.Flush() + if err := w.Error(); err != nil { + return "", err + } + return filename, nil +} - rows := [][]string{ +// writeSnapshotRows writes all CSV sections to w in order: +// header, summary, per-syscall stats, file stats, process stats, histograms. +func writeSnapshotRows(w *csv.Writer, snap *statsengine.Snapshot) error { + summaryRows := [][]string{ {"section", "name", "value1", "value2", "value3"}, - {"summary", "totals", fmt.Sprint(snapValue(snap, func(s *statsengine.Snapshot) uint64 { return s.TotalSyscalls })), fmt.Sprint(snapValue(snap, func(s *statsengine.Snapshot) uint64 { return s.TotalErrors })), fmt.Sprint(snapValue(snap, func(s *statsengine.Snapshot) uint64 { return s.TotalBytes }))}, - {"summary", "rates_per_sec", fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.SyscallRatePerSec })), fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.ReadBytesPerSec })), fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.WriteBytesPerSec }))}, - {"summary", "latency_gap_mean_ns", fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.LatencyMeanNs })), fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.GapMeanNs })), ""}, - {"summary", "trend", trendSummary(snap, func(s *statsengine.Snapshot) statsengine.Trend { return s.LatencyTrend }), trendSummary(snap, func(s *statsengine.Snapshot) statsengine.Trend { return s.GapTrend }), trendSummary(snap, func(s *statsengine.Snapshot) statsengine.Trend { return s.ThroughputTrend })}, + {"summary", "totals", + fmt.Sprint(snapValue(snap, func(s *statsengine.Snapshot) uint64 { return s.TotalSyscalls })), + fmt.Sprint(snapValue(snap, func(s *statsengine.Snapshot) uint64 { return s.TotalErrors })), + fmt.Sprint(snapValue(snap, func(s *statsengine.Snapshot) uint64 { return s.TotalBytes }))}, + {"summary", "rates_per_sec", + fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.SyscallRatePerSec })), + fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.ReadBytesPerSec })), + fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.WriteBytesPerSec }))}, + {"summary", "latency_gap_mean_ns", + fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.LatencyMeanNs })), + fmt.Sprintf("%.2f", snapValueF(snap, func(s *statsengine.Snapshot) float64 { return s.GapMeanNs })), ""}, + {"summary", "trend", + trendSummary(snap, func(s *statsengine.Snapshot) statsengine.Trend { return s.LatencyTrend }), + trendSummary(snap, func(s *statsengine.Snapshot) statsengine.Trend { return s.GapTrend }), + trendSummary(snap, func(s *statsengine.Snapshot) statsengine.Trend { return s.ThroughputTrend })}, } - for _, row := range rows { + for _, row := range summaryRows { if err := w.Write(row); err != nil { - return "", err + return err } } + if snap == nil { + return nil + } + return writeSnapshotDetailRows(w, snap) +} - if snap != nil { - for _, s := range snap.Syscalls() { - if err := w.Write([]string{"syscall", s.Name, fmt.Sprint(s.Count), fmt.Sprintf("%.2f", s.RatePerSec), fmt.Sprint(s.Bytes)}); err != nil { - return "", err - } - if err := w.Write([]string{"syscall_latency_ns", s.Name, fmt.Sprintf("%.2f", s.LatencyMeanNs), fmt.Sprint(s.LatencyMinNs), fmt.Sprint(s.LatencyMaxNs)}); err != nil { - return "", err - } - if err := w.Write([]string{"syscall_percentiles_ns", s.Name, fmt.Sprint(s.LatencyP50Ns), fmt.Sprint(s.LatencyP95Ns), fmt.Sprint(s.LatencyP99Ns)}); err != nil { - return "", err - } +// writeSnapshotDetailRows writes per-item rows for syscalls, files, processes, +// and histograms. It is called only when snap is non-nil. +func writeSnapshotDetailRows(w *csv.Writer, snap *statsengine.Snapshot) error { + for _, s := range snap.Syscalls() { + if err := w.Write([]string{"syscall", s.Name, fmt.Sprint(s.Count), fmt.Sprintf("%.2f", s.RatePerSec), fmt.Sprint(s.Bytes)}); err != nil { + return err } - for _, r := range snap.Files() { - if err := w.Write([]string{"file", r.Path, fmt.Sprint(r.Accesses), fmt.Sprint(r.BytesRead), fmt.Sprint(r.BytesWritten)}); err != nil { - return "", err - } - if err := w.Write([]string{"file_latency_ns", r.Path, fmt.Sprintf("%.2f", r.AvgLatencyNs), fmt.Sprint(r.MaxLatencyNs), ""}); err != nil { - return "", err - } + if err := w.Write([]string{"syscall_latency_ns", s.Name, fmt.Sprintf("%.2f", s.LatencyMeanNs), fmt.Sprint(s.LatencyMinNs), fmt.Sprint(s.LatencyMaxNs)}); err != nil { + return err } - for _, p := range snap.Processes() { - if err := w.Write([]string{"process", fmt.Sprint(p.PID), fmt.Sprint(p.Syscalls), fmt.Sprintf("%.2f", p.RatePerSec), fmt.Sprint(p.Bytes)}); err != nil { - return "", err - } - if err := w.Write([]string{"process_latency_ns", fmt.Sprint(p.PID), fmt.Sprintf("%.2f", p.AvgLatencyNs), "", ""}); err != nil { - return "", err - } + if err := w.Write([]string{"syscall_percentiles_ns", s.Name, fmt.Sprint(s.LatencyP50Ns), fmt.Sprint(s.LatencyP95Ns), fmt.Sprint(s.LatencyP99Ns)}); err != nil { + return err } - for _, b := range snap.LatencyHistogram.Buckets() { - if err := w.Write([]string{"latency_hist", b.Label, fmt.Sprint(b.Count), fmt.Sprint(b.LowerNs), fmt.Sprint(b.UpperNs)}); err != nil { - return "", err - } + } + for _, r := range snap.Files() { + if err := w.Write([]string{"file", r.Path, fmt.Sprint(r.Accesses), fmt.Sprint(r.BytesRead), fmt.Sprint(r.BytesWritten)}); err != nil { + return err } - for _, b := range snap.GapHistogram.Buckets() { - if err := w.Write([]string{"gap_hist", b.Label, fmt.Sprint(b.Count), fmt.Sprint(b.LowerNs), fmt.Sprint(b.UpperNs)}); err != nil { - return "", err - } + if err := w.Write([]string{"file_latency_ns", r.Path, fmt.Sprintf("%.2f", r.AvgLatencyNs), fmt.Sprint(r.MaxLatencyNs), ""}); err != nil { + return err } } - - w.Flush() - if err := w.Error(); err != nil { - return "", err + for _, p := range snap.Processes() { + if err := w.Write([]string{"process", fmt.Sprint(p.PID), fmt.Sprint(p.Syscalls), fmt.Sprintf("%.2f", p.RatePerSec), fmt.Sprint(p.Bytes)}); err != nil { + return err + } + if err := w.Write([]string{"process_latency_ns", fmt.Sprint(p.PID), fmt.Sprintf("%.2f", p.AvgLatencyNs), "", ""}); err != nil { + return err + } } - return filename, nil + for _, b := range snap.LatencyHistogram.Buckets() { + if err := w.Write([]string{"latency_hist", b.Label, fmt.Sprint(b.Count), fmt.Sprint(b.LowerNs), fmt.Sprint(b.UpperNs)}); err != nil { + return err + } + } + for _, b := range snap.GapHistogram.Buckets() { + if err := w.Write([]string{"gap_hist", b.Label, fmt.Sprint(b.Count), fmt.Sprint(b.LowerNs), fmt.Sprint(b.UpperNs)}); err != nil { + return err + } + } + return nil } func snapValue(snap *statsengine.Snapshot, get func(*statsengine.Snapshot) uint64) uint64 { diff --git a/internal/flags/flags.go b/internal/flags/flags.go index 2544007..a46f6b3 100644 --- a/internal/flags/flags.go +++ b/internal/flags/flags.go @@ -131,6 +131,23 @@ func Parse() (Config, error) { // fresh FlagSet and custom argument slices without touching global state. func parseFromFlagSet(fs *flag.FlagSet, args []string) (Config, error) { cfg := NewFlags() + tpsAttach, tpsExclude, fields := registerFlags(fs, &cfg) + + if err := fs.Parse(args); err != nil { + return Config{}, err + } + if err := resolvePostParseFields(&cfg, tpsAttach, tpsExclude, fields); err != nil { + return Config{}, err + } + if err := validateConfig(cfg); err != nil { + return Config{}, err + } + return cfg, nil +} + +// registerFlags binds all CLI flags to cfg and returns the string pointers for +// fields that require post-parse resolution (tracepoint regexes, collapse fields). +func registerFlags(fs *flag.FlagSet, cfg *Config) (tpsAttach, tpsExclude, fields *string) { validFields := collapse.ValidFields() validCounts := collapse.ValidCountFields() @@ -141,11 +158,10 @@ func parseFromFlagSet(fs *flag.FlagSet, args []string) (Config, error) { fs.StringVar(&cfg.CommFilter, "comm", "", "Command to filter for") fs.StringVar(&cfg.PathFilter, "path", "", "Path to filter for") - fs.BoolVar(&cfg.PprofEnable, "pprof", false, "Enable profiling") - tpsAttach := fs.String("tps", "", "Comma separated list regexes for tracepoints to load") - tpsExclude := fs.String("tpsExclude", "", "Comma separated list regexes for tracepoints to exclude") + tpsAttach = fs.String("tps", "", "Comma separated list regexes for tracepoints to load") + tpsExclude = fs.String("tpsExclude", "", "Comma separated list regexes for tracepoints to exclude") fs.BoolVar(&cfg.PlainMode, "plain", false, "Enable plain CSV output mode (disable TUI)") fs.BoolVar(&cfg.FlamegraphOutput, "flamegraph", false, "Write aggregated .ior.zst output for trace/integration workflows") @@ -158,20 +174,21 @@ func parseFromFlagSet(fs *flag.FlagSet, args []string) (Config, error) { fs.DurationVar(&cfg.ResetTimer, "resetTimer", cfg.ResetTimer, "Auto-reset interval for aggregate dashboard state (flamegraph trie + stats engine); set to 0 to disable") fs.BoolVar(&cfg.ShowVersion, "version", false, "Print version banner and exit") - fields := fs.String("fields", "", + fields = fs.String("fields", "", fmt.Sprintf("Comma separated list of fields to collapse, valid are: %v", validFields)) fs.StringVar(&cfg.CountField, "count", cfg.CountField, fmt.Sprintf("Count field to collapse, valid are: %v", validCounts)) + return tpsAttach, tpsExclude, fields +} - if err := fs.Parse(args); err != nil { - return Config{}, err - } - +// resolvePostParseFields compiles the tracepoint selector and collapse field +// list from the raw string flags that cannot be bound directly to cfg fields. +func resolvePostParseFields(cfg *Config, tpsAttach, tpsExclude, fields *string) error { // Parse the tracepoint include/exclude regex lists into a Selector. // The Selector owns all matching logic; Config is purely a data carrier. sel, err := tracepoints.ParseSelector(*tpsAttach, *tpsExclude) if err != nil { - return Config{}, err + return err } cfg.TracepointSelector = sel @@ -179,7 +196,6 @@ func parseFromFlagSet(fs *flag.FlagSet, args []string) (Config, error) { // As of February 23, 2026, open_by_handle_at and name_to_handle_at were // re-evaluated on newer kernels and do not require CO-RE-based exclusions. // If future kernels regress, add targeted exclusions here. - if *fields == "" { cfg.CollapsedFields = []string{"comm", "tracepoint", "path"} } else { @@ -188,32 +204,33 @@ func parseFromFlagSet(fs *flag.FlagSet, args []string) (Config, error) { for _, field := range cfg.CollapsedFields { if !collapse.IsValidField(field) { - return Config{}, fmt.Errorf("invalid field for collapse: %s", field) + return fmt.Errorf("invalid field for collapse: %s", field) } } - if !collapse.IsValidCountField(cfg.CountField) { - return Config{}, fmt.Errorf("invalid count field: %s", cfg.CountField) + return fmt.Errorf("invalid count field: %s", cfg.CountField) } + return nil +} +// validateConfig checks numeric/duration bounds that cannot be enforced by the +// flag package itself and returns a descriptive error on the first violation. +func validateConfig(cfg Config) error { // A zero or negative duration would cause the trace context to cancel // immediately, capturing no events. Require at least one second. if cfg.Duration <= 0 { - return Config{}, fmt.Errorf("invalid duration: %d (must be > 0)", cfg.Duration) + return fmt.Errorf("invalid duration: %d (must be > 0)", cfg.Duration) } - // A negative reset timer would imply auto-resets in the past, which is // nonsensical. 0 disables, anything positive enables. if cfg.ResetTimer < 0 { - return Config{}, fmt.Errorf("invalid resetTimer: %s (must be >= 0; 0 disables)", cfg.ResetTimer) + return fmt.Errorf("invalid resetTimer: %s (must be >= 0; 0 disables)", cfg.ResetTimer) } - // A non-positive mapSize would wrap to a huge uint32 when cast in // resizeBPFMaps, causing libbpf to fail with a confusing "map too large" // error. Reject it here with a clear diagnostic instead. if cfg.EventMapSize <= 0 { - return Config{}, fmt.Errorf("invalid mapSize: %d (must be > 0)", cfg.EventMapSize) + return fmt.Errorf("invalid mapSize: %d (must be > 0)", cfg.EventMapSize) } - - return cfg, nil + return nil } diff --git a/internal/generate/bpfhandler.go b/internal/generate/bpfhandler.go index cf9a0c9..549e80b 100644 --- a/internal/generate/bpfhandler.go +++ b/internal/generate/bpfhandler.go @@ -62,77 +62,93 @@ func renderHandler(name, ctxStruct, eventStruct, comment, eventTypeConst, extra return b.String() } +// generateExtra returns the kind-specific C body lines for a tracepoint handler, +// dispatching to a per-kind helper so that each case stays concise. func generateExtra(tp GeneratedTracepoint, isEnter bool) string { f := tp.Format - switch tp.Classification.Kind { case KindFd: - if f.Name == "sys_enter_pidfd_getfd" { - return " ev->fd = (__s32)ctx->args[0];\n" - } - fdIdx := f.FieldNumber("fd") - if fdIdx >= 0 { - return fmt.Sprintf(" ev->fd = (__s32)ctx->args[%d];\n", fdIdx) - } - return " ev->fd = (__s32)ctx->args[0];\n" - + return generateExtraFd(f) case KindDup3: return " ev->fd = (__s32)ctx->args[0];\n ev->flags = (__s32)ctx->args[2];\n" - case KindOpenByHandleAt: return " ev->flags = (__s32)ctx->args[2];\n" - case KindOpen: - filenameIdx := f.FieldNumber("filename") - flagsIdx := f.FieldNumber("flags") - var b strings.Builder - b.WriteString(" __builtin_memset(&(ev->filename), 0, sizeof(ev->filename) + sizeof(ev->comm));\n") - fmt.Fprintf(&b, " bpf_probe_read_user_str(ev->filename, sizeof(ev->filename), (void *)ctx->args[%d]);\n", filenameIdx) - b.WriteString(" bpf_get_current_comm(&ev->comm, sizeof(ev->comm));\n") - if flagsIdx > -1 { - fmt.Fprintf(&b, " ev->flags = ctx->args[%d];\n", flagsIdx) - } else { - b.WriteString(" ev->flags = -1; // Probably OK\n") - } - return b.String() - + return generateExtraOpen(f) case KindPathname: - fieldName := tp.Classification.PathnameField - fieldIdx := f.FieldNumber(fieldName) - var b strings.Builder - b.WriteString(" __builtin_memset(&(ev->pathname), 0, sizeof(ev->pathname));\n") - fmt.Fprintf(&b, " bpf_probe_read_user_str(ev->pathname, sizeof(ev->pathname), (void*)ctx->args[%d]);\n", fieldIdx) - return b.String() - + return generateExtraPathname(tp, f) case KindName: - oldIdx := f.FieldNumber("oldname") - newIdx := f.FieldNumber("newname") - var b strings.Builder - b.WriteString(" __builtin_memset(&(ev->oldname), 0, sizeof(ev->oldname) + sizeof(ev->newname));\n") - fmt.Fprintf(&b, " bpf_probe_read_user_str(ev->oldname, sizeof(ev->oldname), (void*)ctx->args[%d]);\n", oldIdx) - fmt.Fprintf(&b, " bpf_probe_read_user_str(ev->newname, sizeof(ev->newname), (void*)ctx->args[%d]);\n", newIdx) - return b.String() - + return generateExtraName(f) case KindFcntl: - fdIdx := f.FieldNumber("fd") - cmdIdx := f.FieldNumber("cmd") - argIdx := f.FieldNumber("arg") - return fmt.Sprintf( - " ev->fd = ctx->args[%d];\n ev->cmd = ctx->args[%d];\n ev->arg = ctx->args[%d];\n", - fdIdx, cmdIdx, argIdx, - ) - + return generateExtraFcntl(f) case KindRet: - classification := ClassifyRet(f.Name) - return fmt.Sprintf(" ev->ret = ctx->ret;\n ev->ret_type = %s;\n", classification) - + return fmt.Sprintf(" ev->ret = ctx->ret;\n ev->ret_type = %s;\n", ClassifyRet(f.Name)) case KindNull: return "" } - return "" } +// generateExtraFd returns the fd-capture lines for fd-family events. +func generateExtraFd(f *Format) string { + if f.Name == "sys_enter_pidfd_getfd" { + return " ev->fd = (__s32)ctx->args[0];\n" + } + fdIdx := f.FieldNumber("fd") + if fdIdx >= 0 { + return fmt.Sprintf(" ev->fd = (__s32)ctx->args[%d];\n", fdIdx) + } + return " ev->fd = (__s32)ctx->args[0];\n" +} + +// generateExtraOpen returns the filename/comm/flags capture lines for open-family events. +func generateExtraOpen(f *Format) string { + filenameIdx := f.FieldNumber("filename") + flagsIdx := f.FieldNumber("flags") + var b strings.Builder + b.WriteString(" __builtin_memset(&(ev->filename), 0, sizeof(ev->filename) + sizeof(ev->comm));\n") + fmt.Fprintf(&b, " bpf_probe_read_user_str(ev->filename, sizeof(ev->filename), (void *)ctx->args[%d]);\n", filenameIdx) + b.WriteString(" bpf_get_current_comm(&ev->comm, sizeof(ev->comm));\n") + if flagsIdx > -1 { + fmt.Fprintf(&b, " ev->flags = ctx->args[%d];\n", flagsIdx) + } else { + b.WriteString(" ev->flags = -1; // Probably OK\n") + } + return b.String() +} + +// generateExtraPathname returns the pathname capture lines for path-family events. +func generateExtraPathname(tp GeneratedTracepoint, f *Format) string { + fieldName := tp.Classification.PathnameField + fieldIdx := f.FieldNumber(fieldName) + var b strings.Builder + b.WriteString(" __builtin_memset(&(ev->pathname), 0, sizeof(ev->pathname));\n") + fmt.Fprintf(&b, " bpf_probe_read_user_str(ev->pathname, sizeof(ev->pathname), (void*)ctx->args[%d]);\n", fieldIdx) + return b.String() +} + +// generateExtraName returns the oldname/newname capture lines for rename/link-family events. +func generateExtraName(f *Format) string { + oldIdx := f.FieldNumber("oldname") + newIdx := f.FieldNumber("newname") + var b strings.Builder + b.WriteString(" __builtin_memset(&(ev->oldname), 0, sizeof(ev->oldname) + sizeof(ev->newname));\n") + fmt.Fprintf(&b, " bpf_probe_read_user_str(ev->oldname, sizeof(ev->oldname), (void*)ctx->args[%d]);\n", oldIdx) + fmt.Fprintf(&b, " bpf_probe_read_user_str(ev->newname, sizeof(ev->newname), (void*)ctx->args[%d]);\n", newIdx) + return b.String() +} + +// generateExtraFcntl returns the fd/cmd/arg capture lines for fcntl events. +func generateExtraFcntl(f *Format) string { + fdIdx := f.FieldNumber("fd") + cmdIdx := f.FieldNumber("cmd") + argIdx := f.FieldNumber("arg") + return fmt.Sprintf( + " ev->fd = ctx->args[%d];\n ev->cmd = ctx->args[%d];\n ev->arg = ctx->args[%d];\n", + fdIdx, cmdIdx, argIdx, + ) +} + // eventStructName returns the C struct name for a TracepointKind. The mapping // is driven by kindRegistry so adding a new kind only requires a registry entry. func eventStructName(kind TracepointKind) string { diff --git a/internal/generate/format.go b/internal/generate/format.go index ea579b6..ef51ba8 100644 --- a/internal/generate/format.go +++ b/internal/generate/format.go @@ -44,52 +44,60 @@ func ParseFormats(r io.Reader) ([]Format, error) { isExternal := false for scanner.Scan() { - line := scanner.Text() - trimmed := strings.TrimSpace(line) - - switch { - case strings.HasPrefix(trimmed, "name:"): - f := Format{} - f.Name = strings.TrimSpace(strings.TrimPrefix(trimmed, "name:")) - formats = append(formats, f) - current = &formats[len(formats)-1] - isExternal = false - - case strings.HasPrefix(trimmed, "ID:"): - if current == nil { - return nil, fmt.Errorf("ID without name") - } - id, err := strconv.Atoi(strings.TrimSpace(strings.TrimPrefix(trimmed, "ID:"))) - if err != nil { - return nil, fmt.Errorf("parsing ID: %w", err) - } - current.ID = id - - case strings.HasPrefix(trimmed, "field:"): - if current == nil { - return nil, fmt.Errorf("field without name") - } - field, err := parseField(trimmed) - if err != nil { - return nil, fmt.Errorf("parsing field in %s: %w", current.Name, err) - } - if field.Name == "__syscall_nr" { - isExternal = true - } - if isExternal { - current.ExternalFields = append(current.ExternalFields, field) - } else { - current.InternalFields = append(current.InternalFields, field) - } + var err error + current, isExternal, err = applyFormatLine(scanner.Text(), formats, current, isExternal, &formats) + if err != nil { + return nil, err } } - if err := scanner.Err(); err != nil { return nil, fmt.Errorf("scanning input: %w", err) } return formats, nil } +// applyFormatLine processes one line from the format file, updating the +// running parse state (current format, isExternal flag, and formats slice). +func applyFormatLine(line string, _ []Format, current *Format, isExternal bool, formats *[]Format) (*Format, bool, error) { + trimmed := strings.TrimSpace(line) + switch { + case strings.HasPrefix(trimmed, "name:"): + f := Format{} + f.Name = strings.TrimSpace(strings.TrimPrefix(trimmed, "name:")) + *formats = append(*formats, f) + current = &(*formats)[len(*formats)-1] + isExternal = false + + case strings.HasPrefix(trimmed, "ID:"): + if current == nil { + return nil, false, fmt.Errorf("ID without name") + } + id, err := strconv.Atoi(strings.TrimSpace(strings.TrimPrefix(trimmed, "ID:"))) + if err != nil { + return nil, false, fmt.Errorf("parsing ID: %w", err) + } + current.ID = id + + case strings.HasPrefix(trimmed, "field:"): + if current == nil { + return nil, false, fmt.Errorf("field without name") + } + field, err := parseField(trimmed) + if err != nil { + return nil, false, fmt.Errorf("parsing field in %s: %w", current.Name, err) + } + if field.Name == "__syscall_nr" { + isExternal = true + } + if isExternal { + current.ExternalFields = append(current.ExternalFields, field) + } else { + current.InternalFields = append(current.InternalFields, field) + } + } + return current, isExternal, nil +} + func parseField(line string) (Field, error) { // Format: "field:TYPE NAME; offset:N; size:N; signed:N;" line = strings.TrimPrefix(line, "field:") diff --git a/internal/ior.go b/internal/ior.go index c8aa47b..decdf12 100644 --- a/internal/ior.go +++ b/internal/ior.go @@ -495,61 +495,92 @@ func finaliseTrace(watcherDone <-chan struct{}, recorder *flamegraph.Recorder, p // is checked by the mode handler (via runnerDeps.getEUID) before calling this // function; the handler is the authoritative place for the EUID gate. func runTraceWithContext(parentCtx context.Context, cfg flags.Config, started chan<- struct{}, configure func(*eventLoop)) error { - verbose := started == nil logln := newLogger(verbose) configure, recorder := maybePrependFlamegraphConfigure(cfg, configure) - bpfModule, mgr, releaseBindings, err := setupBPFModule(parentCtx, cfg) + ch, ctx, cancel, profiling, el, mgr, teardown, err := setupTraceInfra(parentCtx, cfg, started, logln) if err != nil { return err } - defer bpfModule.Close() - // mgr.Close() detaches BPF probes and releases kernel resources; log any - // error so that probe-detach failures are not silently discarded. - defer func() { - if err := mgr.Close(); err != nil { - logln("BPF probe manager close error:", err) - } - }() - defer releaseBindings() + defer teardown() + defer profiling.stop(logln) + defer cancel() + + configureEventLoopOutput(el, mgr, configure) + watcherDone := startTraceShutdownWatcher(ctx, verbose, el, profiling, logln) + + startTime := time.Now() + el.run(ctx, ch) + return finaliseTrace(watcherDone, recorder, profiling, time.Since(startTime), logln) +} + +// setupTraceInfra creates all the BPF/runtime infrastructure for a trace run: +// BPF module + probe manager, event channel + ring buffer, trace context, +// profiling control, and event loop. teardown must be deferred by the caller. +// started is signalled once setup completes (nil in non-TUI modes). +func setupTraceInfra( + parentCtx context.Context, + cfg flags.Config, + started chan<- struct{}, + logln func(...any), +) ( + ch <-chan []byte, + ctx context.Context, + cancel context.CancelFunc, + profiling *profilingControl, + el *eventLoop, + mgr *probemanager.Manager, + teardown func(), + err error, +) { + bpfModule, mgr, releaseBindings, err := setupBPFModule(parentCtx, cfg) + if err != nil { + return nil, nil, nil, nil, nil, nil, func() {}, err + } - ch, rb, err := setupEventChannel(bpfModule) + eventCh, rb, err := setupEventChannel(bpfModule) if err != nil { - return err + bpfModule.Close() + return nil, nil, nil, nil, nil, nil, func() {}, err } - // Stop the ring-buffer polling goroutine before the module is closed. - // rb.Stop() signals the background goroutine, drains the channel, and - // waits for the goroutine to exit; bpfModule.Close() (deferred above) - // then calls rb.Close() which frees the C ring_buffer struct. Both are - // idempotent so double-calling is safe. - defer rb.Stop() + ctx, cancel, stopSignals := setupTraceContext(parentCtx, cfg, logln) - defer cancel() - defer stopSignals() - profiling, err := setupProfiling(ctx, cfg, started) + profiling, err = setupProfiling(ctx, cfg, started) if err != nil { - return err + cancel() + stopSignals() + rb.Stop() + bpfModule.Close() + return nil, nil, nil, nil, nil, nil, func() {}, err } - // Guarantee the profiling file descriptors (cpu/mem/exec-trace profiles) are - // closed even if a later setup step fails before the shutdown watcher is - // registered. profiling.stop is idempotent via sync.Once, so double-calling - // it from the watcher goroutine and from this defer is safe. - defer profiling.stop(logln) signalTraceStarted(started) - el, err := newEventLoop(newEventLoopConfig(cfg)) + el, err = newEventLoop(newEventLoopConfig(cfg)) if err != nil { - return err + cancel() + stopSignals() + rb.Stop() + bpfModule.Close() + return nil, nil, nil, nil, nil, nil, func() {}, err } - configureEventLoopOutput(el, mgr, configure) - watcherDone := startTraceShutdownWatcher(ctx, verbose, el, profiling, logln) - startTime := time.Now() - el.run(ctx, ch) - return finaliseTrace(watcherDone, recorder, profiling, time.Since(startTime), logln) + teardown = func() { + // Stop the ring-buffer polling goroutine before the module is closed. + // rb.Stop() is idempotent; bpfModule.Close() calls rb.Close() for the C struct. + rb.Stop() + // mgr.Close() detaches BPF probes and releases kernel resources; log any + // error so that probe-detach failures are not silently discarded. + if err := mgr.Close(); err != nil { + logln("BPF probe manager close error:", err) + } + releaseBindings() + bpfModule.Close() + stopSignals() + } + return eventCh, ctx, cancel, profiling, el, mgr, teardown, nil } func chainEventLoopConfigure(fns ...func(*eventLoop)) func(*eventLoop) { diff --git a/internal/ior_parquet_sink.go b/internal/ior_parquet_sink.go index 279c9be..6869e1e 100644 --- a/internal/ior_parquet_sink.go +++ b/internal/ior_parquet_sink.go @@ -11,6 +11,7 @@ import ( "ior/internal/flags" "ior/internal/globalfilter" "ior/internal/parquet" + "ior/internal/probemanager" "ior/internal/streamrow" ) @@ -93,52 +94,16 @@ func headlessParquetTraceConfig(cfg flags.Config) flags.Config { // without starting the TUI. Root privilege is checked by the mode handler // (via runnerDeps.getEUID) before this function is invoked. func runHeadlessParquet(cfg flags.Config) error { - cfg = headlessParquetTraceConfig(cfg) logln := newLogger(true) - bpfModule, mgr, releaseBindings, err := setupBPFModule(context.Background(), cfg) + ch, ctx, cancel, profiling, el, mgr, cleanup, err := setupHeadlessParquetInfra(cfg, logln) if err != nil { return err } - defer bpfModule.Close() - // mgr.Close() detaches BPF probes and releases kernel resources; log any - // error so that probe-detach failures are not silently discarded. - defer func() { - if err := mgr.Close(); err != nil { - logln("BPF probe manager close error:", err) - } - }() - defer releaseBindings() - - ch, rb, err := setupEventChannel(bpfModule) - if err != nil { - return err - } - // Stop the ring-buffer polling goroutine before the module is closed. - // rb.Stop() signals the background goroutine, drains the channel, and - // waits for the goroutine to exit; bpfModule.Close() (deferred above) - // then calls rb.Close() which frees the C ring_buffer struct. Both are - // idempotent so double-calling is safe. - defer rb.Stop() - ctx, cancel, stopSignals := setupTraceContext(context.Background(), cfg, logln) - defer cancel() - defer stopSignals() - - profiling, err := setupProfiling(ctx, cfg, nil) - if err != nil { - return err - } - // Guarantee the profiling file descriptors (cpu/mem/exec-trace profiles) are - // closed even if a later setup step fails before the shutdown watcher is - // registered. profiling.stop is idempotent via sync.Once, so double-calling - // it from the watcher goroutine and from this defer is safe. + defer cleanup() defer profiling.stop(logln) - - el, err := newEventLoop(newEventLoopConfig(cfg)) - if err != nil { - return err - } + defer cancel() recorder := parquet.NewRecorder(parquet.RecorderConfig{}) if err := recorder.Start(cfg.ParquetPath, parquet.StartOptions{Metadata: parquet.NewFileMetadata("headless")}); err != nil { @@ -146,6 +111,8 @@ func runHeadlessParquet(cfg flags.Config) error { } sink := newHeadlessParquetSink(recorder, cancel) + // sink.configure wires the event loop's print callback to record each pair + // to Parquet; the mgr filter wraps it to skip inactive probes. configureEventLoopOutput(el, mgr, sink.configure) // startTraceShutdownWatcher returns a done channel that must be drained // before returning to prevent a goroutine leak when ctx is cancelled but @@ -168,7 +135,67 @@ func runHeadlessParquet(cfg flags.Config) error { if stopErr != nil { return stopErr } - logln("Good bye... (unloading BPF tracepoints will take a few seconds...) after", totalDuration) return nil } + +// setupHeadlessParquetInfra creates the BPF module, event channel, trace +// context, profiling control, and event loop for a headless Parquet run. +// mgr is returned so the caller can pass it to configureEventLoopOutput with +// the sink callback after the Parquet recorder has been started. +// cleanup must be deferred by the caller; it stops ring-buffer polling, +// detaches probes, releases BPF bindings, and stops signal handling. +func setupHeadlessParquetInfra(cfg flags.Config, logln func(...any)) ( + ch <-chan []byte, + ctx context.Context, + cancel context.CancelFunc, + profiling *profilingControl, + el *eventLoop, + mgr *probemanager.Manager, + cleanup func(), + err error, +) { + bpfModule, mgr, releaseBindings, err := setupBPFModule(context.Background(), cfg) + if err != nil { + return nil, nil, nil, nil, nil, nil, func() {}, err + } + + eventCh, rb, err := setupEventChannel(bpfModule) + if err != nil { + bpfModule.Close() + return nil, nil, nil, nil, nil, nil, func() {}, err + } + + ctx, cancel, stopSignals := setupTraceContext(context.Background(), cfg, logln) + + profiling, err = setupProfiling(ctx, cfg, nil) + if err != nil { + cancel() + stopSignals() + rb.Stop() + bpfModule.Close() + return nil, nil, nil, nil, nil, nil, func() {}, err + } + + el, err = newEventLoop(newEventLoopConfig(cfg)) + if err != nil { + cancel() + stopSignals() + rb.Stop() + bpfModule.Close() + return nil, nil, nil, nil, nil, nil, func() {}, err + } + + cleanup = func() { + // Stop the ring-buffer polling goroutine before the module is closed. + // rb.Stop() is idempotent; bpfModule.Close() calls rb.Close() for the C struct. + rb.Stop() + if err := mgr.Close(); err != nil { + logln("BPF probe manager close error:", err) + } + releaseBindings() + bpfModule.Close() + stopSignals() + } + return eventCh, ctx, cancel, profiling, el, mgr, cleanup, nil +} diff --git a/internal/ior_profiling.go b/internal/ior_profiling.go index ddae088..77790b9 100644 --- a/internal/ior_profiling.go +++ b/internal/ior_profiling.go @@ -37,50 +37,21 @@ func setupProfiling(ctx context.Context, cfg flags.Config, started chan<- struct } control.enabled = true - isTUIMode := started != nil - cpuProfilePath, memProfilePath, execTracePath, execTraceDuration := profilingFilesForMode(isTUIMode) + cpuProfilePath, memProfilePath, execTracePath, execTraceDuration := profilingFilesForMode(started != nil) - cpuProfile, err := os.Create(cpuProfilePath) + cpuProfile, memProfile, err := openProfilingFiles(cpuProfilePath, memProfilePath) if err != nil { return nil, err } - memProfile, err := os.Create(memProfilePath) - if err != nil { - _ = cpuProfile.Close() - return nil, err - } control.cpuProfile = cpuProfile control.memProfile = memProfile if execTracePath != "" { - execTraceProfile, err := os.Create(execTracePath) - if err != nil { + if err := startExecTrace(ctx, execTracePath, execTraceDuration, control); err != nil { _ = cpuProfile.Close() _ = memProfile.Close() return nil, err } - if err := trace.Start(execTraceProfile); err != nil { - _ = cpuProfile.Close() - _ = memProfile.Close() - _ = execTraceProfile.Close() - return nil, err - } - var stopOnce sync.Once - control.stopExecTrace = func() { - stopOnce.Do(func() { - trace.Stop() - _ = execTraceProfile.Close() - }) - } - go func() { - timer := time.NewTimer(execTraceDuration) - defer timer.Stop() - select { - case <-ctx.Done(): - case <-timer.C: - } - control.stopExecTrace() - }() } if err := pprof.StartCPUProfile(cpuProfile); err != nil { @@ -92,6 +63,52 @@ func setupProfiling(ctx context.Context, cfg flags.Config, started chan<- struct return control, nil } +// openProfilingFiles creates the CPU and memory profile output files. On +// error any successfully opened file is closed before returning. +func openProfilingFiles(cpuPath, memPath string) (*os.File, *os.File, error) { + cpuProfile, err := os.Create(cpuPath) + if err != nil { + return nil, nil, err + } + memProfile, err := os.Create(memPath) + if err != nil { + _ = cpuProfile.Close() + return nil, nil, err + } + return cpuProfile, memProfile, nil +} + +// startExecTrace creates the execution-trace output file, starts the runtime +// tracer, and wires a goroutine that stops it on context cancellation or after +// execTraceDuration, whichever comes first. +func startExecTrace(ctx context.Context, tracePath string, execTraceDuration time.Duration, control *profilingControl) error { + execTraceProfile, err := os.Create(tracePath) + if err != nil { + return err + } + if err := trace.Start(execTraceProfile); err != nil { + _ = execTraceProfile.Close() + return err + } + var stopOnce sync.Once + control.stopExecTrace = func() { + stopOnce.Do(func() { + trace.Stop() + _ = execTraceProfile.Close() + }) + } + go func() { + timer := time.NewTimer(execTraceDuration) + defer timer.Stop() + select { + case <-ctx.Done(): + case <-timer.C: + } + control.stopExecTrace() + }() + return nil +} + func (p *profilingControl) stop(logln func(...any)) { p.stopOnce.Do(func() { if !p.enabled { diff --git a/internal/probemanager/manager.go b/internal/probemanager/manager.go index 677762b..8b15f94 100644 --- a/internal/probemanager/manager.go +++ b/internal/probemanager/manager.go @@ -137,6 +137,7 @@ func (m *Manager) Toggle(syscall string) error { return m.Attach(syscall) } +// Attach attaches enter/exit tracepoints for a registered syscall. // Attach attaches enter/exit tracepoints for a registered syscall. func (m *Manager) Attach(syscall string) error { if syscall == "" { @@ -153,25 +154,47 @@ func (m *Manager) Attach(syscall string) error { entry.attachMu.Lock() defer entry.attachMu.Unlock() + // Re-acquire the lock after the per-entry mutex to prevent races with + // concurrent Detach calls on the same syscall. + enterTP, exitTP, attacher, err := m.snapshotAttachParams(syscall, entry) + if err != nil { + return err + } + if attacher == nil { + return nil // entry was already active + } + + enterLink, exitLink, attachErr := attachPair(attacher, enterTP, exitTP) + return m.commitAttach(syscall, entry, enterLink, exitLink, attachErr) +} + +// snapshotAttachParams re-validates the entry under the manager lock and +// returns the tracepoint names and attacher needed for attachPair. It returns +// (nil attacher, nil error) when the probe is already active. +func (m *Manager) snapshotAttachParams(syscall string, entry *probeEntry) (enterTP, exitTP string, attacher Attacher, err error) { m.mu.Lock() entry, err = m.entryLocked(syscall) if err != nil { m.mu.Unlock() - return err + return "", "", nil, err } if entry.active { m.mu.Unlock() - return nil + return "", "", nil, nil } - enterTP := entry.enterTP - exitTP := entry.exitTP - attacher := m.attacher + enterTP = entry.enterTP + exitTP = entry.exitTP + attacher = m.attacher m.mu.Unlock() + return enterTP, exitTP, attacher, nil +} - enterLink, exitLink, attachErr := attachPair(attacher, enterTP, exitTP) - +// commitAttach stores the newly attached link pair in entry under the manager +// lock, recording any attach error or cleaning up on a concurrent manager close. +func (m *Manager) commitAttach(syscall string, entry *probeEntry, enterLink, exitLink Link, attachErr error) error { m.mu.Lock() defer m.mu.Unlock() + var err error entry, err = m.entryLocked(syscall) if err != nil { return errors.Join( @@ -180,13 +203,11 @@ func (m *Manager) Attach(syscall string) error { destroyLink(fmt.Sprintf("cleanup exit %s", syscall), exitLink), ) } - if attachErr != nil { entry.lastErr = attachErr entry.active = entry.enterLink != nil || entry.exitLink != nil return attachErr } - entry.enterLink = enterLink entry.exitLink = exitLink entry.lastErr = nil @@ -210,6 +231,8 @@ func (m *Manager) Detach(syscall string) error { entry.attachMu.Lock() defer entry.attachMu.Unlock() + // Re-acquire the lock after the per-entry mutex to prevent races with + // concurrent Attach calls on the same syscall. m.mu.Lock() entry, err = m.entryLocked(syscall) if err != nil { @@ -220,22 +243,31 @@ func (m *Manager) Detach(syscall string) error { exitLink := entry.exitLink m.mu.Unlock() - var errs []string - enterErr := error(nil) + enterErr, exitErr, errs := destroyLinkPair(syscall, enterLink, exitLink) + return m.commitDetach(entry, enterErr, exitErr, errs) +} + +// destroyLinkPair destroys both BPF links and collects any errors into a slice. +// It returns each link's error separately so partial-success can be recorded. +func destroyLinkPair(syscall string, enterLink, exitLink Link) (enterErr, exitErr error, errs []string) { if enterLink != nil { if err := enterLink.Destroy(); err != nil { enterErr = err errs = append(errs, fmt.Sprintf("detach enter %s: %v", syscall, err)) } } - exitErr := error(nil) if exitLink != nil { if err := exitLink.Destroy(); err != nil { exitErr = err errs = append(errs, fmt.Sprintf("detach exit %s: %v", syscall, err)) } } + return enterErr, exitErr, errs +} +// commitDetach updates entry link pointers and active flag under the manager +// lock, then returns a combined error if any link destroy failed. +func (m *Manager) commitDetach(entry *probeEntry, enterErr, exitErr error, errs []string) error { m.mu.Lock() defer m.mu.Unlock() if enterErr == nil { @@ -313,20 +345,40 @@ func (m *Manager) IsActive(syscall string) bool { } // Close detaches all registered probes and marks the manager closed. +// It returns the first detach error encountered (subsequent errors are +// recorded on the probe entry but not returned). func (m *Manager) Close() error { if m == nil { return nil } + entries, ok := m.snapshotAndMarkClosed() + if !ok { + return nil // already closed + } + + var firstErr error + for _, item := range entries { + if err := m.detachProbeEntry(item); err != nil && firstErr == nil { + firstErr = err + } + } + return firstErr +} + +// pairEntry groups a probe entry with its syscall name for use during Close. +type pairEntry struct { + syscall string + entry *probeEntry + hasLinks bool +} +// snapshotAndMarkClosed atomically marks the manager as closed and returns a +// snapshot of all probe entries. Returns (nil, false) if already closed. +func (m *Manager) snapshotAndMarkClosed() ([]pairEntry, bool) { m.mu.Lock() + defer m.mu.Unlock() if m.closed { - m.mu.Unlock() - return nil - } - type pairEntry struct { - syscall string - entry *probeEntry - hasLinks bool + return nil, false } entries := make([]pairEntry, 0, len(m.probes)) for syscall, entry := range m.probes { @@ -337,47 +389,39 @@ func (m *Manager) Close() error { }) } m.closed = true - m.mu.Unlock() + return entries, true +} - var firstErr error - for _, item := range entries { - if item.hasLinks { - item.entry.attachMu.Lock() - } - var errForSyscall error - m.mu.Lock() - enterLink := item.entry.enterLink - exitLink := item.entry.exitLink - item.entry.enterLink = nil - item.entry.exitLink = nil - item.entry.active = false - item.entry.lastErr = nil - m.mu.Unlock() +// detachProbeEntry destroys the BPF links for a single probe entry under its +// per-entry mutex, clears the link pointers, and records any error. +func (m *Manager) detachProbeEntry(item pairEntry) error { + if item.hasLinks { + item.entry.attachMu.Lock() + defer item.entry.attachMu.Unlock() + } - if enterLink != nil { - if err := enterLink.Destroy(); err != nil { - errForSyscall = err - if firstErr == nil { - firstErr = err - } - } - } - if exitLink != nil { - if err := exitLink.Destroy(); err != nil { - if errForSyscall == nil { - errForSyscall = err - } - if firstErr == nil { - firstErr = err - } - } + m.mu.Lock() + enterLink := item.entry.enterLink + exitLink := item.entry.exitLink + item.entry.enterLink = nil + item.entry.exitLink = nil + item.entry.active = false + item.entry.lastErr = nil + m.mu.Unlock() + + var errForSyscall error + if enterLink != nil { + if err := enterLink.Destroy(); err != nil { + errForSyscall = err } - m.setLastError(item.syscall, errForSyscall) - if item.hasLinks { - item.entry.attachMu.Unlock() + } + if exitLink != nil { + if err := exitLink.Destroy(); err != nil && errForSyscall == nil { + errForSyscall = err } } - return firstErr + m.setLastError(item.syscall, errForSyscall) + return errForSyscall } func (m *Manager) entryLocked(syscall string) (*probeEntry, error) { diff --git a/internal/tui/common/keys.go b/internal/tui/common/keys.go index d1f26cf..e50ee94 100644 --- a/internal/tui/common/keys.go +++ b/internal/tui/common/keys.go @@ -98,40 +98,35 @@ func (k KeyMap) DashboardStatusHelp() []key.Binding { // DashboardStatusHelpSections returns grouped bindings for dashboard status bars. func (k KeyMap) DashboardStatusHelpSections() []HelpSection { - global := []key.Binding{ + return []HelpSection{ + {Title: "Global", Bindings: k.globalStatusBindings()}, + {Title: "Dashboard", Bindings: dashboardStatusBindings(k)}, + } +} + +// globalStatusBindings returns the global key bindings shown in the status bar, +// appending the optional export binding when it has a non-empty label. +func (k KeyMap) globalStatusBindings() []key.Binding { + bindings := []key.Binding{ helpTextBinding("H", "toggle help"), - k.Tab, - k.ShiftTab, - k.One, - k.Two, - k.Three, - k.Four, - k.Five, - k.Six, - k.Seven, - k.Visualize, - k.Metric, - k.Sort, - k.ReverseSort, - k.Filter, - k.FilterUndo, - k.SelectPID, - k.SelectTID, - k.Probes, - k.Record, - k.Refresh, - k.AutoReset, - k.Quit, + k.Tab, k.ShiftTab, + k.One, k.Two, k.Three, k.Four, k.Five, k.Six, k.Seven, + k.Visualize, k.Metric, k.Sort, k.ReverseSort, + k.Filter, k.FilterUndo, + k.SelectPID, k.SelectTID, + k.Probes, k.Record, k.Refresh, k.AutoReset, k.Quit, } if help := k.Export.Help(); help.Key != "" || help.Desc != "" { - global = append(global, k.Export) + bindings = append(bindings, k.Export) } - dashboard := []key.Binding{ - k.DirGroup, - k.Visualize, - k.Metric, - k.Sort, - k.ReverseSort, + return bindings +} + +// dashboardStatusBindings returns the dashboard-specific bindings shown in +// the status bar (table navigation, stream controls, and export shortcuts). +func dashboardStatusBindings(k KeyMap) []key.Binding { + return []key.Binding{ + k.DirGroup, k.Visualize, k.Metric, k.Sort, k.ReverseSort, helpTextBinding("space", "stream pause"), helpTextBinding("enter", "selected filter"), helpTextBinding("esc", "stream undo filter"), @@ -147,11 +142,6 @@ func (k KeyMap) DashboardStatusHelpSections() []HelpSection { helpTextBinding("X", "stream export as"), helpTextBinding("E", "stream open last"), } - - return []HelpSection{ - {Title: "Global", Bindings: global}, - {Title: "Dashboard", Bindings: dashboard}, - } } // DashboardFullHelp returns grouped bindings for dashboard overlays. diff --git a/internal/tui/dashboard/bubbles.go b/internal/tui/dashboard/bubbles.go index 96ebea1..eec75a5 100644 --- a/internal/tui/dashboard/bubbles.go +++ b/internal/tui/dashboard/bubbles.go @@ -156,6 +156,9 @@ func (c *bubbleChart) SetDarkMode(isDark bool) { c.isDark = isDark } +// SetData recomputes bubble targets from data and merges them with existing +// animation state so that live updates animate smoothly. Returns true when +// at least one node has motion and a Tick should be scheduled. func (c *bubbleChart) SetData(data []bubbleDatum) bool { targets := buildBubbleTargets(data, c.Metric(), c.width, c.height) @@ -169,6 +172,23 @@ func (c *bubbleChart) SetData(data []bubbleDatum) bool { existing[node.ID] = node } + c.nodes = c.mergeTargetNodes(targets, existing) + if len(c.nodes) == 0 { + c.selected = 0 + c.animating = false + return false + } + c.selected = c.selectIndexByID(selectedID) + c.animating = c.hasMotion() + if c.animating { + c.Tick(0) + } + return c.animating +} + +// mergeTargetNodes converts target positions into live nodes, carrying over +// spring velocities and drift state from existing nodes where available. +func (c *bubbleChart) mergeTargetNodes(targets []bubbleNode, existing map[string]bubbleNode) []bubbleNode { next := make([]bubbleNode, 0, len(targets)) for _, target := range targets { node := bubbleNode{ @@ -188,25 +208,7 @@ func (c *bubbleChart) SetData(data []bubbleDatum) bool { ySpring: harmonica.NewSpring(harmonica.FPS(bubbleFPS), bubbleAngularVelocity, bubbleDamping), } if prev, ok := existing[target.ID]; ok { - node.radius = prev.radius - node.x = prev.x - node.y = prev.y - node.velocityRadius = prev.velocityRadius - node.velocityX = prev.velocityX - node.velocityY = prev.velocityY - node.driftPhase = prev.driftPhase - node.driftSpeed = prev.driftSpeed - node.driftAmpX = prev.driftAmpX - node.driftAmpY = prev.driftAmpY - // New metrics or topology can otherwise produce stale springs. - if node.radius == 0 { - node.radius = target.targetRadius - } - if node.driftSpeed == 0 { - c.initNodeDrift(&node) - } else { - c.updateNodeDriftAmplitude(&node) - } + c.inheritPrevNodeState(&node, prev, target) } else { node.radius = target.targetRadius node.x = target.targetX @@ -216,18 +218,31 @@ func (c *bubbleChart) SetData(data []bubbleDatum) bool { node.applyDrift(c.driftTime, c.width, c.height) next = append(next, node) } - c.nodes = next - if len(c.nodes) == 0 { - c.selected = 0 - c.animating = false - return false - } - c.selected = c.selectIndexByID(selectedID) - c.animating = c.hasMotion() - if c.animating { - c.Tick(0) + return next +} + +// inheritPrevNodeState copies physics and drift state from a previous node +// into node so that the transition animates rather than snapping. +func (c *bubbleChart) inheritPrevNodeState(node *bubbleNode, prev bubbleNode, target bubbleNode) { + node.radius = prev.radius + node.x = prev.x + node.y = prev.y + node.velocityRadius = prev.velocityRadius + node.velocityX = prev.velocityX + node.velocityY = prev.velocityY + node.driftPhase = prev.driftPhase + node.driftSpeed = prev.driftSpeed + node.driftAmpX = prev.driftAmpX + node.driftAmpY = prev.driftAmpY + // New metrics or topology can otherwise produce stale springs. + if node.radius == 0 { + node.radius = target.targetRadius + } + if node.driftSpeed == 0 { + c.initNodeDrift(node) + } else { + c.updateNodeDriftAmplitude(node) } - return c.animating } func (c *bubbleChart) selectIndexByID(id string) int { @@ -641,10 +656,10 @@ func renderBubbleRow(cells []bubbleCell, palette []color.Color) string { return b.String() } +// buildBubbleTargets computes initial target positions and radii for each +// bubble, then runs a short relaxation pass to reduce overlap. Returns nil +// when there is nothing to render. func buildBubbleTargets(data []bubbleDatum, metric bubbleMetric, width, height int) []bubbleNode { - if len(data) == 0 { - return nil - } if width <= 0 { width = 80 } @@ -655,6 +670,20 @@ func buildBubbleTargets(data []bubbleDatum, metric bubbleMetric, width, height i if chartHeight < 4 { chartHeight = 4 } + + filtered := filterAndSortBubbleData(data, metric) + if len(filtered) == 0 { + return nil + } + + targets := placeBubbleNodes(filtered, metric, width, chartHeight) + relaxTargets(targets, width, chartHeight) + return targets +} + +// filterAndSortBubbleData removes datums without an ID, sorts by descending +// metric value (ties broken by label), and caps the result to bubbleMaxItems. +func filterAndSortBubbleData(data []bubbleDatum, metric bubbleMetric) []bubbleDatum { filtered := make([]bubbleD