Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions .github/workflows/code-verify.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
1 change: 0 additions & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 0 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
32 changes: 7 additions & 25 deletions internal/agent/profiler/bpf.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
}

Expand All @@ -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)
Expand Down
5 changes: 5 additions & 0 deletions internal/agent/profiler/bpf_fake.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package profiler

import (
"sync"
"time"

"github.com/nudgebee/application-profiler/internal/agent/job"
Expand All @@ -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
}

Expand Down Expand Up @@ -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++
Expand Down
16 changes: 10 additions & 6 deletions internal/agent/profiler/bpf_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{
Expand All @@ -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())
},
},
}
Expand Down Expand Up @@ -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()
Expand All @@ -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)
},
},
Expand Down
87 changes: 87 additions & 0 deletions internal/agent/profiler/common/pids.go
Original file line number Diff line number Diff line change
@@ -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()
Comment thread
mayankpande88 marked this conversation as resolved.

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 <n>: <reason>; PID <m>: <reason>".
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
}
Loading
Loading