diff options
| author | Paul Buetow <paul@buetow.org> | 2026-05-13 20:04:48 +0300 |
|---|---|---|
| committer | Paul Buetow <paul@buetow.org> | 2026-05-13 20:04:48 +0300 |
| commit | 251894cf3375812564ecf28392179b395cdda9c7 (patch) | |
| tree | 83c3609ab591702e29a375923670e7622a33b5c7 | |
| parent | 78ea9e22e596255c5e23ce445d80641870674ca9 (diff) | |
refactor: break down functions exceeding 50 lines into smaller helpers
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 <noreply@anthropic.com>
| -rw-r--r-- | cmd/ioworkload/scenario_rename.go | 54 | ||||
| -rw-r--r-- | internal/benchutil/eventmix.go | 139 | ||||
| -rw-r--r-- | internal/eventloop_exit.go | 38 | ||||
| -rw-r--r-- | internal/eventloop_runtime.go | 33 | ||||
| -rw-r--r-- | internal/export/snapshot_csv.go | 114 | ||||
| -rw-r--r-- | internal/flags/flags.go | 57 | ||||
| -rw-r--r-- | internal/generate/bpfhandler.go | 120 | ||||
| -rw-r--r-- | internal/generate/format.go | 84 | ||||
| -rw-r--r-- | internal/ior.go | 101 | ||||
| -rw-r--r-- | internal/ior_parquet_sink.go | 107 | ||||
| -rw-r--r-- | internal/ior_profiling.go | 81 | ||||
| -rw-r--r-- | internal/probemanager/manager.go | 154 | ||||
| -rw-r--r-- | internal/tui/common/keys.go | 60 | ||||
| -rw-r--r-- | internal/tui/dashboard/bubbles.go | 121 | ||||
| -rw-r--r-- | internal/tui/dashboard/histogram.go | 63 | ||||
| -rw-r--r-- | internal/tui/dashboard/icicle.go | 88 | ||||
| -rw-r--r-- | internal/tui/dashboard/treemap.go | 91 | ||||
| -rw-r--r-- | internal/tui/eventstream/export.go | 129 | ||||
| -rw-r--r-- | internal/tui/eventstream/model.go | 390 | ||||
| -rw-r--r-- | internal/tui/export/model.go | 69 | ||||
| -rw-r--r-- | internal/tui/probes/model.go | 172 | ||||
| -rw-r--r-- | internal/tui/tracefilter/model.go | 233 |
22 files changed, 1475 insertions, 1023 deletions
diff --git a/cmd/ioworkload/scenario_rename.go b/cmd/ioworkload/scenario_rename.go index 685157b..cdbf758 100644 --- a/cmd/ioworkload/scenario_rename.go +++ b/cmd/ioworkload/scenario_rename.go @@ -127,19 +127,34 @@ func renameRenameat2() error { return fmt.Errorf("close: %w", err) } + errno, err := doRenameat2(dir, oldName, newName, 0) + if err != nil { + return err + } + if errno != 0 { + return fmt.Errorf("renameat2: %w", errno) + } + return nil +} + +// doRenameat2 opens dir as a directory fd, builds name byte pointers, calls +// renameat2(2) with the given flags, and returns the raw errno. Zero means +// success. It is shared by renameRenameat2 and renameNoreplace to avoid +// duplicating the dir-open + byte-pointer + syscall6 boilerplate. +func doRenameat2(dir, oldName, newName string, flags uintptr) (syscall.Errno, error) { dirFD, err := syscall.Open(dir, syscall.O_RDONLY|syscall.O_DIRECTORY, 0) if err != nil { - return fmt.Errorf("open dir: %w", err) + return 0, fmt.Errorf("open dir: %w", err) } defer syscall.Close(dirFD) oldBytes, err := syscall.BytePtrFromString(oldName) if err != nil { - return fmt.Errorf("old name bytes: %w", err) + return 0, fmt.Errorf("old name bytes: %w", err) } newBytes, err := syscall.BytePtrFromString(newName) if err != nil { - return fmt.Errorf("new name bytes: %w", err) + return 0, fmt.Errorf("new name bytes: %w", err) } _, _, errno := syscall.Syscall6( @@ -148,15 +163,12 @@ func renameRenameat2() error { uintptr(unsafe.Pointer(oldBytes)), uintptr(dirFD), uintptr(unsafe.Pointer(newBytes)), - 0, // flags=0: plain rename + flags, 0, ) runtime.KeepAlive(oldBytes) runtime.KeepAlive(newBytes) - if errno != 0 { - return fmt.Errorf("renameat2: %w", errno) - } - return nil + return errno, nil } // renameEnoent attempts to rename a nonexistent file via raw SYS_RENAME. @@ -220,32 +232,10 @@ func renameNoreplace() error { } } - dirFD, err := syscall.Open(dir, syscall.O_RDONLY|syscall.O_DIRECTORY, 0) - if err != nil { - return fmt.Errorf("open dir: %w", err) - } - defer syscall.Close(dirFD) - - srcBytes, err := syscall.BytePtrFromString(srcName) - if err != nil { - return fmt.Errorf("src name bytes: %w", err) - } - dstBytes, err := syscall.BytePtrFromString(dstName) + errno, err := doRenameat2(dir, srcName, dstName, renameNoreplaceFlag) if err != nil { - return fmt.Errorf("dst name bytes: %w", err) + return err } - - _, _, errno := syscall.Syscall6( - sysRenameat2, - uintptr(dirFD), - uintptr(unsafe.Pointer(srcBytes)), - uintptr(dirFD), - uintptr(unsafe.Pointer(dstBytes)), - renameNoreplaceFlag, - 0, - ) - runtime.KeepAlive(srcBytes) - runtime.KeepAlive(dstBytes) if errno == 0 { return fmt.Errorf("expected EEXIST, but renameat2 NOREPLACE succeeded") } 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 "" } - |
