From 639148c658017f5df91dd6060d5c0113ef4509d0 Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Tue, 29 Sep 2026 08:53:16 -0400 Subject: [PATCH 1/2] feat: add bounded CPU profiling --- PROFILING.md | 79 +++++++ cmd/beacon/main.go | 5 + cmd/beacon/profiling.go | 31 +++ env.example | 4 + internal/profiling/cpu_other.go | 8 + internal/profiling/cpu_unix.go | 19 ++ internal/profiling/profiling.go | 302 ++++++++++++++++++++++++++ internal/profiling/profiling_test.go | 305 +++++++++++++++++++++++++++ 8 files changed, 753 insertions(+) create mode 100644 PROFILING.md create mode 100644 cmd/beacon/profiling.go create mode 100644 internal/profiling/cpu_other.go create mode 100644 internal/profiling/cpu_unix.go create mode 100644 internal/profiling/profiling.go create mode 100644 internal/profiling/profiling_test.go diff --git a/PROFILING.md b/PROFILING.md new file mode 100644 index 00000000..dc1a7a14 --- /dev/null +++ b/PROFILING.md @@ -0,0 +1,79 @@ +# Production CPU profiling + +CPU profiling is supported on Linux and macOS and is disabled by default. It writes private files inside the container; +there is no HTTP endpoint or published profiling port. Profiling failures stop the +recorder and log a warning without stopping ingest. + +## Enable a bounded session + +Create a dedicated directory writable by the container user with mode `0700`. +Mount it into the app container and set both environment variables: + +```yaml +services: + app: + environment: + BEACON_CPU_PROFILE_DIR: /profiles + BEACON_CPU_PROFILE_UNTIL: "2026-10-02T12:00:00Z" # Replace with your fixed deadline. + volumes: + - ./profiles:/profiles +``` + +Choose an RFC3339 deadline no more than 72 hours in the future. The same deadline +remains in effect after restarts; expired settings are inactive. Do not generate a +fresh deadline automatically on each startup. Restarting the app is required to +change these settings. Preserve the rest of the deployment configuration. + +Start with a deadline about one minute away to assess the first capture's overhead. +Check process CPU, packet freshness, queue drops and database timeouts against a +comparable period with profiling disabled. After reviewing that capture, an operator +can enable a longer session. Profiling adds overhead while a capture is active. + +## Capture behavior + +- One 30-second sample immediately, then every 30 minutes. +- Route reconfirmation requests an additional sample when the maintenance task starts. + This includes the retention step before route validation. A five-minute cooldown + between capture starts prevents overlap and repeated triggers from increasing load. +- Background task stacks carry a `task` label while profiling is enabled. +- Shutdown or expiry stops the active sample and saves the shorter profile. +- Each profile is limited to 8 MiB. The dedicated directory is limited to 256 MiB + and 512 files, including metadata and files left by interrupted runs. The recorder + reserves space for a full capture before starting and stops when a limit is reached. + Files are never automatically deleted. Existing unrelated files consume the budget. +- Profiles and metadata are written with mode `0600`. A `.partial` file indicates + an interrupted capture and is not a completed profile. + +Only one Beacon process should write to a profiling directory. Keep the mount +private and out of web roots, backups intended for public download, and source control. +Do not run another CPU profiler in the same process during a session. + +## Interpret the output + +Each `.pprof` has a `.pprof.json` sidecar containing UTC start/end times, the trigger, +Go version, process CPU counters (Linux/macOS), goroutine count, and database-pool +counters before and after the capture. Pool acquire duration and counts are cumulative; +use differences between the two snapshots. They do not contain SQL text or credentials. + +Copy completed files privately off the host. Record the container's image revision +and digest alongside them. On a workstation with Go installed: + +```sh +go tool pprof -top capture.pprof +go tool pprof -top -cum capture.pprof +go tool pprof -tags capture.pprof +``` + +CPU samples show work inside Beacon, not CPU used by PostgreSQL or time waiting on +locks. Correlate each UTC interval with separately collected host/container CPU and +I/O, PostgreSQL wait states, observation throughput, and Beacon timeout/drop logs. +Do not enable SQL text logging or run full-table counts for this purpose. The recorder +makes no diagnostic database queries and does not change the ingest path. + +Review captures from quiet periods, bursts, and maintenance separately before +combining them. Additional maintenance samples deliberately bias an aggregate profile. +Thirty-second windows can miss brief stalls; the absence of a stack is not proof that +it never consumes CPU. + +At the deadline, check for `CPU profiling stopped` in the app log. Export the files, +then remove the environment variables and mount during the next planned deployment. diff --git a/cmd/beacon/main.go b/cmd/beacon/main.go index 8e322e6e..9e1b0a90 100644 --- a/cmd/beacon/main.go +++ b/cmd/beacon/main.go @@ -301,6 +301,11 @@ func main() { } tasks = append(tasks, background.ObserverCleanupTask(coalescer, resolved.ObserverDeleteAfter, resolved.CleanupInterval, onDelete)) } + profiles := configureProfiling(ctx, pool) + defer profiles.Stop() + for i := range tasks { + tasks[i].Run = profiles.WrapTask(tasks[i].Name, tasks[i].Run) + } scheduler := background.New(tasks) go scheduler.Start(ctx) diff --git a/cmd/beacon/profiling.go b/cmd/beacon/profiling.go new file mode 100644 index 00000000..d797fe21 --- /dev/null +++ b/cmd/beacon/profiling.go @@ -0,0 +1,31 @@ +// Copyright 2026 Beacon Contributors +// SPDX-License-Identifier: AGPL-3.0-or-later + +package main + +import ( + "context" + "log/slog" + "os" + + "github.com/MeshCore-Beacon/beacon-server/internal/profiling" + "github.com/jackc/pgx/v5/pgxpool" +) + +func configureProfiling(ctx context.Context, pool *pgxpool.Pool) *profiling.Recorder { + r, err := profiling.Start(ctx, os.Getenv("BEACON_CPU_PROFILE_DIR"), os.Getenv("BEACON_CPU_PROFILE_UNTIL"), func() any { + s := pool.Stat() + return struct { + Acquired, Idle, Total, Max int32 + Acquires, EmptyAcquires, CanceledAcquires int64 + AcquireDurationNS int64 + }{ + s.AcquiredConns(), s.IdleConns(), s.TotalConns(), s.MaxConns(), + s.AcquireCount(), s.EmptyAcquireCount(), s.CanceledAcquireCount(), int64(s.AcquireDuration()), + } + }) + if err != nil { + slog.Warn("CPU profiling unavailable", "component", "profiling", "error", err) + } + return r +} diff --git a/env.example b/env.example index 97d7cbc4..b918e6e3 100644 --- a/env.example +++ b/env.example @@ -21,3 +21,7 @@ MQTT_BROKER_1_PASSWORD= MQTT_BROKER_2_URL=wss://mqtt2.meshcore.ca:443 MQTT_BROKER_2_USERNAME= MQTT_BROKER_2_PASSWORD= + +# Optional private CPU captures. Set both; the deadline must be within 72 hours. +# BEACON_CPU_PROFILE_DIR=/profiles +# BEACON_CPU_PROFILE_UNTIL=2026-10-02T12:00:00Z diff --git a/internal/profiling/cpu_other.go b/internal/profiling/cpu_other.go new file mode 100644 index 00000000..9c94f146 --- /dev/null +++ b/internal/profiling/cpu_other.go @@ -0,0 +1,8 @@ +// Copyright 2026 Beacon Contributors +// SPDX-License-Identifier: AGPL-3.0-or-later + +//go:build !linux && !darwin + +package profiling + +func cpuUsage() *processCPU { return nil } diff --git a/internal/profiling/cpu_unix.go b/internal/profiling/cpu_unix.go new file mode 100644 index 00000000..77da9f73 --- /dev/null +++ b/internal/profiling/cpu_unix.go @@ -0,0 +1,19 @@ +// Copyright 2026 Beacon Contributors +// SPDX-License-Identifier: AGPL-3.0-or-later + +//go:build linux || darwin + +package profiling + +import "syscall" + +func cpuUsage() *processCPU { + var usage syscall.Rusage + if syscall.Getrusage(syscall.RUSAGE_SELF, &usage) != nil { + return nil + } + return &processCPU{ + UserSeconds: float64(usage.Utime.Sec) + float64(usage.Utime.Usec)/1e6, + SystemSeconds: float64(usage.Stime.Sec) + float64(usage.Stime.Usec)/1e6, + } +} diff --git a/internal/profiling/profiling.go b/internal/profiling/profiling.go new file mode 100644 index 00000000..8d51cda9 --- /dev/null +++ b/internal/profiling/profiling.go @@ -0,0 +1,302 @@ +// Copyright 2026 Beacon Contributors +// SPDX-License-Identifier: AGPL-3.0-or-later + +// Package profiling records bounded, opt-in CPU samples to private files. +package profiling + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "io" + "log/slog" + "os" + "path/filepath" + "runtime" + "runtime/pprof" + "time" +) + +const ( + captureDuration = 30 * time.Second + captureInterval = 30 * time.Minute + captureCooldown = 5 * time.Minute + maxProfileBytes = 8 << 20 + maxMetadataBytes = 16 << 10 + maxDirectoryBytes = 256 << 20 + maxDirectoryFiles = 512 +) + +type Recorder struct { + root *os.Root + snapshot func() any + requests chan string + cancel context.CancelFunc + done chan struct{} +} + +type snapshot struct { + ProcessCPU *processCPU `json:",omitempty"` + Goroutines int + Pool any `json:",omitempty"` +} + +type metadata struct { + Reason string + StartedAt time.Time + FinishedAt time.Time + GoVersion string + Before snapshot + After snapshot +} + +// Start leaves expired settings inactive so restarts cannot extend a capture window. +func Start(parent context.Context, dir, until string, poolSnapshot func() any) (*Recorder, error) { + if dir == "" && until == "" { + return nil, nil + } + if dir == "" || until == "" { + return nil, errors.New("profiling requires both directory and deadline") + } + deadline, err := time.Parse(time.RFC3339Nano, until) + if err != nil { + return nil, fmt.Errorf("invalid profiling deadline: %w", err) + } + if !deadline.After(time.Now()) { + return nil, nil + } + if time.Until(deadline) > 72*time.Hour { + return nil, errors.New("profiling deadline must be within 72 hours") + } + if runtime.GOOS != "linux" && runtime.GOOS != "darwin" { + return nil, errors.New("private CPU profiling is supported on Linux and macOS") + } + if !filepath.IsAbs(dir) { + return nil, errors.New("profiling directory must be absolute") + } + if err := os.MkdirAll(dir, 0700); err != nil { + return nil, err + } + info, err := os.Lstat(dir) + if err != nil { + return nil, err + } + if !info.IsDir() || info.Mode().Perm()&0077 != 0 { + return nil, errors.New("profiling directory must be private (0700) and not a symlink") + } + root, err := os.OpenRoot(dir) + if err != nil { + return nil, err + } + ctx, cancel := context.WithDeadline(parent, deadline) + r := &Recorder{root: root, snapshot: poolSnapshot, requests: make(chan string, 1), cancel: cancel, done: make(chan struct{})} + slog.Info("CPU profiling enabled", "component", "profiling", "until", deadline.UTC()) + go r.run(ctx) + return r, nil +} + +func (r *Recorder) Stop() { + if r == nil { + return + } + r.cancel() + <-r.done +} + +// Request never holds up maintenance or queues a backlog of captures. +func (r *Recorder) Request(reason string) { + if r == nil || reason != "reconfirm" { + return + } + select { + case <-r.done: + return + default: + } + select { + case r.requests <- reason: + default: + } +} + +func (r *Recorder) WrapTask(name string, run func(context.Context) error) func(context.Context) error { + if r == nil { + return run + } + return func(ctx context.Context) (err error) { + select { + case <-r.done: + return run(ctx) + default: + } + if name == "reconfirm" { + r.Request(name) + } + pprof.Do(ctx, pprof.Labels("task", name), func(ctx context.Context) { err = run(ctx) }) + return err + } +} + +func (r *Recorder) run(ctx context.Context) { + defer close(r.done) + defer r.root.Close() + defer r.cancel() + defer slog.Info("CPU profiling stopped", "component", "profiling") + if err := schedule(ctx, r.requests, r.capture); err != nil { + slog.Warn("CPU profiling stopped after capture failure", "component", "profiling", "error", err) + } +} + +func schedule(ctx context.Context, requests <-chan string, capture func(context.Context, string) error) error { + ticker := time.NewTicker(captureInterval) + defer ticker.Stop() + reason := "periodic" + var last time.Time + for { + if ctx.Err() != nil { + return nil + } + if last.IsZero() || time.Since(last) >= captureCooldown { + last = time.Now() + if err := capture(ctx, reason); err != nil { + return err + } + } + select { + case <-ctx.Done(): + return nil + case <-ticker.C: + reason = "periodic" + case reason = <-requests: + } + } +} + +func (r *Recorder) takeSnapshot() snapshot { + s := snapshot{Goroutines: runtime.NumGoroutine()} + s.ProcessCPU = cpuUsage() + if r.snapshot != nil { + s.Pool = r.snapshot() + } + return s +} + +func (r *Recorder) checkBudget() error { + d, err := r.root.Open(".") + if err != nil { + return err + } + entries, err := d.ReadDir(maxDirectoryFiles + 1) + d.Close() + if err != nil && err != io.EOF { + return err + } + if len(entries) > maxDirectoryFiles-2 { + return errors.New("profiling file budget exhausted") + } + var total int64 + for _, entry := range entries { + info, err := entry.Info() + if err != nil { + return err + } + if !info.Mode().IsRegular() { + return errors.New("profiling directory contains a non-regular file") + } + total += info.Size() + } + if total > maxDirectoryBytes-maxProfileBytes-maxMetadataBytes { + return errors.New("profiling storage budget exhausted") + } + return nil +} + +func (r *Recorder) capture(ctx context.Context, reason string) error { + if err := r.checkBudget(); err != nil { + return err + } + started := time.Now().UTC() + name := started.Format("20060102T150405.000000000Z") + ".pprof" + partial := name + ".partial" + f, err := r.root.OpenFile(partial, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0600) + if err != nil { + return err + } + defer r.root.Remove(partial) + w := &limitedWriter{writer: f, remaining: maxProfileBytes} + m := metadata{Reason: reason, StartedAt: started, GoVersion: runtime.Version(), Before: r.takeSnapshot()} + if err := pprof.StartCPUProfile(w); err != nil { + f.Close() + return err + } + slog.Info("CPU profiling capture started", "component", "profiling", "file", name, "reason", reason) + timer := time.NewTimer(captureDuration) + select { + case <-ctx.Done(): + case <-timer.C: + } + timer.Stop() + pprof.StopCPUProfile() + m.FinishedAt = time.Now().UTC() + m.After = r.takeSnapshot() + closeErr := f.Close() + if w.err != nil { + return w.err + } + if closeErr != nil { + return closeErr + } + if err := r.root.Rename(partial, name); err != nil { + return err + } + data, err := json.MarshalIndent(m, "", " ") + if err != nil { + return err + } + if len(data) > maxMetadataBytes { + return errors.New("profiling metadata budget exhausted") + } + meta, err := r.root.OpenFile(name+".json", os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0600) + if err != nil { + return err + } + _, err = meta.Write(data) + closeErr = meta.Close() + if err != nil { + return err + } + if closeErr != nil { + return closeErr + } + slog.Info("CPU profile saved", "component", "profiling", "file", name, "reason", reason, "duration", m.FinishedAt.Sub(started)) + return nil +} + +type limitedWriter struct { + writer io.Writer + remaining int + err error +} + +func (w *limitedWriter) Write(p []byte) (int, error) { + if w.err != nil { + return 0, w.err + } + if len(p) > w.remaining { + w.err = errors.New("CPU profile exceeded size limit") + return 0, w.err + } + n, err := w.writer.Write(p) + w.remaining -= n + if err == nil && n != len(p) { + err = io.ErrShortWrite + } + w.err = err + return n, err +} + +type processCPU struct { + UserSeconds float64 + SystemSeconds float64 +} diff --git a/internal/profiling/profiling_test.go b/internal/profiling/profiling_test.go new file mode 100644 index 00000000..88afb3f1 --- /dev/null +++ b/internal/profiling/profiling_test.go @@ -0,0 +1,305 @@ +// Copyright 2026 Beacon Contributors +// SPDX-License-Identifier: AGPL-3.0-or-later + +package profiling + +import ( + "bytes" + "compress/gzip" + "context" + "encoding/json" + "errors" + "io" + "log/slog" + "os" + "path/filepath" + "reflect" + "runtime" + "runtime/pprof" + "testing" + "testing/synctest" + "time" +) + +func TestDisabledAndExpired(t *testing.T) { + for _, tc := range []struct{ dir, until string }{ + {}, {filepath.Join(t.TempDir(), "unused"), time.Now().Add(-time.Hour).Format(time.RFC3339)}, + } { + r, err := Start(context.Background(), tc.dir, tc.until, nil) + if err != nil || r != nil { + t.Fatalf("got %v, %v", r, err) + } + if tc.dir != "" { + if _, err := os.Stat(tc.dir); !os.IsNotExist(err) { + t.Fatalf("expired capture created directory: %v", err) + } + } + } +} + +func TestRejectUnsafeConfiguration(t *testing.T) { + for _, tc := range []struct{ name, dir, until string }{ + {"missing deadline", t.TempDir(), ""}, + {"missing directory", "", time.Now().Add(time.Hour).Format(time.RFC3339)}, + {"invalid deadline", t.TempDir(), "tomorrow"}, + {"unbounded window", t.TempDir(), time.Now().Add(73 * time.Hour).Format(time.RFC3339)}, + {"relative directory", "profiles", time.Now().Add(time.Hour).Format(time.RFC3339)}, + } { + t.Run(tc.name, func(t *testing.T) { + r, err := Start(context.Background(), tc.dir, tc.until, nil) + if r != nil { + r.Stop() + } + if err == nil { + t.Fatal("accepted unsafe configuration") + } + }) + } + dir := t.TempDir() + if err := os.Chmod(dir, 0755); err != nil { + t.Fatal(err) + } + r, err := Start(context.Background(), dir, time.Now().Add(time.Hour).Format(time.RFC3339), nil) + if r != nil { + r.Stop() + } + if err == nil { + t.Fatal("accepted public directory") + } +} + +func TestDeadlineProducesPrivateProfileAndMetadata(t *testing.T) { + requireSupportedPlatform(t) + dir := filepath.Join(t.TempDir(), "profiles") + until := time.Now().Add(500 * time.Millisecond) + r, err := Start(context.Background(), dir, until.Format(time.RFC3339Nano), func() any { return map[string]int{"acquired": 2} }) + if err != nil { + t.Fatal(err) + } + select { + case <-r.done: + case <-time.After(3 * time.Second): + r.Stop() + t.Fatal("did not stop at deadline") + } + files, err := filepath.Glob(filepath.Join(dir, "*.pprof")) + if err != nil || len(files) != 1 { + t.Fatalf("profiles %v, %v", files, err) + } + b, err := os.ReadFile(files[0]) + if err != nil { + t.Fatal(err) + } + z, err := gzip.NewReader(bytes.NewReader(b)) + if err != nil { + t.Fatal(err) + } + data, err := io.ReadAll(z) + z.Close() + if err != nil || len(data) == 0 { + t.Fatalf("invalid profile: %v", err) + } + for _, name := range []string{files[0], files[0] + ".json"} { + info, err := os.Stat(name) + if err != nil { + t.Fatal(err) + } + if info.Mode().Perm() != 0600 { + t.Fatalf("permissions %v", info.Mode()) + } + } + var m struct { + Reason string + StartedAt, FinishedAt time.Time + Before, After struct{ Pool map[string]int } + } + b, err = os.ReadFile(files[0] + ".json") + if err != nil { + t.Fatal(err) + } + if err := json.Unmarshal(b, &m); err != nil { + t.Fatal(err) + } + if m.Reason != "periodic" || !m.FinishedAt.After(m.StartedAt) || m.Before.Pool["acquired"] != 2 || m.After.Pool["acquired"] != 2 { + t.Fatalf("bad metadata %+v", m) + } + r.Request("reconfirm") + r.Stop() + files, _ = filepath.Glob(filepath.Join(dir, "*.pprof")) + if len(files) != 1 { + t.Fatal("capture after expiry") + } +} + +func TestStopReleasesProfiler(t *testing.T) { + requireSupportedPlatform(t) + active := make(chan struct{}, 2) + old := slog.Default() + slog.SetDefault(slog.New(captureObserver{Handler: slog.NewTextHandler(io.Discard, nil), active: active})) + defer slog.SetDefault(old) + dir := filepath.Join(t.TempDir(), "profiles") + r, err := Start(context.Background(), dir, time.Now().Add(time.Hour).Format(time.RFC3339), nil) + if err != nil { + t.Fatal(err) + } + select { + case <-active: + case <-time.After(3 * time.Second): + r.Stop() + t.Fatal("capture did not start") + } + r.Stop() + r.Stop() + first, _ := filepath.Glob(filepath.Join(dir, "*.pprof")) + if len(first) != 1 { + t.Fatal("active capture was not saved on shutdown") + } + b, err := os.ReadFile(first[0]) + if err != nil { + t.Fatal(err) + } + z, err := gzip.NewReader(bytes.NewReader(b)) + if err != nil { + t.Fatal(err) + } + if _, err := io.Copy(io.Discard, z); err != nil { + t.Fatal(err) + } + z.Close() + dir2 := filepath.Join(t.TempDir(), "profiles") + r2, err := Start(context.Background(), dir2, time.Now().Add(400*time.Millisecond).Format(time.RFC3339Nano), nil) + if err != nil { + t.Fatal(err) + } + <-r2.done + files, _ := filepath.Glob(filepath.Join(dir2, "*.pprof")) + if len(files) != 1 { + t.Fatal("profiler was not released") + } +} + +func TestProfileWriteLimit(t *testing.T) { + var b bytes.Buffer + w := &limitedWriter{writer: &b, remaining: 4} + if _, err := w.Write([]byte("12345")); err == nil { + t.Fatal("accepted oversized profile") + } + if b.Len() > 4 { + t.Fatal("exceeded file budget") + } + if w.err == nil { + t.Fatal("write failure lost") + } +} + +func TestDirectoryBudgetSurvivesRestart(t *testing.T) { + requireSupportedPlatform(t) + dir := filepath.Join(t.TempDir(), "profiles") + if err := os.Mkdir(dir, 0700); err != nil { + t.Fatal(err) + } + f, err := os.Create(filepath.Join(dir, "old.pprof")) + if err != nil { + t.Fatal(err) + } + if err := f.Truncate(256 << 20); err != nil { + t.Fatal(err) + } + f.Close() + r, err := Start(context.Background(), dir, time.Now().Add(time.Hour).Format(time.RFC3339), nil) + if err != nil { + t.Fatal(err) + } + select { + case <-r.done: + case <-time.After(2 * time.Second): + r.Stop() + t.Fatal("budget did not stop recording") + } + entries, _ := os.ReadDir(dir) + if len(entries) != 1 { + t.Fatal("wrote past directory budget") + } +} + +func TestScheduleBoundsExtraCaptures(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + requests := make(chan string, 1) + type event struct { + at time.Duration + reason string + } + var events []event + start := time.Now() + done := make(chan error, 1) + go func() { + done <- schedule(ctx, requests, func(ctx context.Context, reason string) error { + events = append(events, event{time.Since(start), reason}) + time.Sleep(30 * time.Second) + return nil + }) + }() + synctest.Wait() + time.Sleep(time.Minute) + requests <- "reconfirm" + synctest.Wait() + if len(events) != 1 { + t.Fatalf("ignored cooldown: %+v", events) + } + time.Sleep(4 * time.Minute) + requests <- "reconfirm" + synctest.Wait() + time.Sleep(56 * time.Minute) + cancel() + synctest.Wait() + if err := <-done; err != nil { + t.Fatal(err) + } + want := []event{{0, "periodic"}, {5 * time.Minute, "reconfirm"}, {30 * time.Minute, "periodic"}, {60 * time.Minute, "periodic"}} + if !reflect.DeepEqual(events, want) { + t.Fatalf("captures %+v, want %+v", events, want) + } + }) +} + +func TestWrappedTaskPreservesContextAndError(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + want := errors.New("task failed") + r := &Recorder{requests: make(chan string, 1), done: make(chan struct{})} + run := r.WrapTask("reconfirm", func(got context.Context) error { + if got.Err() != context.Canceled { + t.Fatal("lost cancellation") + } + if label, ok := pprof.Label(got, "task"); !ok || label != "reconfirm" { + t.Fatal("missing task label") + } + return want + }) + for i := 0; i < 10; i++ { + if err := run(ctx); err != want { + t.Fatalf("lost task error: %v", err) + } + } +} + +func requireSupportedPlatform(t *testing.T) { + t.Helper() + if runtime.GOOS != "linux" && runtime.GOOS != "darwin" { + t.Skip("private profiling requires Unix permissions") + } +} + +type captureObserver struct { + slog.Handler + active chan<- struct{} +} + +func (h captureObserver) Handle(ctx context.Context, record slog.Record) error { + if record.Message == "CPU profiling capture started" { + h.active <- struct{}{} + } + return h.Handler.Handle(ctx, record) +} From 317a19f3222708afa8cba765aaab0ac7be47dcb3 Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Wed, 30 Sep 2026 07:40:43 -0400 Subject: [PATCH 2/2] fix(profiling): defer cooldown-skipped samples and configure triggers --- PROFILING.md | 3 +- cmd/beacon/profiling.go | 2 +- internal/profiling/profiling.go | 29 ++++++++---- internal/profiling/profiling_test.go | 68 ++++++++++++++++++++++++---- 4 files changed, 83 insertions(+), 19 deletions(-) diff --git a/PROFILING.md b/PROFILING.md index dc1a7a14..5e3667ee 100644 --- a/PROFILING.md +++ b/PROFILING.md @@ -34,7 +34,8 @@ can enable a longer session. Profiling adds overhead while a capture is active. - One 30-second sample immediately, then every 30 minutes. - Route reconfirmation requests an additional sample when the maintenance task starts. This includes the retention step before route validation. A five-minute cooldown - between capture starts prevents overlap and repeated triggers from increasing load. + between capture starts prevents overlap and repeated triggers from increasing load; + a periodic sample due during the cooldown runs when it ends. - Background task stacks carry a `task` label while profiling is enabled. - Shutdown or expiry stops the active sample and saves the shorter profile. - Each profile is limited to 8 MiB. The dedicated directory is limited to 256 MiB diff --git a/cmd/beacon/profiling.go b/cmd/beacon/profiling.go index d797fe21..df65d4cb 100644 --- a/cmd/beacon/profiling.go +++ b/cmd/beacon/profiling.go @@ -13,7 +13,7 @@ import ( ) func configureProfiling(ctx context.Context, pool *pgxpool.Pool) *profiling.Recorder { - r, err := profiling.Start(ctx, os.Getenv("BEACON_CPU_PROFILE_DIR"), os.Getenv("BEACON_CPU_PROFILE_UNTIL"), func() any { + r, err := profiling.Start(ctx, os.Getenv("BEACON_CPU_PROFILE_DIR"), os.Getenv("BEACON_CPU_PROFILE_UNTIL"), []string{"reconfirm"}, func() any { s := pool.Stat() return struct { Acquired, Idle, Total, Max int32 diff --git a/internal/profiling/profiling.go b/internal/profiling/profiling.go index 8d51cda9..ca836d9d 100644 --- a/internal/profiling/profiling.go +++ b/internal/profiling/profiling.go @@ -31,6 +31,7 @@ const ( type Recorder struct { root *os.Root snapshot func() any + triggers map[string]bool requests chan string cancel context.CancelFunc done chan struct{} @@ -52,7 +53,8 @@ type metadata struct { } // Start leaves expired settings inactive so restarts cannot extend a capture window. -func Start(parent context.Context, dir, until string, poolSnapshot func() any) (*Recorder, error) { +// Tasks named in triggers request an extra capture when they start. +func Start(parent context.Context, dir, until string, triggers []string, poolSnapshot func() any) (*Recorder, error) { if dir == "" && until == "" { return nil, nil } @@ -64,6 +66,7 @@ func Start(parent context.Context, dir, until string, poolSnapshot func() any) ( return nil, fmt.Errorf("invalid profiling deadline: %w", err) } if !deadline.After(time.Now()) { + slog.Info("CPU profiling deadline has passed; profiling stays off", "component", "profiling", "until", deadline.UTC()) return nil, nil } if time.Until(deadline) > 72*time.Hour { @@ -90,7 +93,10 @@ func Start(parent context.Context, dir, until string, poolSnapshot func() any) ( return nil, err } ctx, cancel := context.WithDeadline(parent, deadline) - r := &Recorder{root: root, snapshot: poolSnapshot, requests: make(chan string, 1), cancel: cancel, done: make(chan struct{})} + r := &Recorder{root: root, snapshot: poolSnapshot, triggers: make(map[string]bool, len(triggers)), requests: make(chan string, 1), cancel: cancel, done: make(chan struct{})} + for _, name := range triggers { + r.triggers[name] = true + } slog.Info("CPU profiling enabled", "component", "profiling", "until", deadline.UTC()) go r.run(ctx) return r, nil @@ -104,9 +110,9 @@ func (r *Recorder) Stop() { <-r.done } -// Request never holds up maintenance or queues a backlog of captures. -func (r *Recorder) Request(reason string) { - if r == nil || reason != "reconfirm" { +// request never holds up maintenance or queues a backlog of captures. +func (r *Recorder) request(reason string) { + if r == nil { return } select { @@ -130,8 +136,8 @@ func (r *Recorder) WrapTask(name string, run func(context.Context) error) func(c return run(ctx) default: } - if name == "reconfirm" { - r.Request(name) + if r.triggers[name] { + r.request(name) } pprof.Do(ctx, pprof.Labels("task", name), func(ctx context.Context) { err = run(ctx) }) return err @@ -153,21 +159,28 @@ func schedule(ctx context.Context, requests <-chan string, capture func(context. defer ticker.Stop() reason := "periodic" var last time.Time + var deferred <-chan time.Time for { if ctx.Err() != nil { return nil } - if last.IsZero() || time.Since(last) >= captureCooldown { + if wait := captureCooldown - time.Since(last); last.IsZero() || wait <= 0 { last = time.Now() if err := capture(ctx, reason); err != nil { return err } + } else if reason == "periodic" && deferred == nil { + // Run a periodic sample after the cooldown rather than skipping a whole interval. + deferred = time.After(wait) } select { case <-ctx.Done(): return nil case <-ticker.C: reason = "periodic" + case <-deferred: + deferred = nil + reason = "periodic" case reason = <-requests: } } diff --git a/internal/profiling/profiling_test.go b/internal/profiling/profiling_test.go index 88afb3f1..4a307db6 100644 --- a/internal/profiling/profiling_test.go +++ b/internal/profiling/profiling_test.go @@ -25,7 +25,7 @@ func TestDisabledAndExpired(t *testing.T) { for _, tc := range []struct{ dir, until string }{ {}, {filepath.Join(t.TempDir(), "unused"), time.Now().Add(-time.Hour).Format(time.RFC3339)}, } { - r, err := Start(context.Background(), tc.dir, tc.until, nil) + r, err := Start(context.Background(), tc.dir, tc.until, nil, nil) if err != nil || r != nil { t.Fatalf("got %v, %v", r, err) } @@ -46,7 +46,7 @@ func TestRejectUnsafeConfiguration(t *testing.T) { {"relative directory", "profiles", time.Now().Add(time.Hour).Format(time.RFC3339)}, } { t.Run(tc.name, func(t *testing.T) { - r, err := Start(context.Background(), tc.dir, tc.until, nil) + r, err := Start(context.Background(), tc.dir, tc.until, nil, nil) if r != nil { r.Stop() } @@ -59,7 +59,7 @@ func TestRejectUnsafeConfiguration(t *testing.T) { if err := os.Chmod(dir, 0755); err != nil { t.Fatal(err) } - r, err := Start(context.Background(), dir, time.Now().Add(time.Hour).Format(time.RFC3339), nil) + r, err := Start(context.Background(), dir, time.Now().Add(time.Hour).Format(time.RFC3339), nil, nil) if r != nil { r.Stop() } @@ -72,7 +72,7 @@ func TestDeadlineProducesPrivateProfileAndMetadata(t *testing.T) { requireSupportedPlatform(t) dir := filepath.Join(t.TempDir(), "profiles") until := time.Now().Add(500 * time.Millisecond) - r, err := Start(context.Background(), dir, until.Format(time.RFC3339Nano), func() any { return map[string]int{"acquired": 2} }) + r, err := Start(context.Background(), dir, until.Format(time.RFC3339Nano), nil, func() any { return map[string]int{"acquired": 2} }) if err != nil { t.Fatal(err) } @@ -123,7 +123,7 @@ func TestDeadlineProducesPrivateProfileAndMetadata(t *testing.T) { if m.Reason != "periodic" || !m.FinishedAt.After(m.StartedAt) || m.Before.Pool["acquired"] != 2 || m.After.Pool["acquired"] != 2 { t.Fatalf("bad metadata %+v", m) } - r.Request("reconfirm") + r.request("reconfirm") r.Stop() files, _ = filepath.Glob(filepath.Join(dir, "*.pprof")) if len(files) != 1 { @@ -138,7 +138,7 @@ func TestStopReleasesProfiler(t *testing.T) { slog.SetDefault(slog.New(captureObserver{Handler: slog.NewTextHandler(io.Discard, nil), active: active})) defer slog.SetDefault(old) dir := filepath.Join(t.TempDir(), "profiles") - r, err := Start(context.Background(), dir, time.Now().Add(time.Hour).Format(time.RFC3339), nil) + r, err := Start(context.Background(), dir, time.Now().Add(time.Hour).Format(time.RFC3339), nil, nil) if err != nil { t.Fatal(err) } @@ -167,7 +167,7 @@ func TestStopReleasesProfiler(t *testing.T) { } z.Close() dir2 := filepath.Join(t.TempDir(), "profiles") - r2, err := Start(context.Background(), dir2, time.Now().Add(400*time.Millisecond).Format(time.RFC3339Nano), nil) + r2, err := Start(context.Background(), dir2, time.Now().Add(400*time.Millisecond).Format(time.RFC3339Nano), nil, nil) if err != nil { t.Fatal(err) } @@ -206,7 +206,7 @@ func TestDirectoryBudgetSurvivesRestart(t *testing.T) { t.Fatal(err) } f.Close() - r, err := Start(context.Background(), dir, time.Now().Add(time.Hour).Format(time.RFC3339), nil) + r, err := Start(context.Background(), dir, time.Now().Add(time.Hour).Format(time.RFC3339), nil, nil) if err != nil { t.Fatal(err) } @@ -264,11 +264,61 @@ func TestScheduleBoundsExtraCaptures(t *testing.T) { }) } +func TestScheduleDefersPeriodicCaptureInCooldown(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + requests := make(chan string, 1) + var at []time.Duration + var reasons []string + start := time.Now() + done := make(chan error, 1) + go func() { + done <- schedule(ctx, requests, func(ctx context.Context, reason string) error { + at = append(at, time.Since(start)) + reasons = append(reasons, reason) + time.Sleep(30 * time.Second) + return nil + }) + }() + synctest.Wait() + time.Sleep(28 * time.Minute) + requests <- "reconfirm" + synctest.Wait() + time.Sleep(33 * time.Minute) + cancel() + synctest.Wait() + if err := <-done; err != nil { + t.Fatal(err) + } + wantAt := []time.Duration{0, 28 * time.Minute, 33 * time.Minute, 60 * time.Minute} + wantReasons := []string{"periodic", "reconfirm", "periodic", "periodic"} + if !reflect.DeepEqual(at, wantAt) || !reflect.DeepEqual(reasons, wantReasons) { + t.Fatalf("captures %v %v, want %v %v", at, reasons, wantAt, wantReasons) + } + }) +} + +func TestTriggersAreConfigurable(t *testing.T) { + r := &Recorder{triggers: map[string]bool{"custom": true}, requests: make(chan string, 1), done: make(chan struct{})} + noop := func(context.Context) error { return nil } + _ = r.WrapTask("reconfirm", noop)(context.Background()) + select { + case got := <-r.requests: + t.Fatalf("untriggered task requested %q", got) + default: + } + _ = r.WrapTask("custom", noop)(context.Background()) + if got := <-r.requests; got != "custom" { + t.Fatalf("got %q", got) + } +} + func TestWrappedTaskPreservesContextAndError(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) cancel() want := errors.New("task failed") - r := &Recorder{requests: make(chan string, 1), done: make(chan struct{})} + r := &Recorder{triggers: map[string]bool{"reconfirm": true}, requests: make(chan string, 1), done: make(chan struct{})} run := r.WrapTask("reconfirm", func(got context.Context) error { if got.Err() != context.Canceled { t.Fatal("lost cancellation")