diff --git a/util/desktop/bundle/content.go b/util/desktop/bundle/content.go index 7a8a8fd41a29..54cc261e3e01 100644 --- a/util/desktop/bundle/content.go +++ b/util/desktop/bundle/content.go @@ -1,7 +1,9 @@ package bundle import ( + "bufio" "context" + "io" "github.com/containerd/containerd/v2/core/content" cerrdefs "github.com/containerd/errdefs" @@ -9,6 +11,35 @@ import ( ocispecs "github.com/opencontainers/image-spec/specs-go/v1" ) +// readChunkSize is the buffer size for blob reads. Each ReadAt on a remote +// content store starts a separate Content/Read RPC. +const readChunkSize = 4 * 1024 * 1024 // 4 MiB + +func newChunkedReader(ra content.ReaderAt) io.Reader { + return bufio.NewReaderSize(io.NewSectionReader(ra, 0, ra.Size()), int(min(ra.Size(), readChunkSize))) +} + +// chunkedProvider exposes buffered readers through content.NewReader's Reader hook. +type chunkedProvider struct { + content.InfoReaderProvider +} + +func (p *chunkedProvider) ReaderAt(ctx context.Context, desc ocispecs.Descriptor) (content.ReaderAt, error) { + ra, err := p.InfoReaderProvider.ReaderAt(ctx, desc) + if err != nil { + return nil, err + } + return &chunkedReaderAt{ReaderAt: ra}, nil +} + +type chunkedReaderAt struct { + content.ReaderAt +} + +func (r *chunkedReaderAt) Reader() io.Reader { + return newChunkedReader(r.ReaderAt) +} + type nsFallbackStore struct { main content.Store fb content.Store diff --git a/util/desktop/bundle/content_test.go b/util/desktop/bundle/content_test.go new file mode 100644 index 000000000000..4b86c968b6d1 --- /dev/null +++ b/util/desktop/bundle/content_test.go @@ -0,0 +1,49 @@ +package bundle + +import ( + "bytes" + "context" + "encoding/json" + "io" + "testing" + + "github.com/containerd/containerd/v2/core/content" + imgarchive "github.com/containerd/containerd/v2/core/images/archive" + "github.com/moby/buildkit/util/contentutil" + "github.com/opencontainers/go-digest" + "github.com/opencontainers/image-spec/specs-go" + ocispecs "github.com/opencontainers/image-spec/specs-go/v1" + "github.com/stretchr/testify/require" +) + +func writeTestBlob(t *testing.T, buf contentutil.Buffer, mt string, dt []byte) ocispecs.Descriptor { + t.Helper() + + desc := ocispecs.Descriptor{ + MediaType: mt, + Digest: digest.FromBytes(dt), + Size: int64(len(dt)), + } + require.NoError(t, content.WriteBlob(context.TODO(), buf, desc.Digest.String(), bytes.NewReader(dt), desc)) + return desc +} + +func TestArchiveExportReadsLayerOnce(t *testing.T) { + buf := contentutil.NewBuffer() + + logs := writeTestBlob(t, buf, "application/vnd.buildkit.status.v0", bytes.Repeat([]byte("#5 [stage 3/9] RUN apt-get install\n"), 52000)) + config := writeTestBlob(t, buf, HistoryRecordMediaTypeV0, []byte(`{}`)) + mfst, err := json.Marshal(ocispecs.Manifest{ + Versioned: specs.Versioned{SchemaVersion: 2}, + MediaType: ocispecs.MediaTypeImageManifest, + Config: config, + Layers: []ocispecs.Descriptor{logs}, + }) + require.NoError(t, err) + mfstDesc := writeTestBlob(t, buf, ocispecs.MediaTypeImageManifest, mfst) + + p := newCountingProvider(buf) + err = imgarchive.Export(context.TODO(), &chunkedProvider{p}, io.Discard, imgarchive.WithManifest(mfstDesc), imgarchive.WithSkipDockerManifest()) + require.NoError(t, err) + require.Equal(t, 1, p.reads[logs.Digest]) +} diff --git a/util/desktop/bundle/export.go b/util/desktop/bundle/export.go index ad716f6fbd8a..8a1dc5d02155 100644 --- a/util/desktop/bundle/export.go +++ b/util/desktop/bundle/export.go @@ -59,7 +59,7 @@ func Export(ctx context.Context, c []*client.Client, w io.Writer, records []*Rec gz := gzip.NewWriter(w) defer gz.Close() - if err := imgarchive.Export(ctx, mp, gz, imgarchive.WithManifest(desc), imgarchive.WithSkipDockerManifest()); err != nil { + if err := imgarchive.Export(ctx, &chunkedProvider{mp}, gz, imgarchive.WithManifest(desc), imgarchive.WithSkipDockerManifest()); err != nil { return errors.Wrap(err, "failed to create dockerbuild archive") } diff --git a/util/desktop/bundle/trace.go b/util/desktop/bundle/trace.go index 63f920372bbd..fc984cc8f2fe 100644 --- a/util/desktop/bundle/trace.go +++ b/util/desktop/bundle/trace.go @@ -45,7 +45,7 @@ func sanitizeTrace(ctx context.Context, mp *contentutil.MultiProvider, desc ocis defer ra.Close() buf := &bytes.Buffer{} - dec := json.NewDecoder(io.NewSectionReader(ra, 0, ra.Size())) + dec := json.NewDecoder(newChunkedReader(ra)) enc := json.NewEncoder(buf) enc.SetIndent("", " ") for { diff --git a/util/desktop/bundle/trace_test.go b/util/desktop/bundle/trace_test.go new file mode 100644 index 000000000000..e6f4df24370e --- /dev/null +++ b/util/desktop/bundle/trace_test.go @@ -0,0 +1,70 @@ +package bundle + +import ( + "bytes" + "context" + "os" + "testing" + + "github.com/containerd/containerd/v2/core/content" + "github.com/moby/buildkit/util/contentutil" + "github.com/opencontainers/go-digest" + ocispecs "github.com/opencontainers/image-spec/specs-go/v1" + "github.com/stretchr/testify/require" +) + +// countingProvider counts ReadAt calls per blob, each of which is one +// Content/Read RPC when the provider is a remote BuildKit content store. +type countingProvider struct { + content.InfoReaderProvider + reads map[digest.Digest]int +} + +func newCountingProvider(p content.InfoReaderProvider) *countingProvider { + return &countingProvider{InfoReaderProvider: p, reads: map[digest.Digest]int{}} +} + +func (p *countingProvider) ReaderAt(ctx context.Context, desc ocispecs.Descriptor) (content.ReaderAt, error) { + ra, err := p.InfoReaderProvider.ReaderAt(ctx, desc) + if err != nil { + return nil, err + } + return &countingReaderAt{ReaderAt: ra, count: func() { p.reads[desc.Digest]++ }}, nil +} + +type countingReaderAt struct { + content.ReaderAt + count func() +} + +func (r *countingReaderAt) ReadAt(p []byte, off int64) (int, error) { + r.count() + return r.ReaderAt.ReadAt(p, off) +} + +func TestSanitizeTraceReadsBlobOnce(t *testing.T) { + ctx := context.TODO() + + dt, err := os.ReadFile("../../otelutil/fixtures/bktraces.json") + require.NoError(t, err) + dt = bytes.Replace(dt, []byte(`"Value": "grpc"`), []byte(`"Value": "git fetch token=hunter2"`), 1) + + desc := ocispecs.Descriptor{ + MediaType: "application/json", + Digest: digest.FromBytes(dt), + Size: int64(len(dt)), + } + buf := contentutil.NewBuffer() + require.NoError(t, content.WriteBlob(ctx, buf, "trace", bytes.NewReader(dt), desc)) + + p := newCountingProvider(buf) + mp := contentutil.NewMultiProvider(p) + out, err := sanitizeTrace(ctx, mp, desc) + require.NoError(t, err) + require.Equal(t, 1, p.reads[desc.Digest]) + + sanitized, err := content.ReadBlob(ctx, mp, *out) + require.NoError(t, err) + require.Contains(t, string(sanitized), "token=xxxxx") + require.NotContains(t, string(sanitized), "hunter2") +}