Skip to content
Open
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
31 changes: 31 additions & 0 deletions util/desktop/bundle/content.go
Original file line number Diff line number Diff line change
@@ -1,14 +1,45 @@
package bundle

import (
"bufio"
"context"
"io"

"github.com/containerd/containerd/v2/core/content"
cerrdefs "github.com/containerd/errdefs"
digest "github.com/opencontainers/go-digest"
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
Expand Down
49 changes: 49 additions & 0 deletions util/desktop/bundle/content_test.go
Original file line number Diff line number Diff line change
@@ -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])
}
2 changes: 1 addition & 1 deletion util/desktop/bundle/export.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}

Expand Down
2 changes: 1 addition & 1 deletion util/desktop/bundle/trace.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
70 changes: 70 additions & 0 deletions util/desktop/bundle/trace_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
Loading