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
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ require (
github.com/oklog/run v1.2.0
github.com/olekukonko/tablewriter v0.0.5
github.com/parquet-go/parquet-go v0.24.0
github.com/pierrec/lz4/v4 v4.1.25
github.com/planetscale/vtprotobuf v0.6.1-0.20250313105119-ba97887b0a25
github.com/polarsignals/frostdb v0.0.0-20260121113628-9e5cfe0171ad
github.com/polarsignals/iceberg-go v0.0.0-20240502213135-2ee70b71e76b
Expand Down Expand Up @@ -219,7 +220,6 @@ require (
github.com/oracle/oci-go-sdk/v65 v65.41.1 // indirect
github.com/ovh/go-ovh v1.8.0 // indirect
github.com/paulmach/orb v0.12.0 // indirect
github.com/pierrec/lz4/v4 v4.1.25 // indirect
github.com/pkg/browser v0.0.0-20240102092130-5ac0b6a4141c // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect
Expand Down
53 changes: 49 additions & 4 deletions pkg/scrape/scrape.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,9 @@
package scrape

import (
"bufio"
"bytes"
"compress/gzip"
"context"
"errors"
"fmt"
Expand All @@ -28,6 +30,8 @@ import (
"github.com/go-kit/log"
"github.com/go-kit/log/level"
"github.com/google/pprof/profile"
"github.com/klauspost/compress/zstd"
"github.com/pierrec/lz4/v4"
"github.com/prometheus/client_golang/prometheus"
commonconfig "github.com/prometheus/common/config"
"github.com/prometheus/common/version"
Expand Down Expand Up @@ -365,16 +369,57 @@ func (s *targetScraper) scrape(ctx context.Context, w io.Writer, profileType str
case ProfileTraceType:
return fmt.Errorf("unimplemented")
default:
b, err := io.ReadAll(io.TeeReader(resp.Body, w))
if err != nil {
if err := readProfile(resp.Body, w); err != nil {
return fmt.Errorf("failed to read body: %w", err)
}
}

return nil
}

if len(b) == 0 {
return fmt.Errorf("empty %s profile from %s", profileType, s.req.URL.String())
const maxProfileSize = 100 * 1024 * 1024

func readProfile(r io.Reader, w io.Writer) error {
br := bufio.NewReader(r)
magic, err := br.Peek(4)
if err != nil && !errors.Is(err, io.EOF) {
return err
}

var profile io.Reader = br
var cleanup func()
switch {
case len(magic) >= 2 && magic[0] == 0x1f && magic[1] == 0x8b:
gz, err := gzip.NewReader(br)
if err != nil {
return err
}
profile = gz
cleanup = func() { _ = gz.Close() }
case len(magic) >= 4 && magic[0] == 0x04 && magic[1] == 0x22 && magic[2] == 0x4d && magic[3] == 0x18:
profile = lz4.NewReader(br)
case len(magic) >= 4 && magic[0] == 0x28 && magic[1] == 0xb5 && magic[2] == 0x2f && magic[3] == 0xfd:
decoder, err := zstd.NewReader(br)
if err != nil {
return err
}
profile = decoder
cleanup = decoder.Close
}
if cleanup != nil {
defer cleanup()
}

n, err := io.Copy(w, io.LimitReader(profile, maxProfileSize+1))
if err != nil {
return err
}
if n == 0 {
return errors.New("empty profile")
}
if n > maxProfileSize {
return fmt.Errorf("profile exceeds %d bytes", maxProfileSize)
}
return nil
}

Expand Down
45 changes: 45 additions & 0 deletions pkg/scrape/scrape_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,14 +14,59 @@
package scrape

import (
"bytes"
"compress/gzip"
"errors"
"testing"

"github.com/klauspost/compress/zstd"
"github.com/pierrec/lz4/v4"
"github.com/stretchr/testify/require"

profilepb "github.com/parca-dev/parca/gen/proto/go/parca/profilestore/v1alpha1"
)

func TestReadProfile(t *testing.T) {
raw := []byte("profile")
cases := map[string]func(*testing.T, *bytes.Buffer){
"raw": func(_ *testing.T, b *bytes.Buffer) {},
"gzip": func(t *testing.T, b *bytes.Buffer) {
w := gzip.NewWriter(b)
_, err := w.Write(raw)
require.NoError(t, err)
require.NoError(t, w.Close())
},
"lz4": func(t *testing.T, b *bytes.Buffer) {
w := lz4.NewWriter(b)
_, err := w.Write(raw)
require.NoError(t, err)
require.NoError(t, w.Close())
},
"zstd": func(t *testing.T, b *bytes.Buffer) {
w, err := zstd.NewWriter(b)
require.NoError(t, err)
_, err = w.Write(raw)
require.NoError(t, err)
require.NoError(t, w.Close())
},
}

for name, encode := range cases {
t.Run(name, func(t *testing.T) {
var encoded, decoded bytes.Buffer
if name != "raw" {
encode(t, &encoded)
}
input := bytes.NewReader(encoded.Bytes())
if name == "raw" {
input = bytes.NewReader(raw)
}
require.NoError(t, readProfile(input, &decoded))
require.Equal(t, raw, decoded.Bytes())
})
}
}

func TestParseExecutableInfo(t *testing.T) {
testCases := []struct {
name string
Expand Down
Loading