/
rustwizard
/
pgtrace
Обзор
Документация
Войти
/
rustwizard
/
pgtrace
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
internal/bpf/offcpu.go
296 строк
7 KB
Rust Wizard
fix(offcpu): capture stack at block time via kprobe:schedule
08 авг 2026, 10:55
Верифицирован
08 авг 2026, 10:55
d84b1f0
Код
Авторство
О чём код?
// Package bpf loads the pgtrace eBPF programs, attaches them to syscall // tracepoints and streams syscall latency events from the kernel ring buffer. package bpf import ( "bytes" "context" "encoding/binary" "errors" "fmt" "log/slog" "os" "strconv" "strings" "github.com/cilium/ebpf/link" "github.com/cilium/ebpf/ringbuf" ) // WaitType names, keyed by the stack classifier output. Kept in sync with the // enum in bpf/offcpu.c. const ( WaitLwlock = "lwlock" WaitLock = "lock" WaitIO = "io" WaitWalWrite = "wal_write" WaitNetwork = "network" WaitIdle = "idle" WaitOther = "other" ) // OffCPUConfig configures the kernel-side off-CPU filter. type OffCPUConfig struct { // ThresholdNs: only off-CPU durations longer than this are reported. ThresholdNs uint64 // Comm: only processes with this comm are traced. Comm string } // OffCPUEvent is a single off-CPU observation for a postgres process. type OffCPUEvent struct { PID uint32 WaitType string DurationNs uint64 Ts uint64 // kernel monotonic timestamp, ns } // OffCPUTracer owns the loaded off-CPU BPF objects, the kprobe + tracepoint // links and the ring buffer reader. type OffCPUTracer struct { objs offcpuObjects links []link.Link reader *ringbuf.Reader kallsyms *kallsyms } // kallsyms resolves kernel addresses to symbol names. Kernel addresses from // bpf_get_stack may point inside a function, so symbols are kept sorted by // address and resolved to the nearest symbol at or below the address. type kallsyms struct { addrs []uint64 syms []string } func newKallsyms() (*kallsyms, error) { data, err := os.ReadFile("/proc/kallsyms") if err != nil { return nil, fmt.Errorf("read /proc/kallsyms: %w", err) } ks := &kallsyms{} for line := range strings.SplitSeq(string(data), "\n") { fields := strings.Fields(line) if len(fields) < 3 { continue } addr, err := strconv.ParseUint(fields[0], 16, 64) if err != nil || addr == 0 { continue } ks.addrs = append(ks.addrs, addr) ks.syms = append(ks.syms, fields[2]) } return ks, nil } // resolve returns the symbol containing addr (nearest symbol at or below), or // a hex fallback when nothing is found. func (k *kallsyms) resolve(addr uint64) string { lo, hi := 0, len(k.addrs) for lo < hi { mid := (lo + hi) / 2 if k.addrs[mid] <= addr { lo = mid + 1 } else { hi = mid } } if lo > 0 { return k.syms[lo-1] } return fmt.Sprintf("0x%x", addr) } // stackMatches reports whether any frame in stack matches one of the patterns. func stackMatches(stack []string, patterns ...string) bool { for _, fn := range stack { for _, p := range patterns { if strings.Contains(fn, p) { return true } } } return false } // classify maps a kernel stack (top frame = leaf) to a wait type by scanning // for known kernel functions. // // Postgres lightweight/heavyweight locks both block in the futex syscall, so // kernel stacks alone cannot tell them apart; both are reported as "lock". // Distinguishing LWLock from heavyweight lock would need user-space stacks. func classify(stack []string) string { switch { case stackMatches(stack, "io_schedule", "io_schedule_timeout", "lock_page", "wait_on_page_bit", "wait_on_buffer", "submit_bio", "blk_mq", "blk_queue"): return WaitIO case stackMatches(stack, "futex_wait", "do_futex", "futex_wait_queue_me"): return WaitLock case stackMatches(stack, "ep_poll", "ep_wait", "epoll_wait", "do_epoll_wait", "sock_poll", "tcp_poll"): return WaitIdle case stackMatches(stack, "tcp_", "inet_", "sock_", "sk_stream", "unix_stream", "sk_wait"): return WaitNetwork case stackMatches(stack, "xlog", "wal_writer", "XLogWrite"): return WaitWalWrite default: // schedule, schedule_hrtimeout_range, do_nanosleep, ... — sleeping. return WaitOther } } // NewOffCPUTracer loads the off-CPU BPF objects, writes the config, attaches // the sched_switch tracepoint and opens the ring buffer reader. func NewOffCPUTracer(cfg OffCPUConfig) (t *OffCPUTracer, err error) { var objs offcpuObjects if err := loadOffcpuObjects(&objs, nil); err != nil { return nil, fmt.Errorf("load offcpu objects: %w", err) } t = &OffCPUTracer{objs: objs} defer func() { if t.reader == nil { if cerr := t.Close(); cerr != nil { err = errors.Join(err, fmt.Errorf("cleanup on failed init: %w", cerr)) } } }() kcfg := struct { ThresholdNs uint64 Comm [CommLen]byte }{ThresholdNs: cfg.ThresholdNs} copy(kcfg.Comm[:], cfg.Comm) if err := objs.OffcpuConfigMap.Put(uint32(0), kcfg); err != nil { return nil, fmt.Errorf("write offcpu config: %w", err) } l, err := link.Tracepoint("sched", "sched_switch", objs.TraceSchedSwitch, nil) if err != nil { return nil, fmt.Errorf("attach sched_switch: %w", err) } t.links = append(t.links, l) kl, err := link.Kprobe("schedule", objs.TraceSchedule, nil) if err != nil { return nil, fmt.Errorf("attach kprobe schedule: %w", err) } t.links = append(t.links, kl) reader, err := ringbuf.NewReader(objs.OffcpuEvents) if err != nil { return nil, fmt.Errorf("open offcpu ring buffer: %w", err) } t.reader = reader ks, err := newKallsyms() if err != nil { slog.Warn("offcpu: kallsyms unavailable, wait types limited", "err", err) } t.kallsyms = ks return t, nil } // Run streams off-CPU events until ctx is cancelled or the tracer is closed. // The returned channel is closed when the reader stops. func (t *OffCPUTracer) Run(ctx context.Context) <-chan OffCPUEvent { out := make(chan OffCPUEvent, 256) go func() { defer close(out) go func() { <-ctx.Done() if err := t.reader.Close(); err != nil { slog.Warn("close offcpu ring buffer on shutdown", "err", err) } }() var ev offcpuOffcpuEvent for { rec, err := t.reader.Read() if err != nil { if !errors.Is(err, ringbuf.ErrClosed) { slog.Error("offcpu ring buffer read", "err", err) } return } if err := binary.Read(bytes.NewReader(rec.RawSample), binary.LittleEndian, &ev); err != nil { slog.Error("decode offcpu event", "err", err) continue } out <- t.convert(ev) } }() return out } // convert decodes a kernel event into an OffCPUEvent, resolving the stack to // a wait type when symbolization is available. func (t *OffCPUTracer) convert(ev offcpuOffcpuEvent) OffCPUEvent { out := OffCPUEvent{ PID: ev.Pid, WaitType: WaitOther, DurationNs: ev.DurationNs, Ts: ev.Ts, } if t.kallsyms == nil { return out } stack := make([]string, 0, len(ev.Stack)) for _, addr := range ev.Stack { if addr == 0 { continue } stack = append(stack, t.kallsyms.resolve(addr)) } out.WaitType = classify(stack) return out } // Close detaches the tracepoint and releases the BPF objects. func (t *OffCPUTracer) Close() error { var errs []error if t.reader != nil { if err := t.reader.Close(); err != nil { errs = append(errs, err) } } for _, l := range t.links { if err := l.Close(); err != nil { errs = append(errs, err) } } if err := t.objs.Close(); err != nil { errs = append(errs, err) } return errors.Join(errs...) }