From 760a9227a0f5ccb7bb906ceadb67e7c02ff0d374 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Sat, 3 Oct 2026 23:49:45 +0530 Subject: [PATCH 1/4] fix(profiler): keep the profiles of PIDs that succeed when another fails Every profiler fanned out over the container's leaf PIDs and returned the first error, which cancelled the PIDs not yet started and failed the whole run. A container whose process tree has a shell, a helper process or anything else the tool cannot attach to next to the real workload could therefore never be profiled without --pid, even though the workload's own PID would have produced a result. Run every PID through one shared helper instead. Each successful PID still publishes its own result event; the run succeeds when at least one PID did and emits a notice naming the skipped PIDs and why. It fails only when every PID failed, with each reason in the error: "PID : ; PID : ". An empty PID list is now an error rather than a success that publishes nothing. The Go pprof profiler also set the PID on the shared job from concurrent tasks, so a later PID could leak into an earlier PID's scrape; each PID now gets its own copy of the job. The helper uses a WaitGroup, so the worker-pool dependency is dropped. The test fakes for the per-PID managers now lock in invoke, which the helper calls from one goroutine per PID. --- go.mod | 1 - go.sum | 2 - internal/agent/profiler/bpf.go | 27 +-- internal/agent/profiler/bpf_fake.go | 5 + internal/agent/profiler/bpf_test.go | 10 +- internal/agent/profiler/common/pids.go | 87 +++++++++ internal/agent/profiler/common/pids_test.go | 167 ++++++++++++++++++ internal/agent/profiler/go_pprof.go | 37 ++-- internal/agent/profiler/jvm/async_profiler.go | 27 +-- .../agent/profiler/jvm/async_profiler_fake.go | 5 + .../agent/profiler/jvm/async_profiler_test.go | 10 +- internal/agent/profiler/jvm/jcmd.go | 27 +-- internal/agent/profiler/jvm/jcmd_fake.go | 5 + internal/agent/profiler/jvm/jcmd_test.go | 10 +- internal/agent/profiler/perf.go | 27 +-- internal/agent/profiler/perf_fake.go | 5 + internal/agent/profiler/perf_test.go | 10 +- internal/agent/profiler/python.go | 27 +-- internal/agent/profiler/python_austin.go | 27 +-- internal/agent/profiler/python_fake.go | 5 + internal/agent/profiler/python_test.go | 10 +- internal/agent/profiler/ruby.go | 28 +-- internal/agent/profiler/ruby_fake.go | 5 + internal/agent/profiler/ruby_test.go | 10 +- 24 files changed, 364 insertions(+), 210 deletions(-) create mode 100644 internal/agent/profiler/common/pids.go create mode 100644 internal/agent/profiler/common/pids_test.go diff --git a/go.mod b/go.mod index 2448690..751193e 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,6 @@ go 1.26.0 require ( github.com/agrison/go-commons-lang v0.0.0-20240106075236-2e001e6401ef - github.com/alitto/pond v1.9.2 github.com/golang/snappy v1.0.0 github.com/json-iterator/go v1.1.12 github.com/klauspost/compress v1.18.6 diff --git a/go.sum b/go.sum index 39f69d2..1194214 100644 --- a/go.sum +++ b/go.sum @@ -6,8 +6,6 @@ github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1 github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= github.com/agrison/go-commons-lang v0.0.0-20240106075236-2e001e6401ef h1:KkznClyESbRaLmRo7Oam4vv5L4oknDK+mixJ9mypl6E= github.com/agrison/go-commons-lang v0.0.0-20240106075236-2e001e6401ef/go.mod h1:u+Zwm0OKtJAGx+DXcmp2NNwZ0GKtV80ipbF/uhKhQdw= -github.com/alitto/pond v1.9.2 h1:9Qb75z/scEZVCoSU+osVmQ0I0JOeLfdTDafrbcJ8CLs= -github.com/alitto/pond v1.9.2/go.mod h1:xQn3P/sHTYcU/1BR3i86IGIrilcrGC2LiS+E2+CJWsI= github.com/armon/go-socks5 v0.0.0-20160902184237-e75332964ef5 h1:0CwZNZbxp69SHPdPJAN/hZIm0C4OItdklCFmMRWYpio= github.com/armon/go-socks5 v0.0.0-20160902184237-e75332964ef5/go.mod h1:wHh0iHkYZB8zMSxRWpUBQtwG5a7fFgvEO+odwuTv2gs= github.com/blang/semver/v4 v4.0.0 h1:1PFHFE6yCCTv8C1TeyNNarDzntLi7wMI5i/pzqYIsAM= diff --git a/internal/agent/profiler/bpf.go b/internal/agent/profiler/bpf.go index 065ac82..d041277 100644 --- a/internal/agent/profiler/bpf.go +++ b/internal/agent/profiler/bpf.go @@ -2,14 +2,12 @@ package profiler import ( "bytes" - "context" "fmt" "os/exec" "strconv" "time" "github.com/agrison/go-commons-lang/stringUtils" - "github.com/alitto/pond" "github.com/nudgebee/application-profiler/api" "github.com/nudgebee/application-profiler/internal/agent/config" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -77,27 +75,10 @@ func (b *BpfProfiler) SetUp(job *job.ProfilingJob) error { func (b *BpfProfiler) Invoke(job *job.ProfilingJob) (error, time.Duration) { start := time.Now() - - pool := pond.New(len(b.targetPIDs), 0, pond.MinWorkers(len(b.targetPIDs))) - defer pool.StopAndWait() - - // create a task group associated to a context - group, _ := pool.GroupContext(context.Background()) - - // submit tasks to profile - for _, pid := range b.targetPIDs { - pid := pid - group.Submit(func() error { - err, _ := b.invoke(job, pid) - return err - }) - // wait a bit between jobs for not overloading the system - time.Sleep(b.delay) - } - - // wait for all tasks to finish - err := group.Wait() - + err := common.ProfilePIDs(b.targetPIDs, b.delay, func(pid string) error { + err, _ := b.invoke(job, pid) + return err + }) return err, time.Since(start) } diff --git a/internal/agent/profiler/bpf_fake.go b/internal/agent/profiler/bpf_fake.go index 0a27996..77fd50b 100644 --- a/internal/agent/profiler/bpf_fake.go +++ b/internal/agent/profiler/bpf_fake.go @@ -1,6 +1,7 @@ package profiler import ( + "sync" "time" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -17,6 +18,8 @@ type FakeBpfManager interface { // fakeBpfManager is an implementation of the FakeBpfManager interface type fakeBpfManager struct { + // mu serialises invoke, which Invoke calls concurrently, once per PID. + mu sync.Mutex fakeMethods map[string]*fakeBpfManagerMethod } @@ -60,6 +63,8 @@ func (f *fakeBpfManagerMethod) InvokedTimes() int { } func (f *fakeBpfManager) invoke(*job.ProfilingJob, string) (error, time.Duration) { + f.mu.Lock() + defer f.mu.Unlock() var err error var duration time.Duration f.fakeMethods["invoke"].invokes++ diff --git a/internal/agent/profiler/bpf_test.go b/internal/agent/profiler/bpf_test.go index e68cfc6..266593b 100644 --- a/internal/agent/profiler/bpf_test.go +++ b/internal/agent/profiler/bpf_test.go @@ -167,10 +167,12 @@ func TestBpfProfiler_Invoke(t *testing.T) { }, }, { - name: "should invoke fail when invoke fail", + name: "should invoke fail when invoke fails for every PID", given: func() (fields, args) { bpfManager := newFakeBpfManager() - bpfManager.On("invoke").Return(errors.New("fake invoke error"), time.Duration(0)) + bpfManager.On("invoke"). + Return(errors.New("fake invoke error"), time.Duration(0)). + Return(errors.New("fake invoke error"), time.Duration(0)) return fields{ BpfProfiler: &BpfProfiler{ @@ -192,8 +194,8 @@ func TestBpfProfiler_Invoke(t *testing.T) { }, then: func(t *testing.T, err error, fields fields) { require.Error(t, err) - assert.EqualError(t, err, "fake invoke error") - assert.Equal(t, 1, fields.BpfProfiler.BpfManager.(FakeBpfManager).On("invoke").InvokedTimes()) + assert.EqualError(t, err, "PID 1000: fake invoke error; PID 2000: fake invoke error") + assert.Equal(t, 2, fields.BpfProfiler.BpfManager.(FakeBpfManager).On("invoke").InvokedTimes()) }, }, } diff --git a/internal/agent/profiler/common/pids.go b/internal/agent/profiler/common/pids.go new file mode 100644 index 0000000..7f7ecc5 --- /dev/null +++ b/internal/agent/profiler/common/pids.go @@ -0,0 +1,87 @@ +package common + +import ( + "fmt" + "strings" + "sync" + "time" + + "github.com/nudgebee/application-profiler/api" + "github.com/nudgebee/application-profiler/pkg/util/log" + "github.com/pkg/errors" +) + +// ProfilePIDs runs profile once for every PID, concurrently, starting each one +// delay after the previous so a container with many processes is not hit all +// at once. Each PID publishes its own result. +// +// The PIDs are the leaf processes of the container's process tree, which +// routinely include ones the tool cannot profile: a shell, a sidecar-style +// helper, a process that is not the target's runtime, one that exits +// mid-run. One of those failing must not throw away the profiles the others +// produced, so the run succeeds as long as at least one PID did, and a notice +// names the PIDs that were skipped and why. It fails only when every PID +// failed, with each one's reason. +func ProfilePIDs(pids []string, delay time.Duration, profile func(pid string) error) error { + if len(pids) == 0 { + // Nothing would be published, so succeeding here would end the run + // without a result. + return errors.New("no PIDs to profile") + } + + errs := make([]error, len(pids)) + var wg sync.WaitGroup + for i, pid := range pids { + if i > 0 { + // wait a bit between jobs for not overloading the system + time.Sleep(delay) + } + wg.Go(func() { errs[i] = profile(pid) }) + } + wg.Wait() + + var failed pidErrors + for i, err := range errs { + if err != nil { + failed = append(failed, pidError{pid: pids[i], err: err}) + } + } + switch len(failed) { + case 0: + return nil + case len(pids): + return failed + } + + _ = log.EventLn(api.Notice, &api.NoticeData{ + Time: time.Now(), + Msg: fmt.Sprintf("Profiled %d of %d PIDs; skipped %s", + len(pids)-len(failed), len(pids), failed.Error()), + }) + return nil +} + +// pidError is the reason one PID could not be profiled. +type pidError struct { + pid string + err error +} + +// pidErrors reads "PID : ; PID : ". +type pidErrors []pidError + +func (e pidErrors) Error() string { + reasons := make([]string, len(e)) + for i, f := range e { + reasons[i] = fmt.Sprintf("PID %s: %s", f.pid, f.err) + } + return strings.Join(reasons, "; ") +} + +func (e pidErrors) Unwrap() []error { + errs := make([]error, len(e)) + for i, f := range e { + errs[i] = f.err + } + return errs +} diff --git a/internal/agent/profiler/common/pids_test.go b/internal/agent/profiler/common/pids_test.go new file mode 100644 index 0000000..60b6696 --- /dev/null +++ b/internal/agent/profiler/common/pids_test.go @@ -0,0 +1,167 @@ +package common + +import ( + "bufio" + "errors" + "io" + "os" + "sync" + "testing" + "time" + + "github.com/nudgebee/application-profiler/api" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// captureEvents returns the events printed on stdout while fn runs. +func captureEvents(t *testing.T, fn func()) []string { + t.Helper() + r, w, err := os.Pipe() + require.NoError(t, err) + old := os.Stdout + os.Stdout = w + defer func() { os.Stdout = old }() + + var out []string + done := make(chan struct{}) + go func() { + defer close(done) + scanner := bufio.NewScanner(r) + for scanner.Scan() { + out = append(out, scanner.Text()) + } + _, _ = io.Copy(io.Discard, r) + }() + fn() + require.NoError(t, w.Close()) + <-done + return out +} + +func notices(t *testing.T, events []string) []string { + t.Helper() + var msgs []string + for _, e := range events { + data, err := api.ParseEvent(e) + require.NoError(t, err) + if n, ok := data.(*api.NoticeData); ok { + msgs = append(msgs, n.Msg) + } + } + return msgs +} + +// profileWith fakes a per-PID profile: a PID listed in failures fails with +// that reason, any other succeeds. It records which PIDs ran. +func profileWith(failures map[string]string) (func(string) error, func() []string) { + var mu sync.Mutex + var ran []string + return func(pid string) error { + mu.Lock() + ran = append(ran, pid) + mu.Unlock() + if reason, ok := failures[pid]; ok { + return errors.New(reason) + } + return nil + }, func() []string { + mu.Lock() + defer mu.Unlock() + return append([]string(nil), ran...) + } +} + +func TestProfilePIDs(t *testing.T) { + t.Run("one PID failing does not fail the run, and is reported in a notice", func(t *testing.T) { + profile, ran := profileWith(map[string]string{"8": "not a Python process"}) + var err error + + events := captureEvents(t, func() { + err = ProfilePIDs([]string{"7", "8", "9"}, 0, profile) + }) + + require.NoError(t, err) + assert.ElementsMatch(t, []string{"7", "8", "9"}, ran(), "every PID should be tried") + assert.Equal(t, []string{"Profiled 2 of 3 PIDs; skipped PID 8: not a Python process"}, notices(t, events)) + }) + + t.Run("every PID failing fails the run with each reason", func(t *testing.T) { + profile, ran := profileWith(map[string]string{"7": "permission denied", "8": "no such process"}) + var err error + + events := captureEvents(t, func() { + err = ProfilePIDs([]string{"7", "8"}, 0, profile) + }) + + require.Error(t, err) + assert.EqualError(t, err, "PID 7: permission denied; PID 8: no such process") + assert.ElementsMatch(t, []string{"7", "8"}, ran()) + assert.Empty(t, notices(t, events), "a failed run reports through its error, not a notice") + }) + + t.Run("the per-PID errors stay reachable", func(t *testing.T) { + cause := errors.New("cause") + err := ProfilePIDs([]string{"7"}, 0, func(string) error { return cause }) + + assert.EqualError(t, err, "PID 7: cause") + assert.ErrorIs(t, err, cause) + }) + + t.Run("every PID succeeding needs no notice", func(t *testing.T) { + profile, ran := profileWith(nil) + var err error + + events := captureEvents(t, func() { + err = ProfilePIDs([]string{"7", "8"}, 0, profile) + }) + + require.NoError(t, err) + assert.ElementsMatch(t, []string{"7", "8"}, ran()) + assert.Empty(t, notices(t, events)) + }) + + t.Run("no PIDs is a failure, not a run without a result", func(t *testing.T) { + profile, ran := profileWith(nil) + + err := ProfilePIDs(nil, 0, profile) + + assert.EqualError(t, err, "no PIDs to profile") + assert.Empty(t, ran()) + }) + + t.Run("PIDs are started delay apart but run concurrently", func(t *testing.T) { + const delay = 50 * time.Millisecond + var mu sync.Mutex + started := map[string]time.Time{} + startedLate := false + release := make(chan struct{}) + go func() { + // Hold every PID until well after the last should have started: a + // sequential run would not start the second until this fires. + time.Sleep(5 * delay) + close(release) + }() + + err := ProfilePIDs([]string{"7", "8", "9"}, delay, func(pid string) error { + mu.Lock() + started[pid] = time.Now() + select { + case <-release: + startedLate = true + default: + } + mu.Unlock() + <-release + return nil + }) + + require.NoError(t, err) + require.Len(t, started, 3) + assert.False(t, startedLate, "every PID should start while the others are still running") + // Scheduling jitter can shave a little off the gap between two + // goroutines starting; the stagger itself is a sleep of delay. + assert.Greater(t, started["8"].Sub(started["7"]), delay/2) + assert.Greater(t, started["9"].Sub(started["8"]), delay/2) + }) +} diff --git a/internal/agent/profiler/go_pprof.go b/internal/agent/profiler/go_pprof.go index ff5d9b2..71f3d6a 100644 --- a/internal/agent/profiler/go_pprof.go +++ b/internal/agent/profiler/go_pprof.go @@ -3,7 +3,6 @@ package profiler import ( "bufio" "bytes" - "context" "fmt" "os/exec" "strconv" @@ -11,7 +10,6 @@ import ( "time" "github.com/agrison/go-commons-lang/stringUtils" - "github.com/alitto/pond" "github.com/nudgebee/application-profiler/api" "github.com/nudgebee/application-profiler/internal/agent/config" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -24,10 +22,13 @@ import ( "github.com/pkg/errors" ) +const goPprofDelayBetweenJobs = 2 * time.Second + // GoPprofProfiler uses Go's pprof HTTP server to collect profiles. type GoPprofProfiler struct { manager *goPprofManager targetPIDs []string + delay time.Duration } type goPprofManager struct { @@ -37,7 +38,10 @@ type goPprofManager struct { // NewGoPprofProfiler creates a new profiler with the given publisher. func NewGoPprofProfiler(commander executil.Commander, publisher publish.Publisher) *GoPprofProfiler { - return &GoPprofProfiler{manager: &goPprofManager{commander: commander, publisher: publisher}} + return &GoPprofProfiler{ + manager: &goPprofManager{commander: commander, publisher: publisher}, + delay: goPprofDelayBetweenJobs, + } } // SetUp ensures the job has a valid PID. @@ -59,26 +63,17 @@ func (p *GoPprofProfiler) SetUp(job *job.ProfilingJob) error { // Invoke runs the profiling job and returns execution time. func (p *GoPprofProfiler) Invoke(job *job.ProfilingJob) (error, time.Duration) { start := time.Now() - pool := pond.New(len(p.targetPIDs), 0, pond.MinWorkers(len(p.targetPIDs))) - defer pool.StopAndWait() - // create a task group associated to a context - group, _ := pool.GroupContext(context.Background()) - // submit tasks to profile - for _, pid := range p.targetPIDs { - pid := pid - group.Submit(func() error { - job.PID = pid - err := p.manager.fetchProfileFromPID(job) - return err - }) - // wait a bit between jobs for not overloading the system - time.Sleep(2 * time.Second) - } - // wait for all tasks to finish - err := group.Wait() - + err := common.ProfilePIDs(p.targetPIDs, p.delay, func(pid string) error { + // The PIDs run concurrently, so each needs its own copy of the job: + // setting PID on the shared one let a later PID's value leak into an + // earlier PID's scrape. + pidJob := *job + pidJob.PID = pid + return p.manager.fetchProfileFromPID(&pidJob) + }) return err, time.Since(start) } + func (p *goPprofManager) heapProfile(job *job.ProfilingJob, port string, fileName string) error { var out bytes.Buffer var stderr bytes.Buffer diff --git a/internal/agent/profiler/jvm/async_profiler.go b/internal/agent/profiler/jvm/async_profiler.go index afd4daa..05b49ef 100644 --- a/internal/agent/profiler/jvm/async_profiler.go +++ b/internal/agent/profiler/jvm/async_profiler.go @@ -2,7 +2,6 @@ package jvm import ( "bytes" - "context" "fmt" "io" "os" @@ -12,7 +11,6 @@ import ( "time" "github.com/agrison/go-commons-lang/stringUtils" - "github.com/alitto/pond" "github.com/nudgebee/application-profiler/api" "github.com/nudgebee/application-profiler/internal/agent/config" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -168,27 +166,10 @@ func (j *asyncProfilerManager) selectProfilerLibrary(targetFs string) error { func (j *AsyncProfiler) Invoke(job *job.ProfilingJob) (error, time.Duration) { start := time.Now() - - pool := pond.New(len(j.targetPIDs), 0, pond.MinWorkers(len(j.targetPIDs))) - defer pool.StopAndWait() - - // create a task group associated to a context - group, _ := pool.GroupContext(context.Background()) - - // submit tasks to profile - for _, pid := range j.targetPIDs { - pid := pid - group.Submit(func() error { - err, _ := j.invoke(job, pid) - return err - }) - // wait a bit between jobs for not overloading the system - time.Sleep(j.delay) - } - - // wait for all tasks to finish - err := group.Wait() - + err := common.ProfilePIDs(j.targetPIDs, j.delay, func(pid string) error { + err, _ := j.invoke(job, pid) + return err + }) return err, time.Since(start) } diff --git a/internal/agent/profiler/jvm/async_profiler_fake.go b/internal/agent/profiler/jvm/async_profiler_fake.go index 73ef552..c62ab32 100644 --- a/internal/agent/profiler/jvm/async_profiler_fake.go +++ b/internal/agent/profiler/jvm/async_profiler_fake.go @@ -1,6 +1,7 @@ package jvm import ( + "sync" "time" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -16,6 +17,8 @@ type FakeAsyncProfilerManager interface { // fakeAsyncProfilerManager is an implementation of the FakeAsyncProfilerManager interface type fakeAsyncProfilerManager struct { + // mu serialises invoke, which Invoke calls concurrently, once per PID. + mu sync.Mutex fakeMethods map[string]*fakeAsyncProfilerManagerMethod } @@ -59,6 +62,8 @@ func (f *fakeAsyncProfilerManagerMethod) InvokedTimes() int { } func (f *fakeAsyncProfilerManager) invoke(*job.ProfilingJob, string) (error, time.Duration) { + f.mu.Lock() + defer f.mu.Unlock() var err error var duration time.Duration f.fakeMethods["invoke"].invokes++ diff --git a/internal/agent/profiler/jvm/async_profiler_test.go b/internal/agent/profiler/jvm/async_profiler_test.go index 33cc3cc..e1c1667 100644 --- a/internal/agent/profiler/jvm/async_profiler_test.go +++ b/internal/agent/profiler/jvm/async_profiler_test.go @@ -310,10 +310,12 @@ func TestAsyncProfiler_Invoke(t *testing.T) { }, }, { - name: "should invoke fail when invoke fail", + name: "should invoke fail when invoke fails for every PID", given: func() (fields, args) { asyncProfilerManager := newFakeAsyncProfilerManager() - asyncProfilerManager.On("invoke").Return(errors.New("fake invoke error"), time.Duration(0)) + asyncProfilerManager.On("invoke"). + Return(errors.New("fake invoke error"), time.Duration(0)). + Return(errors.New("fake invoke error"), time.Duration(0)) return fields{ AsyncProfiler: AsyncProfiler{ @@ -335,8 +337,8 @@ func TestAsyncProfiler_Invoke(t *testing.T) { }, then: func(t *testing.T, err error, fields fields) { require.Error(t, err) - assert.EqualError(t, err, "fake invoke error") - assert.Equal(t, 1, fields.AsyncProfiler.AsyncProfilerManager.(FakeAsyncProfilerManager).On("invoke").InvokedTimes()) + assert.EqualError(t, err, "PID 1000: fake invoke error; PID 2000: fake invoke error") + assert.Equal(t, 2, fields.AsyncProfiler.AsyncProfilerManager.(FakeAsyncProfilerManager).On("invoke").InvokedTimes()) }, }, } diff --git a/internal/agent/profiler/jvm/jcmd.go b/internal/agent/profiler/jvm/jcmd.go index c6f354a..87f125a 100644 --- a/internal/agent/profiler/jvm/jcmd.go +++ b/internal/agent/profiler/jvm/jcmd.go @@ -2,7 +2,6 @@ package jvm import ( "bytes" - "context" "fmt" "io" "os" @@ -12,7 +11,6 @@ import ( "time" "github.com/agrison/go-commons-lang/stringUtils" - "github.com/alitto/pond" "github.com/nudgebee/application-profiler/api" "github.com/nudgebee/application-profiler/internal/agent/config" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -147,27 +145,10 @@ func (j *jcmdManager) copyJfrSettingsToTmpDir() error { func (j *JcmdProfiler) Invoke(job *job.ProfilingJob) (error, time.Duration) { start := time.Now() - - pool := pond.New(len(j.targetPIDs), 0, pond.MinWorkers(len(j.targetPIDs))) - defer pool.StopAndWait() - - // create a task group associated to a context - group, _ := pool.GroupContext(context.Background()) - - // submit tasks to profile - for _, pid := range j.targetPIDs { - pid := pid - group.Submit(func() error { - err, _ := j.invoke(job, pid) - return err - }) - // wait a bit between jobs for not overloading the system - time.Sleep(j.delay) - } - - // wait for all tasks to finish - err := group.Wait() - + err := common.ProfilePIDs(j.targetPIDs, j.delay, func(pid string) error { + err, _ := j.invoke(job, pid) + return err + }) return err, time.Since(start) } diff --git a/internal/agent/profiler/jvm/jcmd_fake.go b/internal/agent/profiler/jvm/jcmd_fake.go index 0f64b05..3b3e8a7 100644 --- a/internal/agent/profiler/jvm/jcmd_fake.go +++ b/internal/agent/profiler/jvm/jcmd_fake.go @@ -2,6 +2,7 @@ package jvm import ( "bytes" + "sync" "time" "github.com/nudgebee/application-profiler/api" @@ -19,6 +20,8 @@ type FakeJcmdManager interface { // fakeJcmdManager is an implementation of the FakeJcmdManager interface type fakeJcmdManager struct { + // mu serialises invoke, which Invoke calls concurrently, once per PID. + mu sync.Mutex fakeMethods map[string]*fakeJcmdManagerMethod } @@ -62,6 +65,8 @@ func (f *fakeJcmdManagerMethod) InvokedTimes() int { } func (f *fakeJcmdManager) invoke(*job.ProfilingJob, string) (error, time.Duration) { + f.mu.Lock() + defer f.mu.Unlock() var err error var duration time.Duration f.fakeMethods["invoke"].invokes++ diff --git a/internal/agent/profiler/jvm/jcmd_test.go b/internal/agent/profiler/jvm/jcmd_test.go index a32b2fb..4649dd3 100644 --- a/internal/agent/profiler/jvm/jcmd_test.go +++ b/internal/agent/profiler/jvm/jcmd_test.go @@ -306,10 +306,12 @@ func TestJcmdProfiler_Invoke(t *testing.T) { }, }, { - name: "should invoke fail when invoke fail", + name: "should invoke fail when invoke fails for every PID", given: func() (fields, args) { jcmdManager := newFakeJcmdManager() - jcmdManager.On("invoke").Return(errors.New("fake invoke error"), time.Duration(0)) + jcmdManager.On("invoke"). + Return(errors.New("fake invoke error"), time.Duration(0)). + Return(errors.New("fake invoke error"), time.Duration(0)) return fields{ JcmdProfiler: JcmdProfiler{ @@ -331,8 +333,8 @@ func TestJcmdProfiler_Invoke(t *testing.T) { }, then: func(t *testing.T, err error, fields fields) { require.Error(t, err) - assert.EqualError(t, err, "fake invoke error") - assert.Equal(t, 1, fields.JcmdProfiler.JcmdManager.(FakeJcmdManager).On("invoke").InvokedTimes()) + assert.EqualError(t, err, "PID 1000: fake invoke error; PID 2000: fake invoke error") + assert.Equal(t, 2, fields.JcmdProfiler.JcmdManager.(FakeJcmdManager).On("invoke").InvokedTimes()) }, }, } diff --git a/internal/agent/profiler/perf.go b/internal/agent/profiler/perf.go index fdf5f8a..fd0c222 100644 --- a/internal/agent/profiler/perf.go +++ b/internal/agent/profiler/perf.go @@ -2,14 +2,12 @@ package profiler import ( "bytes" - "context" "fmt" "os" "strconv" "time" "github.com/agrison/go-commons-lang/stringUtils" - "github.com/alitto/pond" "github.com/nudgebee/application-profiler/api" "github.com/nudgebee/application-profiler/internal/agent/config" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -77,27 +75,10 @@ func (p *PerfProfiler) SetUp(job *job.ProfilingJob) error { func (p *PerfProfiler) Invoke(job *job.ProfilingJob) (error, time.Duration) { start := time.Now() - - pool := pond.New(len(p.targetPIDs), 0, pond.MinWorkers(len(p.targetPIDs))) - defer pool.StopAndWait() - - // create a task group associated to a context - group, _ := pool.GroupContext(context.Background()) - - // submit tasks to profile - for _, pid := range p.targetPIDs { - pid := pid - group.Submit(func() error { - err, _ := p.invoke(job, pid) - return err - }) - // wait a bit between jobs for not overloading the system - time.Sleep(p.delay) - } - - // wait for all tasks to finish - err := group.Wait() - + err := common.ProfilePIDs(p.targetPIDs, p.delay, func(pid string) error { + err, _ := p.invoke(job, pid) + return err + }) return err, time.Since(start) } diff --git a/internal/agent/profiler/perf_fake.go b/internal/agent/profiler/perf_fake.go index 3c58115..d413e4f 100644 --- a/internal/agent/profiler/perf_fake.go +++ b/internal/agent/profiler/perf_fake.go @@ -1,6 +1,7 @@ package profiler import ( + "sync" "time" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -17,6 +18,8 @@ type FakePerfManager interface { // fakePerfManager is an implementation of the FakePerfManager interface type fakePerfManager struct { + // mu serialises invoke, which Invoke calls concurrently, once per PID. + mu sync.Mutex fakeMethods map[string]*fakePerfManagerMethod } @@ -60,6 +63,8 @@ func (f *fakePerfManagerMethod) InvokedTimes() int { } func (f *fakePerfManager) invoke(*job.ProfilingJob, string) (error, time.Duration) { + f.mu.Lock() + defer f.mu.Unlock() var err error var duration time.Duration f.fakeMethods["invoke"].invokes++ diff --git a/internal/agent/profiler/perf_test.go b/internal/agent/profiler/perf_test.go index a25d8b1..1b895e7 100644 --- a/internal/agent/profiler/perf_test.go +++ b/internal/agent/profiler/perf_test.go @@ -167,10 +167,12 @@ func TestPerfProfiler_Invoke(t *testing.T) { }, }, { - name: "should invoke fail when invoke fail", + name: "should invoke fail when invoke fails for every PID", given: func() (fields, args) { fakePerfManager := newFakePerfManager() - fakePerfManager.On("invoke").Return(errors.New("fake invoke error"), time.Duration(0)) + fakePerfManager.On("invoke"). + Return(errors.New("fake invoke error"), time.Duration(0)). + Return(errors.New("fake invoke error"), time.Duration(0)) return fields{ PerfProfiler: &PerfProfiler{ @@ -192,8 +194,8 @@ func TestPerfProfiler_Invoke(t *testing.T) { }, then: func(t *testing.T, err error, fields fields) { require.Error(t, err) - assert.EqualError(t, err, "fake invoke error") - assert.Equal(t, 1, fields.PerfProfiler.PerfManager.(FakePerfManager).On("invoke").InvokedTimes()) + assert.EqualError(t, err, "PID 1000: fake invoke error; PID 2000: fake invoke error") + assert.Equal(t, 2, fields.PerfProfiler.PerfManager.(FakePerfManager).On("invoke").InvokedTimes()) }, }, } diff --git a/internal/agent/profiler/python.go b/internal/agent/profiler/python.go index 1af18c9..15ee557 100644 --- a/internal/agent/profiler/python.go +++ b/internal/agent/profiler/python.go @@ -2,14 +2,12 @@ package profiler import ( "bytes" - "context" "fmt" "os/exec" "strconv" "time" "github.com/agrison/go-commons-lang/stringUtils" - "github.com/alitto/pond" "github.com/nudgebee/application-profiler/api" "github.com/nudgebee/application-profiler/internal/agent/config" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -92,27 +90,10 @@ func (p *PythonProfiler) SetUp(job *job.ProfilingJob) error { func (p *PythonProfiler) Invoke(job *job.ProfilingJob) (error, time.Duration) { start := time.Now() - - pool := pond.New(len(p.targetPIDs), 0, pond.MinWorkers(len(p.targetPIDs))) - defer pool.StopAndWait() - - // create a task group associated to a context - group, _ := pool.GroupContext(context.Background()) - - // submit tasks to profile - for _, pid := range p.targetPIDs { - pid := pid - group.Submit(func() error { - err, _ := p.invoke(job, pid) - return err - }) - // wait a bit between jobs for not overloading the system - time.Sleep(p.delay) - } - - // wait for all tasks to finish - err := group.Wait() - + err := common.ProfilePIDs(p.targetPIDs, p.delay, func(pid string) error { + err, _ := p.invoke(job, pid) + return err + }) return err, time.Since(start) } diff --git a/internal/agent/profiler/python_austin.go b/internal/agent/profiler/python_austin.go index 557a6ac..1fb2728 100644 --- a/internal/agent/profiler/python_austin.go +++ b/internal/agent/profiler/python_austin.go @@ -3,7 +3,6 @@ package profiler import ( "bufio" "bytes" - "context" "fmt" "io" "os" @@ -16,7 +15,6 @@ import ( "unicode" "github.com/agrison/go-commons-lang/stringUtils" - "github.com/alitto/pond" "github.com/nudgebee/application-profiler/api" "github.com/nudgebee/application-profiler/internal/agent/config" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -215,27 +213,10 @@ func (p *AustinPythonProfiler) SetUp(job *job.ProfilingJob) error { func (p *AustinPythonProfiler) Invoke(job *job.ProfilingJob) (error, time.Duration) { start := time.Now() - - pool := pond.New(len(p.targetPIDs), 0, pond.MinWorkers(len(p.targetPIDs))) - defer pool.StopAndWait() - - // create a task group associated to a context - group, _ := pool.GroupContext(context.Background()) - - // submit tasks to profile - for _, pid := range p.targetPIDs { - pid := pid - group.Submit(func() error { - err, _ := p.invoke(job, pid) - return err - }) - // wait a bit between jobs for not overloading the system - time.Sleep(p.delay) - } - - // wait for all tasks to finish - err := group.Wait() - + err := common.ProfilePIDs(p.targetPIDs, p.delay, func(pid string) error { + err, _ := p.invoke(job, pid) + return err + }) return err, time.Since(start) } diff --git a/internal/agent/profiler/python_fake.go b/internal/agent/profiler/python_fake.go index 1e6c890..39d1eed 100644 --- a/internal/agent/profiler/python_fake.go +++ b/internal/agent/profiler/python_fake.go @@ -1,6 +1,7 @@ package profiler import ( + "sync" "time" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -17,6 +18,8 @@ type FakePythonManager interface { // fakePythonManager is an implementation of the FakePythonManager interface type fakePythonManager struct { + // mu serialises invoke, which Invoke calls concurrently, once per PID. + mu sync.Mutex fakeMethods map[string]*fakePythonManagerMethod } @@ -60,6 +63,8 @@ func (f *fakePythonManagerMethod) InvokedTimes() int { } func (f *fakePythonManager) invoke(*job.ProfilingJob, string) (error, time.Duration) { + f.mu.Lock() + defer f.mu.Unlock() var err error var duration time.Duration f.fakeMethods["invoke"].invokes++ diff --git a/internal/agent/profiler/python_test.go b/internal/agent/profiler/python_test.go index 10b408b..5542cf6 100644 --- a/internal/agent/profiler/python_test.go +++ b/internal/agent/profiler/python_test.go @@ -167,10 +167,12 @@ func TestPythonProfiler_Invoke(t *testing.T) { }, }, { - name: "should invoke fail when invoke fail", + name: "should invoke fail when invoke fails for every PID", given: func() (fields, args) { pythonManager := newFakePythonManager() - pythonManager.On("invoke").Return(errors.New("fake invoke error"), time.Duration(0)) + pythonManager.On("invoke"). + Return(errors.New("fake invoke error"), time.Duration(0)). + Return(errors.New("fake invoke error"), time.Duration(0)) return fields{ PythonProfiler: &PythonProfiler{ @@ -192,8 +194,8 @@ func TestPythonProfiler_Invoke(t *testing.T) { }, then: func(t *testing.T, err error, fields fields) { require.Error(t, err) - assert.EqualError(t, err, "fake invoke error") - assert.Equal(t, 1, fields.PythonProfiler.PythonManager.(FakePythonManager).On("invoke").InvokedTimes()) + assert.EqualError(t, err, "PID 1000: fake invoke error; PID 2000: fake invoke error") + assert.Equal(t, 2, fields.PythonProfiler.PythonManager.(FakePythonManager).On("invoke").InvokedTimes()) }, }, } diff --git a/internal/agent/profiler/ruby.go b/internal/agent/profiler/ruby.go index 122ea0a..c123b96 100644 --- a/internal/agent/profiler/ruby.go +++ b/internal/agent/profiler/ruby.go @@ -2,14 +2,12 @@ package profiler import ( "bytes" - "context" "fmt" "os/exec" "strconv" "time" "github.com/agrison/go-commons-lang/stringUtils" - "github.com/alitto/pond" "github.com/nudgebee/application-profiler/internal/agent/config" "github.com/nudgebee/application-profiler/internal/agent/job" "github.com/nudgebee/application-profiler/internal/agent/profiler/common" @@ -78,28 +76,10 @@ func (r *RubyProfiler) SetUp(job *job.ProfilingJob) error { func (r *RubyProfiler) Invoke(job *job.ProfilingJob) (error, time.Duration) { start := time.Now() - - // create a pool of workers - pool := pond.New(len(r.targetPIDs), 0, pond.MinWorkers(len(r.targetPIDs))) - defer pool.StopAndWait() - - // create a task group associated to a context - group, _ := pool.GroupContext(context.Background()) - - // submit tasks to profile - for _, pid := range r.targetPIDs { - pid := pid - group.Submit(func() error { - err, _ := r.invoke(job, pid) - return err - }) - // wait a bit between jobs for not overloading the system - time.Sleep(r.delay) - } - - // wait for all tasks to finish - err := group.Wait() - + err := common.ProfilePIDs(r.targetPIDs, r.delay, func(pid string) error { + err, _ := r.invoke(job, pid) + return err + }) return err, time.Since(start) } diff --git a/internal/agent/profiler/ruby_fake.go b/internal/agent/profiler/ruby_fake.go index c41753f..5d277de 100644 --- a/internal/agent/profiler/ruby_fake.go +++ b/internal/agent/profiler/ruby_fake.go @@ -1,6 +1,7 @@ package profiler import ( + "sync" "time" "github.com/nudgebee/application-profiler/internal/agent/job" @@ -17,6 +18,8 @@ type FakeRubyManager interface { // fakeRubyManager is an implementation of the FakeRubyManager interface type fakeRubyManager struct { + // mu serialises invoke, which Invoke calls concurrently, once per PID. + mu sync.Mutex fakeMethods map[string]*fakeRubyManagerMethod } @@ -60,6 +63,8 @@ func (f *fakeRubyManagerMethod) InvokedTimes() int { } func (f *fakeRubyManager) invoke(*job.ProfilingJob, string) (error, time.Duration) { + f.mu.Lock() + defer f.mu.Unlock() var err error var duration time.Duration f.fakeMethods["invoke"].invokes++ diff --git a/internal/agent/profiler/ruby_test.go b/internal/agent/profiler/ruby_test.go index 7f668a6..9e19aa0 100644 --- a/internal/agent/profiler/ruby_test.go +++ b/internal/agent/profiler/ruby_test.go @@ -166,10 +166,12 @@ func TestRubyProfiler_Invoke(t *testing.T) { }, }, { - name: "should invoke fail when invoke fail", + name: "should invoke fail when invoke fails for every PID", given: func() (fields, args) { rubyManager := newFakeRubyManager() - rubyManager.On("invoke").Return(errors.New("fake invoke error"), time.Duration(0)) + rubyManager.On("invoke"). + Return(errors.New("fake invoke error"), time.Duration(0)). + Return(errors.New("fake invoke error"), time.Duration(0)) return fields{ @@ -192,8 +194,8 @@ func TestRubyProfiler_Invoke(t *testing.T) { }, then: func(t *testing.T, err error, fields fields) { require.Error(t, err) - assert.EqualError(t, err, "fake invoke error") - assert.Equal(t, 1, fields.RubyProfiler.RubyManager.(FakeRubyManager).On("invoke").InvokedTimes()) + assert.EqualError(t, err, "PID 1000: fake invoke error; PID 2000: fake invoke error") + assert.Equal(t, 2, fields.RubyProfiler.RubyManager.(FakeRubyManager).On("invoke").InvokedTimes()) }, }, } From 871e50dcd17a7c4d594a892fb86e1f54e7e92001 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Sat, 3 Oct 2026 23:50:56 +0530 Subject: [PATCH 2/4] fix(profiler): fail the PID when its flamegraph cannot be rendered bpf, perf, py-spy and austin logged a flamegraph rendering error and returned nil without publishing anything. With a single PID the run then ended as a success with no result, leaving the caller waiting for a file that would never come. Return the error instead, so the PID counts as failed: the run still succeeds when another PID published, and otherwise fails with the reason. --- internal/agent/profiler/bpf.go | 5 ++-- internal/agent/profiler/bpf_test.go | 6 +++-- internal/agent/profiler/perf.go | 5 ++-- internal/agent/profiler/perf_test.go | 6 +++-- internal/agent/profiler/python.go | 5 ++-- internal/agent/profiler/python_austin.go | 5 ++-- .../profiler/python_austin_output_test.go | 27 +++++++++++++++++++ internal/agent/profiler/python_test.go | 6 +++-- 8 files changed, 51 insertions(+), 14 deletions(-) diff --git a/internal/agent/profiler/bpf.go b/internal/agent/profiler/bpf.go index d041277..cf06bc8 100644 --- a/internal/agent/profiler/bpf.go +++ b/internal/agent/profiler/bpf.go @@ -105,8 +105,9 @@ func (b *bpfManager) invoke(job *job.ProfilingJob, pid string) (error, time.Dura err = b.handleFlamegraph(job, flamegraph.Get(job), fileName, resultFileName) if err != nil { - log.ErrorLogLn(fmt.Sprintf("could not generate flamegraph (PID: %s): %s", pid, err.Error())) - return nil, time.Since(start) + // Not nil: nothing is published for this PID, and a run that + // ends without a result must not report success. + return errors.Wrap(err, "could not generate flamegraph"), time.Since(start) } return b.publisher.Do(job.Compressor, resultFileName, job.OutputType), time.Since(start) diff --git a/internal/agent/profiler/bpf_test.go b/internal/agent/profiler/bpf_test.go index 266593b..212f74d 100644 --- a/internal/agent/profiler/bpf_test.go +++ b/internal/agent/profiler/bpf_test.go @@ -359,7 +359,7 @@ func Test_bpfManager_invoke(t *testing.T) { }, }, { - name: "should invoke return nil when fail handle flamegraph", + name: "should invoke fail when handle flamegraph fails", given: func() (fields, args) { log.SetPrintLogs(true) commander := executil.NewFakeCommander() @@ -385,7 +385,9 @@ func Test_bpfManager_invoke(t *testing.T) { return fields.BpfProfiler.invoke(args.job, args.pid) }, then: func(t *testing.T, fields fields, err error) { - require.NoError(t, err) + // nothing is published, so the PID must count as failed + require.Error(t, err) + assert.ErrorContains(t, err, "could not generate flamegraph") assert.True(t, fields.BpfProfiler.BpfManager.(*bpfManager).publisher.(*publish.Fake).On("Do").InvokedTimes() == 0) }, }, diff --git a/internal/agent/profiler/perf.go b/internal/agent/profiler/perf.go index fd0c222..1d9a95e 100644 --- a/internal/agent/profiler/perf.go +++ b/internal/agent/profiler/perf.go @@ -104,8 +104,9 @@ func (m *perfManager) invoke(job *job.ProfilingJob, pid string) (error, time.Dur err = m.handleFlamegraph(job, flamegraph.Get(job), fileName, resultFileName) if err != nil { - log.ErrorLogLn(fmt.Sprintf("could not generate flamegraph (PID: %s): %s", pid, err.Error())) - return nil, time.Since(start) + // Not nil: nothing is published for this PID, and a run that + // ends without a result must not report success. + return errors.Wrap(err, "could not generate flamegraph"), time.Since(start) } return m.publisher.Do(job.Compressor, resultFileName, job.OutputType), time.Since(start) diff --git a/internal/agent/profiler/perf_test.go b/internal/agent/profiler/perf_test.go index 1b895e7..9efe007 100644 --- a/internal/agent/profiler/perf_test.go +++ b/internal/agent/profiler/perf_test.go @@ -433,7 +433,7 @@ func Test_perfManager_invoke(t *testing.T) { }, }, { - name: "should invoke return nil when fail handle flamegraph", + name: "should invoke fail when handle flamegraph fails", given: func() (fields, args) { log.SetPrintLogs(true) commander := executil.NewFakeCommander() @@ -462,7 +462,9 @@ func Test_perfManager_invoke(t *testing.T) { return fields.PerfProfiler.invoke(args.job, args.pid) }, then: func(t *testing.T, fields fields, err error) { - require.NoError(t, err) + // nothing is published, so the PID must count as failed + require.Error(t, err) + assert.ErrorContains(t, err, "could not generate flamegraph") assert.True(t, fields.PerfProfiler.PerfManager.(*perfManager).publisher.(*publish.Fake).On("Do").InvokedTimes() == 0) assert.True(t, fields.PerfProfiler.PerfManager.(*perfManager).commander.(*executil.Fake).On("Command").InvokedTimes() == 3) }, diff --git a/internal/agent/profiler/python.go b/internal/agent/profiler/python.go index 15ee557..5612abb 100644 --- a/internal/agent/profiler/python.go +++ b/internal/agent/profiler/python.go @@ -123,8 +123,9 @@ func (p *pythonManager) invoke(job *job.ProfilingJob, pid string) (error, time.D } else { err = p.handleFlamegraph(job, flamegraph.Get(job), fileName, resultFileName) if err != nil { - log.ErrorLogLn(fmt.Sprintf("could not generate flamegraph (PID: %s): %s", pid, err.Error())) - return nil, time.Since(start) + // Not nil: nothing is published for this PID, and a run that + // ends without a result must not report success. + return errors.Wrap(err, "could not generate flamegraph"), time.Since(start) } } diff --git a/internal/agent/profiler/python_austin.go b/internal/agent/profiler/python_austin.go index 1fb2728..9f4a95f 100644 --- a/internal/agent/profiler/python_austin.go +++ b/internal/agent/profiler/python_austin.go @@ -292,8 +292,9 @@ func (p *austinPythonManager) invoke(job *job.ProfilingJob, pid string) (error, if job.OutputType == api.FlameGraph { err = p.handleFlamegraph(job, flamegraph.Get(job), fileName, resultFileName) if err != nil { - log.ErrorLogLn(fmt.Sprintf("could not generate flamegraph (PID: %s): %s", pid, err.Error())) - return nil, time.Since(start) + // Not nil: nothing is published for this PID, and a run that + // ends without a result must not report success. + return errors.Wrap(err, "could not generate flamegraph"), time.Since(start) } } // Nothing to do for the other output types: austin has already written diff --git a/internal/agent/profiler/python_austin_output_test.go b/internal/agent/profiler/python_austin_output_test.go index 8d84746..b0f9da0 100644 --- a/internal/agent/profiler/python_austin_output_test.go +++ b/internal/agent/profiler/python_austin_output_test.go @@ -262,3 +262,30 @@ func Test_invoke_noMemoryGrowth(t *testing.T) { assert.Equal(t, 0, publisher.On("Do").InvokedTimes()) }) } + +// Test_invoke_flamegraphFailure — a flamegraph that cannot be rendered leaves +// nothing to publish, so the PID must fail rather than end the run as a +// success without a result. +func Test_invoke_flamegraphFailure(t *testing.T) { + tmp := t.TempDir() + oldTmp := common.TmpDir + common.TmpDir = func() string { return tmp } + t.Cleanup(func() { common.TmpDir = oldTmp }) + useProcDir(t, t.TempDir()) + + // No language: flamegraph.Get returns a renderer that always fails. + j := &job.ProfilingJob{Tool: api.Austin, OutputType: api.FlameGraph, Interval: 30 * time.Second, Iteration: 1} + raw := common.GetResultFile(tmp, j.Tool, api.Raw, "42", j.Iteration) + commander := executil.NewFakeCommander() + commander.On("Command").Return(exec.Command("sh", "-c", + "printf '# austin: 4.0.0\\n# mode: memory\\nP1;T0:1;/app.py::6 131072\\n' > "+raw)) + publisher := publish.NewFakePublisher() + publisher.On("Do").Return(nil) + m := &austinPythonManager{commander: commander, publisher: publisher} + + err, _ := m.invoke(j, "42") + + require.Error(t, err) + assert.EqualError(t, err, "could not generate flamegraph: could not convert raw format to flamegraph: StackSamplesToFlameGraph with error") + assert.Equal(t, 0, publisher.On("Do").InvokedTimes()) +} diff --git a/internal/agent/profiler/python_test.go b/internal/agent/profiler/python_test.go index 5542cf6..ea8067a 100644 --- a/internal/agent/profiler/python_test.go +++ b/internal/agent/profiler/python_test.go @@ -396,7 +396,7 @@ func Test_pythonManager_invoke(t *testing.T) { }, }, { - name: "should invoke return nil when fail handle flamegraph", + name: "should invoke fail when handle flamegraph fails", given: func() (fields, args) { log.SetPrintLogs(true) commander := executil.NewFakeCommander() @@ -422,7 +422,9 @@ func Test_pythonManager_invoke(t *testing.T) { return fields.PythonProfiler.invoke(args.job, args.pid) }, then: func(t *testing.T, fields fields, err error) { - require.NoError(t, err) + // nothing is published, so the PID must count as failed + require.Error(t, err) + assert.ErrorContains(t, err, "could not generate flamegraph") assert.True(t, fields.PythonProfiler.PythonManager.(*pythonManager).publisher.(*publish.Fake).On("Do").InvokedTimes() == 0) }, }, From 2bb038ea7e5cd1b98aee449e33445e82e216c147 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Sat, 3 Oct 2026 23:52:28 +0530 Subject: [PATCH 3/4] fix(go): say why a pprof scrape failed The pprof profiler reported every failed scrape as the raw nsenter+wget command line and stderr, and when it could not find the target's listening port it silently scraped :8080 instead, so a failure there read as if the target's own pprof port were down. Map the failures busybox wget reports to the cause: - HTTP 401/403: the pprof endpoint on : requires authentication - HTTP 404: no /debug/pprof handler on : (is net/http/pprof registered?) - connection refused: nothing listening on : Anything else is passed through as before. :8080 is still tried when the port cannot be detected, but a failure there now says the port could not be detected and the default was tried. The scrape now runs through the profiler's commander, like the other profilers, so it can be exercised in tests. --- internal/agent/profiler/go_pprof.go | 126 ++++++++++++------ internal/agent/profiler/go_pprof_test.go | 157 +++++++++++++++++++++++ 2 files changed, 241 insertions(+), 42 deletions(-) create mode 100644 internal/agent/profiler/go_pprof_test.go diff --git a/internal/agent/profiler/go_pprof.go b/internal/agent/profiler/go_pprof.go index 71f3d6a..a780fd5 100644 --- a/internal/agent/profiler/go_pprof.go +++ b/internal/agent/profiler/go_pprof.go @@ -5,6 +5,7 @@ import ( "bytes" "fmt" "os/exec" + "regexp" "strconv" "strings" "time" @@ -22,7 +23,16 @@ import ( "github.com/pkg/errors" ) -const goPprofDelayBetweenJobs = 2 * time.Second +const ( + goPprofDelayBetweenJobs = 2 * time.Second + // defaultPprofPort is scraped when the target's listening port cannot be + // detected. + defaultPprofPort = "8080" +) + +// wgetHTTPStatus finds the status code in busybox wget's report of an HTTP +// error, e.g. "wget: server returned error: HTTP/1.1 401 Unauthorized". +var wgetHTTPStatus = regexp.MustCompile(`HTTP/\d(?:\.\d)? (\d{3})`) // GoPprofProfiler uses Go's pprof HTTP server to collect profiles. type GoPprofProfiler struct { @@ -75,8 +85,6 @@ func (p *GoPprofProfiler) Invoke(job *job.ProfilingJob) (error, time.Duration) { } func (p *goPprofManager) heapProfile(job *job.ProfilingJob, port string, fileName string) error { - var out bytes.Buffer - var stderr bytes.Buffer // A heap profile is a snapshot of what the process is holding right now, so // it takes no duration. Passing ?seconds= makes Go return a DELTA over that // window instead — it samples, waits, samples again and subtracts — which @@ -88,45 +96,63 @@ func (p *goPprofManager) heapProfile(job *job.ProfilingJob, port string, fileNam // rather than whatever the last GC cycle happened to leave behind; without // it a process that has not GC'd yet reports nothing at all. targetURL := fmt.Sprintf("http://127.0.0.1:%s/debug/pprof/heap?gc=1", port) - // for local testing - // cmd := exec.Command( - // "curl", targetURL, "-o", fileName, - // ) - cmd := exec.Command( - "nsenter", "-t", job.PID, "-n", "wget", "-qO", fileName, targetURL, - ) - cmd.Stdout = &out - cmd.Stderr = &stderr - err := cmd.Run() - if err != nil { - log.ErrorLogLn(out.String()) - return errors.Wrapf(err, "failed to nsenter+wget %q error %s", targetURL, stderr.String()) - } - return nil + return p.scrape(job.PID, port, targetURL, fileName) } func (p *goPprofManager) cpuProfile(job *job.ProfilingJob, port string, fileName string) error { - var out bytes.Buffer - var stderr bytes.Buffer targetURL := fmt.Sprintf( "http://127.0.0.1:%s/debug/pprof/%s?seconds=%d", port, "profile", int(job.Interval.Seconds()), ) + return p.scrape(job.PID, port, targetURL, fileName) +} + +// scrape downloads targetURL into fileName from inside the PID's network +// namespace, where 127.0.0.1 is the target's own loopback. +func (p *goPprofManager) scrape(pid string, port string, targetURL string, fileName string) error { + var out bytes.Buffer + var stderr bytes.Buffer // for local testing // cmd := exec.Command( // "curl", targetURL, "-o", fileName, // ) - cmd := exec.Command( - "nsenter", "-t", job.PID, "-n", "wget", "-qO", fileName, targetURL, + cmd := p.commander.Command( + "nsenter", "-t", pid, "-n", "wget", "-qO", fileName, targetURL, ) cmd.Stdout = &out cmd.Stderr = &stderr err := cmd.Run() if err != nil { log.ErrorLogLn(out.String()) - return errors.Wrapf(err, "failed to nsenter+wget %q error %s", targetURL, stderr.String()) + return pprofScrapeError(port, targetURL, stderr.String(), err) } return nil } + +// pprofScrapeError turns a failed scrape into a reason the user can act on. +// wget exits 1 for every failure, so the cause can only be read from its +// stderr. The images ship busybox wget, which prints the failure even with +// -q: +// +// wget: server returned error: HTTP/1.1 401 Unauthorized +// wget: can't connect to remote host (127.0.0.1): Connection refused +// +// Anything else is passed through as wget reported it. +func pprofScrapeError(port string, targetURL string, stderr string, err error) error { + msg := strings.TrimSpace(stderr) + if m := wgetHTTPStatus.FindStringSubmatch(msg); m != nil { + switch m[1] { + case "401", "403": + return errors.Errorf("the pprof endpoint on :%s requires authentication", port) + case "404": + return errors.Errorf("no /debug/pprof handler on :%s (is net/http/pprof registered?)", port) + } + } + if strings.Contains(msg, "Connection refused") { + return errors.Errorf("nothing listening on :%s", port) + } + return errors.Wrapf(err, "failed to nsenter+wget %q error %s", targetURL, msg) +} + func (p *goPprofManager) convertPprofToRaw(pprofFilePath string) (string, error) { // Convert the pprof output to raw format using go tool pprof var out bytes.Buffer @@ -144,53 +170,69 @@ func (p *goPprofManager) convertPprofToRaw(pprofFilePath string) (string, error) } func (m *goPprofManager) fetchProfileFromPID(job *job.ProfilingJob) error { - port, err := findListeningPortForPID(job.PID) - if err != nil { - log.ErrorLogLn(fmt.Sprintf("failed to find listening port for PID %s: %s", job.PID, err)) - port = "8080" + port, portErr := findListeningPortForPID(m.commander, job.PID) + if portErr != nil { + log.ErrorLogLn(fmt.Sprintf("failed to find listening port for PID %s: %s", job.PID, portErr)) + port = defaultPprofPort log.DebugLogLn(fmt.Sprintf("using default port %s", port)) } + err := m.fetchProfile(job, port) + if err != nil && portErr != nil { + // The port was a guess, not one the target listens on. Say so, or + // "nothing listening on :8080" reads as if the target's own pprof + // port were down. + return errors.Wrapf(err, "could not detect the listening port (%s), so tried the default :%s", + portErr, port) + } + return err +} + +// fetchProfile scrapes the profile the job asks for from port and publishes +// it. The PID is not repeated in the errors: the caller reports each PID's +// failure under its PID. +func (m *goPprofManager) fetchProfile(job *job.ProfilingJob, port string) error { rawFilePath := common.GetResultFile(common.TmpDir(), job.Tool, job.OutputType, job.PID, job.Iteration) - if job.OutputType == api.HeapDump { - err = m.heapProfile(job, port, rawFilePath) + switch job.OutputType { + case api.HeapDump: + err := m.heapProfile(job, port, rawFilePath) if err != nil { - return errors.Wrapf(err, "failed to fetch heap profile for PID %s", job.PID) + return errors.Wrap(err, "failed to fetch heap profile") } - } else if job.OutputType == api.Pprof { - err = m.cpuProfile(job, port, rawFilePath) + case api.Pprof: + err := m.cpuProfile(job, port, rawFilePath) if err != nil { - return errors.Wrapf(err, "failed to fetch CPU profile for PID %s", job.PID) + return errors.Wrap(err, "failed to fetch CPU profile") } - } else if job.OutputType == api.Raw { + case api.Raw: profileFilePath := common.GetResultFile(common.TmpDir(), job.Tool, "cpu", job.PID, job.Iteration) - err = m.cpuProfile(job, port, profileFilePath) + err := m.cpuProfile(job, port, profileFilePath) if err != nil { - return errors.Wrapf(err, "failed to create CPU profile for PID %s", job.PID) + return errors.Wrap(err, "failed to create CPU profile") } profileRaw, err := m.convertPprofToRaw(profileFilePath) if err != nil { - return errors.Wrapf(err, "failed to convert CPU profile for PID %s", job.PID) + return errors.Wrap(err, "failed to convert CPU profile") } heapFilePath := common.GetResultFile(common.TmpDir(), job.Tool, "heap", job.PID, job.Iteration) err = m.heapProfile(job, port, heapFilePath) if err != nil { - return errors.Wrapf(err, "failed to fetch heap profile for PID %s", job.PID) + return errors.Wrap(err, "failed to fetch heap profile") } heapRaw, err := m.convertPprofToRaw(heapFilePath) if err != nil { - return errors.Wrapf(err, "failed to convert heap profile for PID %s", job.PID) + return errors.Wrap(err, "failed to convert heap profile") } file.Write(rawFilePath, fmt.Sprintf("heap dump\n %s \n cpu dump\n %s", heapRaw, profileRaw)) - } else { + default: return errors.New("unsupported output type for Go pprof profiler") } // Finally, publish the file as before return m.publisher.Do(job.Compressor, rawFilePath, job.OutputType) } -func findListeningPortForPID(pid string) (string, error) { +func findListeningPortForPID(commander executil.Commander, pid string) (string, error) { // nsenter into the PID's network namespace and list listening TCP sockets - cmd := exec.Command("nsenter", "-t", pid, "-n", "ss", "-tulnp") + cmd := commander.Command("nsenter", "-t", pid, "-n", "ss", "-tulnp") output, err := cmd.Output() if err != nil { return "", errors.Wrap(err, "failed to run ss") diff --git a/internal/agent/profiler/go_pprof_test.go b/internal/agent/profiler/go_pprof_test.go new file mode 100644 index 0000000..4156398 --- /dev/null +++ b/internal/agent/profiler/go_pprof_test.go @@ -0,0 +1,157 @@ +package profiler + +import ( + "errors" + "os/exec" + "testing" + "time" + + "github.com/nudgebee/application-profiler/api" + "github.com/nudgebee/application-profiler/internal/agent/job" + "github.com/nudgebee/application-profiler/internal/agent/profiler/common" + executil "github.com/nudgebee/application-profiler/internal/agent/util/exec" + "github.com/nudgebee/application-profiler/internal/agent/util/publish" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Test_pprofScrapeError uses the stderr busybox wget (what the images ship) +// actually prints with -q for each failure. +func Test_pprofScrapeError(t *testing.T) { + exitErr := errors.New("exit status 1") + const url = "http://127.0.0.1:6060/debug/pprof/profile?seconds=30" + tests := []struct { + name string + stderr string + want string + }{ + { + name: "401 asks for credentials", + stderr: "wget: server returned error: HTTP/1.1 401 Unauthorized\n", + want: "the pprof endpoint on :6060 requires authentication", + }, + { + name: "403 is treated the same", + stderr: "wget: server returned error: HTTP/1.1 403 Forbidden\n", + want: "the pprof endpoint on :6060 requires authentication", + }, + { + name: "404 means pprof is not registered on that port", + stderr: "wget: server returned error: HTTP/1.0 404 Not Found\n", + want: "no /debug/pprof handler on :6060 (is net/http/pprof registered?)", + }, + { + name: "connection refused means nothing listens there", + stderr: "wget: can't connect to remote host (127.0.0.1): Connection refused\n", + want: "nothing listening on :6060", + }, + { + name: "anything else is passed through", + stderr: "wget: server returned error: HTTP/1.1 500 Internal Server Error\n", + want: `failed to nsenter+wget "` + url + `" error wget: server returned error: ` + + `HTTP/1.1 500 Internal Server Error: exit status 1`, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.EqualError(t, pprofScrapeError("6060", url, tt.stderr, exitErr), tt.want) + }) + } +} + +// ssOutput fakes `ss -tulnp` run in the target's network namespace. +func ssOutput(lines ...string) *exec.Cmd { + args := append([]string{"%s\n", "Netid State Recv-Q Send-Q Local Address:Port Peer Address:Port Process"}, lines...) + return exec.Command("printf", args...) +} + +// wgetFailing fakes wget failing with the given stderr. +func wgetFailing(stderr string) *exec.Cmd { + return exec.Command("sh", "-c", "printf '%s\\n' \"$1\" >&2; exit 1", "sh", stderr) +} + +const listeningOn6060 = `tcp LISTEN 0 4096 *:6060 *:* users:(("server",pid=42,fd=3))` + +func Test_goPprofManager_fetchProfileFromPID(t *testing.T) { + setup := func(t *testing.T, output api.OutputType, cmds ...*exec.Cmd) (*goPprofManager, publish.FakePublisher, *job.ProfilingJob) { + t.Helper() + tmp := t.TempDir() + oldTmp := common.TmpDir + common.TmpDir = func() string { return tmp } + t.Cleanup(func() { common.TmpDir = oldTmp }) + + commander := executil.NewFakeCommander() + for _, c := range cmds { + commander.On("Command").Return(c) + } + publisher := publish.NewFakePublisher() + publisher.On("Do").Return(nil) + j := &job.ProfilingJob{Tool: api.PProf, OutputType: output, Interval: 30 * time.Second, Iteration: 1, PID: "42"} + return &goPprofManager{commander: commander, publisher: publisher}, publisher, j + } + + t.Run("an endpoint behind auth says so", func(t *testing.T) { + m, publisher, j := setup(t, api.Pprof, + ssOutput(listeningOn6060), wgetFailing("wget: server returned error: HTTP/1.1 401 Unauthorized")) + + err := m.fetchProfileFromPID(j) + + assert.EqualError(t, err, "failed to fetch CPU profile: the pprof endpoint on :6060 requires authentication") + assert.Equal(t, 0, publisher.On("Do").InvokedTimes()) + }) + + t.Run("a heap profile without pprof registered says so", func(t *testing.T) { + m, publisher, j := setup(t, api.HeapDump, + ssOutput(listeningOn6060), wgetFailing("wget: server returned error: HTTP/1.1 404 Not Found")) + + err := m.fetchProfileFromPID(j) + + assert.EqualError(t, err, + "failed to fetch heap profile: no /debug/pprof handler on :6060 (is net/http/pprof registered?)") + assert.Equal(t, 0, publisher.On("Do").InvokedTimes()) + }) + + t.Run("an undetected port says the default was a guess", func(t *testing.T) { + m, publisher, j := setup(t, api.Pprof, + ssOutput(), wgetFailing("wget: can't connect to remote host (127.0.0.1): Connection refused")) + + err := m.fetchProfileFromPID(j) + + assert.EqualError(t, err, "could not detect the listening port (no listening port found for PID), "+ + "so tried the default :8080: failed to fetch CPU profile: nothing listening on :8080") + assert.Equal(t, 0, publisher.On("Do").InvokedTimes()) + }) + + t.Run("an undetected port still profiles the default when it answers", func(t *testing.T) { + m, publisher, j := setup(t, api.Pprof, ssOutput(), exec.Command("true")) + + require.NoError(t, m.fetchProfileFromPID(j)) + assert.Equal(t, 1, publisher.On("Do").InvokedTimes()) + }) +} + +// TestGoPprofProfiler_Invoke — the reason reaches the run's error under its +// PID, and the shared job is left alone (each PID works on its own copy). +func TestGoPprofProfiler_Invoke(t *testing.T) { + tmp := t.TempDir() + oldTmp := common.TmpDir + common.TmpDir = func() string { return tmp } + t.Cleanup(func() { common.TmpDir = oldTmp }) + + commander := executil.NewFakeCommander() + commander.On("Command"). + Return(ssOutput()). + Return(wgetFailing("wget: can't connect to remote host (127.0.0.1): Connection refused")) + publisher := publish.NewFakePublisher() + publisher.On("Do").Return(nil) + p := NewGoPprofProfiler(commander, publisher) + p.delay = 0 + p.targetPIDs = []string{"42"} + j := &job.ProfilingJob{Tool: api.PProf, OutputType: api.Pprof, Interval: 30 * time.Second, Iteration: 1} + + err, _ := p.Invoke(j) + + assert.EqualError(t, err, "PID 42: could not detect the listening port (no listening port found for PID), "+ + "so tried the default :8080: failed to fetch CPU profile: nothing listening on :8080") + assert.Empty(t, j.PID) +} From 65be9ae06a955883c604526dce53c3e5ee15f4f5 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Sun, 4 Oct 2026 00:01:46 +0530 Subject: [PATCH 4/4] test(e2e): profile multi-process Python targets and real Go servers python-austin-multi-pid.sh runs a Python program next to a non-Python process, in both start orders, and lets the agent find the container's leaf PIDs itself. The Python PID's profile must be published, the other PID named in a notice, and the run must not fail. go-pprof.sh builds a small Go server and scrapes it with the bpf image's own wget: with pprof registered a profile is published; behind basic auth, without a pprof handler, and with nothing listening, the agent's error must give the matching reason. Both run in Code Verify. --- .github/workflows/code-verify.yml | 19 ++++ test/e2e/go-pprof.sh | 151 ++++++++++++++++++++++++++++ test/e2e/python-austin-multi-pid.sh | 130 ++++++++++++++++++++++++ 3 files changed, 300 insertions(+) create mode 100755 test/e2e/go-pprof.sh create mode 100755 test/e2e/python-austin-multi-pid.sh diff --git a/.github/workflows/code-verify.yml b/.github/workflows/code-verify.yml index b3f6e62..78091cc 100644 --- a/.github/workflows/code-verify.yml +++ b/.github/workflows/code-verify.yml @@ -76,3 +76,22 @@ jobs: - name: Profile real Python targets run: test/e2e/python-austin.sh application-profiler-python:e2e + + # A PID the profiler cannot attach to (a shell or helper process next + # to the workload) used to fail the whole run. + - name: Profile a Python target next to a non-Python process + run: test/e2e/python-austin-multi-pid.sh application-profiler-python:e2e + + # The pprof profiler reads the cause of a failed scrape from wget's stderr, + # which unit tests can only fake. Scrape real Go servers with the image's + # own wget instead. + go-pprof-e2e: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v6 + + - name: Build the bpf profiler image + run: docker build -f docker/bpf/Dockerfile -t application-profiler-bpf:e2e . + + - name: Profile real Go servers + run: test/e2e/go-pprof.sh application-profiler-bpf:e2e diff --git a/test/e2e/go-pprof.sh b/test/e2e/go-pprof.sh new file mode 100755 index 0000000..b7af884 --- /dev/null +++ b/test/e2e/go-pprof.sh @@ -0,0 +1,151 @@ +#!/usr/bin/env bash +# End-to-end check of Go pprof profiling against real Go servers. +# +# Usage: test/e2e/go-pprof.sh +# +# The agent scrapes /debug/pprof with wget from inside the target's network +# namespace, so what it can report depends on what the image's wget prints. +# Each target is a small Go server built here; the agent runs from the +# profiler image with --pid=container:, as in python-austin.sh: +# - pprof registered: a profile is published +# - pprof behind basic auth: the error says authentication is required +# - no pprof handler: the error says net/http/pprof is not registered +# - nothing listening: the error says the port could not be detected and +# the default was tried +set -euo pipefail + +image=${1:?usage: $0 } +server=application-profiler-e2e-go-server +work=$(mktemp -d) +containers=() +cleanup() { + if [ ${#containers[@]} -gt 0 ]; then docker rm -f "${containers[@]}" >/dev/null 2>&1 || true; fi + rm -rf "$work" +} +trap cleanup EXIT + +# server +cat >"$work/main.go" <<'GO' +package main + +import ( + "log" + "net/http" + "net/http/pprof" + "os" +) + +func main() { + mode, addr := os.Args[1], os.Args[2] + + profiles := http.NewServeMux() + profiles.HandleFunc("/debug/pprof/", pprof.Index) + profiles.HandleFunc("/debug/pprof/profile", pprof.Profile) + + mux := http.NewServeMux() + mux.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) { _, _ = w.Write([]byte("ok")) }) + switch mode { + case "open": + mux.Handle("/debug/pprof/", profiles) + case "auth": + mux.Handle("/debug/pprof/", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if user, pass, ok := r.BasicAuth(); !ok || user != "user" || pass != "secret" { + w.Header().Set("WWW-Authenticate", `Basic realm="pprof"`) + http.Error(w, "unauthorized", http.StatusUnauthorized) + return + } + profiles.ServeHTTP(w, r) + })) + } + log.Fatal(http.ListenAndServe(addr, mux)) +} +GO +cat >"$work/Dockerfile" <<'DOCKERFILE' +FROM golang:1.26.5-alpine AS build +WORKDIR /src +COPY main.go . +RUN go mod init example.com/server && CGO_ENABLED=0 go build -o /server . + +FROM alpine:3.23.4 +COPY --from=build /server /server +ENTRYPOINT ["/server"] +DOCKERFILE +docker build -q -t "$server" "$work" >/dev/null + +# event : prints data. of every agent event of that type. +event() { + jq -rR --arg type "$1" --arg field "$2" 'fromjson? | select(.type == $type) | .data[$field]' +} + +start_target() { # name args... + local name=$1 + shift + docker run -d --name "$name" "$@" >/dev/null + containers+=("$name") +} + +# profile : prints the published file, or the agent's error reason on +# stderr and returns 1. The file is copied to $work/.gz. +profile() { + local agent logs file + agent=$(docker run -d --privileged --pid="container:$1" --entrypoint /app/agent "$image" \ + --target-container-id e2e --target-container-runtime containerd \ + --target-container-runtime-path /run/containerd --pid 1 --lang go \ + --profiling-tool pprof --output-type pprof --duration 5s \ + --grace-period-ending 30s --compressor-type gzip) + containers+=("$agent") + for _ in $(seq 1 60); do + docker logs "$agent" 2>&1 | grep -q '"stage":"ended"\|"type":"error"' && break + sleep 1 + done + logs=$(docker logs "$agent" 2>&1) + if grep -q '"type":"error"' <<<"$logs"; then + event error reason <<<"$logs" >&2 + return 1 + fi + file=$(event result file <<<"$logs") + if [ -z "$file" ]; then + echo "neither a result nor an error within 60s" >&2 + return 1 + fi + docker cp -q "$agent:$file" "$work/$1.gz" + echo "$file" +} + +start_target t-pprof-open "$server" open :6060 +start_target t-pprof-auth "$server" auth :6060 +start_target t-pprof-none "$server" none :9090 +start_target t-no-listener alpine:3.23.4 sleep 100000 +sleep 2 + +failed=0 + +# A pprof profile is itself gzipped protobuf, which the agent gzips again. +if file=$(profile t-pprof-open 2>"$work/err") && gzip -dc "$work/t-pprof-open.gz" | gzip -t; then + echo "ok t-pprof-open: published $file" +else + echo "FAIL t-pprof-open: no profile: $(cat "$work/err")" + failed=1 +fi + +expect_error() { # target want + local reason + if profile "$1" >/dev/null 2>"$work/err"; then + echo "FAIL $1: profiled, want an error containing: $2" + failed=1 + return + fi + reason=$(cat "$work/err") + if [[ $reason == *"$2"* ]]; then + echo "ok $1: $reason" + else + echo "FAIL $1: error '$reason', want it to contain: $2" + failed=1 + fi +} + +expect_error t-pprof-auth "PID 1: failed to fetch CPU profile: the pprof endpoint on :6060 requires authentication" +expect_error t-pprof-none "PID 1: failed to fetch CPU profile: no /debug/pprof handler on :9090 (is net/http/pprof registered?)" +expect_error t-no-listener "could not detect the listening port (no listening port found for PID), so tried the default :8080: failed to fetch CPU profile: nothing listening on :8080" + +exit "$failed" diff --git a/test/e2e/python-austin-multi-pid.sh b/test/e2e/python-austin-multi-pid.sh new file mode 100755 index 0000000..ed4bfc7 --- /dev/null +++ b/test/e2e/python-austin-multi-pid.sh @@ -0,0 +1,130 @@ +#!/usr/bin/env bash +# End-to-end check that a PID the profiler cannot attach to does not fail the +# run, using Python memory profiling (austin). +# +# Usage: test/e2e/python-austin-multi-pid.sh +# +# Each target runs a Python program next to a non-Python process, so the +# container's leaf processes — what the agent profiles when no --pid is +# given — are one austin can read and one it cannot. The run must publish the +# Python PID's profile, name the other PID in a notice, and not fail. PIDs +# are started one after another, so both orders are checked. +# +# The agent finds the container's root PID through the containerd runtime +# directory. A minimal stand-in for it, pointing at PID 1 of the target, is +# mounted where the agent looks; the agent runs with --pid=container: +# as in python-austin.sh, so that PID 1 is the target's. +set -euo pipefail + +image=${1:?usage: $0 } +work=$(mktemp -d) +containers=() +cleanup() { + if [ ${#containers[@]} -gt 0 ]; then docker rm -f "${containers[@]}" >/dev/null 2>&1 || true; fi + rm -rf "$work" +} +trap cleanup EXIT + +# Grows by 64 KiB every 5 ms, resetting at ~128 MiB so the container stays small. +cat >"$work/grow.py" <<'PY' +import time +hold = [] +while True: + hold.append(bytearray(64 * 1024)) + if len(hold) >= 2000: + hold = [] + time.sleep(0.005) +PY +chmod 0644 "$work/grow.py" + +# The agent reads the root PID from /io.containerd.runtime.v2.task/k8s.io//init.pid. +mkdir -p "$work/runtime/io.containerd.runtime.v2.task/k8s.io/e2e" +printf 1 >"$work/runtime/io.containerd.runtime.v2.task/k8s.io/e2e/init.pid" + +# event : prints data. of every agent event of that type. +event() { + jq -rR --arg type "$1" --arg field "$2" 'fromjson? | select(.type == $type) | .data[$field]' +} + +failed=0 +fail() { + echo "FAIL $1: $2" + failed=$((failed + 1)) +} + +# check : profiles the target, whose command starts a python +# and a sleep process, without --pid. +check() { + local target=$1 agent logs detected p py="" other="" files samples notice err + local -a pids + docker run -d --name "$target" -v "$work/grow.py:/app.py:ro" python:3.12-slim sh -c "$2" >/dev/null + containers+=("$target") + sleep 3 + + agent=$(docker run -d --privileged --pid="container:$target" \ + -v "$work/runtime:/run/containerd:ro" --entrypoint /app/agent "$image" \ + --target-container-id e2e --target-container-runtime containerd \ + --target-container-runtime-path /run/containerd --lang python \ + --profiling-tool austin --output-type raw --duration 5s \ + --grace-period-ending 30s --compressor-type gzip) + containers+=("$agent") + for _ in $(seq 1 60); do + docker logs "$agent" 2>&1 | grep -q '"stage":"ended"\|"type":"error"' && break + sleep 1 + done + logs=$(docker logs "$agent" 2>&1) + local before=$failed + + # Without two leaf PIDs there is nothing to tolerate, and the checks below + # would pass for the wrong reason. + detected=$(event notice msg <<<"$logs" | grep -o 'Detected more than one PID to profile: \[[0-9 ]*\]' | grep -o '[0-9 ]*\]' | tr -d ']' || true) + read -r -a pids <<<"$detected" + if [ ${#pids[@]} -eq 2 ]; then + for p in "${pids[@]}"; do + case $(docker exec "$target" cat "/proc/$p/comm") in + python*) py=$p ;; + sleep) other=$p ;; + esac + done + fi + if [ -z "$py" ] || [ -z "$other" ]; then + fail "$target" "want a python and a sleep leaf PID, the agent detected [${detected}]" + fi + + err=$(event error reason <<<"$logs") + if [ -n "$err" ]; then + fail "$target" "the run failed: $err" + fi + + files=$(event result file <<<"$logs") + if [ "$(grep -c . <<<"$files")" -ne 1 ]; then + fail "$target" "want one result, got: ${files:-none}" + elif [[ $files != *"-$py-"* ]]; then + fail "$target" "the result $files is not the python PID $py" + else + docker cp -q "$agent:$files" "$work/profile.gz" + samples=$(gzip -dc "$work/profile.gz" | grep -c '^P' || true) + if [ "$samples" -eq 0 ]; then + fail "$target" "the python PID's profile has no samples" + else + echo "ok $target: PID $py (python) published $files, $samples samples" + fi + fi + + notice=$(event notice msg <<<"$logs" | grep "^Profiled 1 of 2 PIDs; skipped PID $other: " || true) + if [ -z "$notice" ]; then + fail "$target" "no notice naming the skipped PID $other" + else + echo "ok $target: PID $other (sleep) skipped with a notice: $notice" + fi + + if [ "$failed" -ne "$before" ]; then + echo "--- $target: agent events" >&2 + grep -v '"type":"log"' <<<"$logs" >&2 || true + fi +} + +check t-python-first 'python /app.py & sleep 100000 & wait' +check t-sleep-first 'sleep 100000 & python /app.py & wait' + +exit $((failed > 0))