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/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..cf06bc8 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) } @@ -124,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_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..212f74d 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()) }, }, } @@ -357,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() @@ -383,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/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..a780fd5 100644 --- a/internal/agent/profiler/go_pprof.go +++ b/internal/agent/profiler/go_pprof.go @@ -3,15 +3,14 @@ package profiler import ( "bufio" "bytes" - "context" "fmt" "os/exec" + "regexp" "strconv" "strings" "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 +23,22 @@ import ( "github.com/pkg/errors" ) +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 { manager *goPprofManager targetPIDs []string + delay time.Duration } type goPprofManager struct { @@ -37,7 +48,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,29 +73,18 @@ 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 // 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 @@ -93,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 @@ -149,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) +} 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..1d9a95e 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) } @@ -123,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_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..9efe007 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()) }, }, } @@ -431,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() @@ -460,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 1af18c9..5612abb 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) } @@ -142,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 557a6ac..9f4a95f 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) } @@ -311,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_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..ea8067a 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()) }, }, } @@ -394,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() @@ -420,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) }, }, 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()) }, }, } 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))