From 8f2182458a96736a0d8531b557b1c6c8f1c24fe8 Mon Sep 17 00:00:00 2001 From: Edward McFarlane Date: Wed, 9 Sep 2026 22:07:10 +0100 Subject: [PATCH 1/2] Deduplicate Ref fetches --- CHANGELOG.md | 3 + private/buf/buffetch/buffetch.go | 3 + private/buf/buffetch/internal/reader.go | 135 ++++++-- private/buf/buffetch/internal/reader_cache.go | 90 +++++ .../buffetch/internal/reader_cache_test.go | 309 ++++++++++++++++++ private/pkg/git/git.go | 13 + 6 files changed, 526 insertions(+), 27 deletions(-) create mode 100644 private/buf/buffetch/internal/reader_cache.go create mode 100644 private/buf/buffetch/internal/reader_cache_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 76c945721b..943f207f03 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,9 @@ - Add `--stdin-filepath` flag to `buf format`, which reads a single `.proto` file from stdin and writes the formatted result to stdout. The path is not read from disk, and is only used to report parse errors and diffs. +- Deduplicate remote input fetches within a single command invocation, so that multiple + `inputs` in a `buf.gen.yaml` that resolve to the same archive, git repository, or image + are fetched once instead of once per input. ## [v1.73.0] - 2026-09-11 diff --git a/private/buf/buffetch/buffetch.go b/private/buf/buffetch/buffetch.go index a6758060e3..d8cdcf6334 100644 --- a/private/buf/buffetch/buffetch.go +++ b/private/buf/buffetch/buffetch.go @@ -414,6 +414,9 @@ type ModuleFetcher interface { } // Reader is a reader for Buf. +// +// A Reader fetches a given remote file or git repository at most once, for the +// lifetime of the Reader. type Reader interface { MessageReader SourceReader diff --git a/private/buf/buffetch/internal/reader.go b/private/buf/buffetch/internal/reader.go index eb3d8fb69a..2020696317 100644 --- a/private/buf/buffetch/internal/reader.go +++ b/private/buf/buffetch/internal/reader.go @@ -32,6 +32,7 @@ import ( "github.com/bufbuild/buf/private/buf/buftarget" "github.com/bufbuild/buf/private/bufpkg/bufmodule" "github.com/bufbuild/buf/private/bufpkg/bufparse" + "github.com/bufbuild/buf/private/pkg/cache" "github.com/bufbuild/buf/private/pkg/git" "github.com/bufbuild/buf/private/pkg/httpauth" "github.com/bufbuild/buf/private/pkg/normalpath" @@ -61,6 +62,11 @@ type reader struct { moduleEnabled bool moduleKeyProvider bufmodule.ModuleKeyProvider + + // Keyed so that Refs sharing a fetch share an entry, see reader_cache.go. + fileDataCache cache.Cache[fileDataCacheKey, []byte] + archiveBucketCache cache.Cache[archiveBucketCacheKey, storage.ReadBucket] + gitBucketCache cache.Cache[gitBucketCacheKey, storage.ReadBucket] } func newReader( @@ -244,11 +250,47 @@ func (r *reader) getArchiveBucket( targetPaths []string, targetExcludePaths []string, terminateFunc buftarget.TerminateFunc, -) (_ ReadBucketCloser, _ buftarget.BucketTargeting, retErr error) { - readCloser, size, err := r.getFileReadCloserAndSize(ctx, container, archiveRef, false) +) (ReadBucketCloser, buftarget.BucketTargeting, error) { + readBucket, err := r.getArchiveReadBucket(ctx, container, archiveRef) if err != nil { return nil, nil, err } + return getReadBucketCloserForBucket( + ctx, + r.logger, + storage.NopReadBucketCloser(readBucket), + archiveRef.SubDirPath(), + targetPaths, + targetExcludePaths, + terminateFunc, + ) +} + +func (r *reader) getArchiveReadBucket( + ctx context.Context, + container app.EnvStdinContainer, + archiveRef ArchiveRef, +) (storage.ReadBucket, error) { + if !isRemoteFileScheme(archiveRef.FileScheme()) { + return r.unarchive(ctx, container, archiveRef) + } + return r.archiveBucketCache.GetOrAdd( + newArchiveBucketCacheKey(archiveRef), + func() (storage.ReadBucket, error) { + return r.unarchive(ctx, container, archiveRef) + }, + ) +} + +func (r *reader) unarchive( + ctx context.Context, + container app.EnvStdinContainer, + archiveRef ArchiveRef, +) (_ storage.ReadBucket, retErr error) { + readCloser, size, err := r.getFileReadCloserAndSizeUncached(ctx, container, archiveRef, false) + if err != nil { + return nil, err + } defer func() { retErr = errors.Join(retErr, readCloser.Close()) }() @@ -263,21 +305,21 @@ func (r *reader) getArchiveBucket( archiveRef.StripComponents(), ), ); err != nil { - return nil, nil, err + return nil, err } case ArchiveTypeZip: var readerAt io.ReaderAt if size < 0 { data, err := io.ReadAll(readCloser) if err != nil { - return nil, nil, err + return nil, err } readerAt = bytes.NewReader(data) size = int64(len(data)) } else { readerAt, err = xio.ReaderAtForReader(readCloser) if err != nil { - return nil, nil, err + return nil, err } } if err := storagearchive.Unzip( @@ -289,20 +331,12 @@ func (r *reader) getArchiveBucket( archiveRef.StripComponents(), ), ); err != nil { - return nil, nil, err + return nil, err } default: - return nil, nil, fmt.Errorf("unknown ArchiveType: %v", archiveType) + return nil, fmt.Errorf("unknown ArchiveType: %v", archiveType) } - return getReadBucketCloserForBucket( - ctx, - r.logger, - storage.NopReadBucketCloser(readWriteBucket), - archiveRef.SubDirPath(), - targetPaths, - targetExcludePaths, - terminateFunc, - ) + return readWriteBucket, nil } func (r *reader) getDirBucket( @@ -359,10 +393,35 @@ func (r *reader) getGitBucket( if r.gitCloner == nil { return nil, nil, errors.New("git cloner is nil") } - gitURL, err := getGitURL(gitRef) + readBucket, err := r.gitBucketCache.GetOrAdd( + newGitBucketCacheKey(gitRef), + func() (storage.ReadBucket, error) { + return r.clone(ctx, container, gitRef) + }, + ) if err != nil { return nil, nil, err } + return getReadBucketCloserForBucket( + ctx, + r.logger, + storage.NopReadBucketCloser(readBucket), + gitRef.SubDirPath(), + targetPaths, + targetExcludePaths, + terminateFunc, + ) +} + +func (r *reader) clone( + ctx context.Context, + container app.EnvStdinContainer, + gitRef GitRef, +) (storage.ReadBucket, error) { + gitURL, err := getGitURL(gitRef) + if err != nil { + return nil, err + } readWriteBucket := storagemem.NewReadWriteBucket() if err := r.gitCloner.CloneToBucket( ctx, @@ -377,17 +436,9 @@ func (r *reader) getGitBucket( Filter: gitRef.Filter(), }, ); err != nil { - return nil, nil, fmt.Errorf("could not clone %s: %v", gitURL, err) + return nil, fmt.Errorf("could not clone %s: %v", gitURL, err) } - return getReadBucketCloserForBucket( - ctx, - r.logger, - storage.NopReadBucketCloser(readWriteBucket), - gitRef.SubDirPath(), - targetPaths, - targetExcludePaths, - terminateFunc, - ) + return readWriteBucket, nil } func (r *reader) getModuleKey( @@ -420,6 +471,36 @@ func (r *reader) getFileReadCloserAndSize( container app.EnvStdinContainer, fileRef FileRef, keepFileCompression bool, +) (io.ReadCloser, int64, error) { + if !isRemoteFileScheme(fileRef.FileScheme()) { + return r.getFileReadCloserAndSizeUncached(ctx, container, fileRef, keepFileCompression) + } + data, err := r.fileDataCache.GetOrAdd( + fileDataCacheKey{ + path: fileRef.Path(), + fileScheme: fileRef.FileScheme(), + compressionType: fileRef.CompressionType(), + keepFileCompression: keepFileCompression, + }, + func() ([]byte, error) { + readCloser, _, err := r.getFileReadCloserAndSizeUncached(ctx, container, fileRef, keepFileCompression) + if err != nil { + return nil, err + } + return xio.ReadAllAndClose(readCloser) + }, + ) + if err != nil { + return nil, -1, err + } + return io.NopCloser(bytes.NewReader(data)), int64(len(data)), nil +} + +func (r *reader) getFileReadCloserAndSizeUncached( + ctx context.Context, + container app.EnvStdinContainer, + fileRef FileRef, + keepFileCompression bool, ) (_ io.ReadCloser, _ int64, retErr error) { readCloser, size, err := r.getFileReadCloserAndSizePotentiallyCompressed(ctx, container, fileRef) if err != nil { diff --git a/private/buf/buffetch/internal/reader_cache.go b/private/buf/buffetch/internal/reader_cache.go new file mode 100644 index 0000000000..b83b97b378 --- /dev/null +++ b/private/buf/buffetch/internal/reader_cache.go @@ -0,0 +1,90 @@ +// Copyright 2020-2026 Buf Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package internal + +import "github.com/bufbuild/buf/private/pkg/git" + +// A cache key must contain everything that affects what is fetched, and nothing +// that is applied to the result afterwards, such as SubDirPath and target +// paths, so that Refs differing only in those still share a fetch. + +// fileDataCacheKey identifies the contents of a fetched file. +type fileDataCacheKey struct { + path string + fileScheme FileScheme + compressionType CompressionType + keepFileCompression bool +} + +// archiveBucketCacheKey identifies the contents of an unarchived archive. The +// unarchived bucket is cached rather than the archive itself so that the bytes +// are held once, and so that the archive is expanded once. +type archiveBucketCacheKey struct { + path string + fileScheme FileScheme + archiveType ArchiveType + compressionType CompressionType + stripComponents uint32 +} + +func newArchiveBucketCacheKey(archiveRef ArchiveRef) archiveBucketCacheKey { + return archiveBucketCacheKey{ + path: archiveRef.Path(), + fileScheme: archiveRef.FileScheme(), + archiveType: archiveRef.ArchiveType(), + compressionType: archiveRef.CompressionType(), + stripComponents: archiveRef.StripComponents(), + } +} + +// gitBucketCacheKey identifies the contents of a cloned git repository. +type gitBucketCacheKey struct { + path string + gitScheme GitScheme + cloneBranch string + checkout string + depth uint32 + recurseSubmodules bool + filter string + // Set only when filter is set, the one case where it changes what is cloned: + // the clone is then a sparse checkout of subDirPath. + subDirPath string +} + +func newGitBucketCacheKey(gitRef GitRef) gitBucketCacheKey { + cloneBranch, checkout := git.FetchIdentity(gitRef.GitName()) + filter := gitRef.Filter() + var subDirPath string + if filter != "" { + subDirPath = gitRef.SubDirPath() + } + return gitBucketCacheKey{ + path: gitRef.Path(), + gitScheme: gitRef.GitScheme(), + cloneBranch: cloneBranch, + checkout: checkout, + depth: gitRef.Depth(), + recurseSubmodules: gitRef.RecurseSubmodules(), + filter: filter, + subDirPath: subDirPath, + } +} + +// isRemoteFileScheme returns whether a FileRef with the given FileScheme is +// fetched over the network. Local files have no fetch to save, and streams +// cannot be re-read at all, so neither is cached. +func isRemoteFileScheme(fileScheme FileScheme) bool { + return fileScheme == FileSchemeHTTP || fileScheme == FileSchemeHTTPS +} diff --git a/private/buf/buffetch/internal/reader_cache_test.go b/private/buf/buffetch/internal/reader_cache_test.go new file mode 100644 index 0000000000..c55081b2a5 --- /dev/null +++ b/private/buf/buffetch/internal/reader_cache_test.go @@ -0,0 +1,309 @@ +// Copyright 2020-2026 Buf Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package internal + +import ( + "bytes" + "context" + "io" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "sync/atomic" + "testing" + + "buf.build/go/app" + "github.com/bufbuild/buf/private/pkg/git" + "github.com/bufbuild/buf/private/pkg/httpauth" + "github.com/bufbuild/buf/private/pkg/slogtestext" + "github.com/bufbuild/buf/private/pkg/storage" + "github.com/bufbuild/buf/private/pkg/storage/storagearchive" + "github.com/bufbuild/buf/private/pkg/storage/storagemem" + "github.com/bufbuild/buf/private/pkg/storage/storageos" + "github.com/klauspost/compress/gzip" + "github.com/stretchr/testify/require" +) + +func TestReaderArchiveFetchDeduplication(t *testing.T) { + t.Parallel() + ctx := t.Context() + server, requestCount := newTestArchiveServer(t) + reader := newTestHTTPReader(t) + // Reads differing only in what is applied after the fetch share one fetch. + // StripComponents is applied while unarchiving, so it is not one of those: + // the last read fetches again. + for _, read := range []archiveTestRead{ + {subDirPath: "svc-a"}, + {subDirPath: "svc-b"}, + {subDirPath: "svc-a"}, + {stripComponents: 1}, + } { + archiveRef, err := newArchiveRef( + "targz", + server.URL+"/archive.tar.gz", + ArchiveTypeTar, + CompressionTypeGzip, + read.stripComponents, + read.subDirPath, + ) + require.NoError(t, err) + readBucketCloser, bucketTargeting, err := reader.GetReadBucketCloser( + ctx, + newTestStdinContainer(), + archiveRef, + ) + require.NoError(t, err) + if read.subDirPath != "" { + // Each read is still scoped to its own subdirectory. + require.Equal(t, read.subDirPath, bucketTargeting.SubDirPath()) + data, err := storage.ReadPath(ctx, readBucketCloser, read.subDirPath+"/test.proto") + require.NoError(t, err) + require.Equal(t, read.subDirPath, string(data)) + } + require.NoError(t, readBucketCloser.Close()) + } + require.Equal(t, int64(2), requestCount.Load()) +} + +func TestReaderFileFetchDeduplication(t *testing.T) { + t.Parallel() + ctx := t.Context() + var requestCount atomic.Int64 + server := httptest.NewServer( + http.HandlerFunc(func(responseWriter http.ResponseWriter, request *http.Request) { + requestCount.Add(1) + _, err := responseWriter.Write([]byte("image")) + require.NoError(t, err) + }), + ) + t.Cleanup(server.Close) + reader := newTestHTTPReader(t) + for range 3 { + singleRef, err := newSingleRef("binpb", server.URL+"/image.binpb", CompressionTypeNone, nil) + require.NoError(t, err) + readCloser, err := reader.GetFile(ctx, newTestStdinContainer(), singleRef) + require.NoError(t, err) + data, err := io.ReadAll(readCloser) + require.NoError(t, err) + require.NoError(t, readCloser.Close()) + require.Equal(t, "image", string(data)) + } + require.Equal(t, int64(1), requestCount.Load()) +} + +func TestReaderFileFetchNoDeduplicationForStdin(t *testing.T) { + t.Parallel() + ctx := t.Context() + reader := NewReader(slogtestext.NewLogger(t), storageos.NewProvider(), WithReaderStdio()) + container := app.NewContainer(nil, strings.NewReader("image"), nil, nil) + singleRef, err := newSingleRef("binpb", "-", CompressionTypeNone, nil) + require.NoError(t, err) + readCloser, err := reader.GetFile(ctx, container, singleRef) + require.NoError(t, err) + data, err := io.ReadAll(readCloser) + require.NoError(t, err) + require.NoError(t, readCloser.Close()) + require.Equal(t, "image", string(data)) + // Stdin is a stream, so it is not cached: the second read is exhausted. + readCloser, err = reader.GetFile(ctx, container, singleRef) + require.NoError(t, err) + data, err = io.ReadAll(readCloser) + require.NoError(t, err) + require.NoError(t, readCloser.Close()) + require.Empty(t, string(data)) +} + +func TestReaderGitCloneDeduplication(t *testing.T) { + t.Parallel() + for _, testCase := range []struct { + name string + reads []gitTestRead + expectedClones int + }{ + { + name: "same repository different subdirs is one clone", + reads: []gitTestRead{ + {gitName: git.NewBranchName("main"), subDirPath: "svc-a"}, + {gitName: git.NewBranchName("main"), subDirPath: "svc-b"}, + }, + expectedClones: 1, + }, + { + name: "no name is one clone", + reads: []gitTestRead{ + {subDirPath: "svc-a"}, + {subDirPath: "svc-b"}, + }, + expectedClones: 1, + }, + { + // These share a String, so the key holds the fetch identity. + name: "ref and branch with the same value are separate clones", + reads: []gitTestRead{ + {gitName: git.NewBranchName("main")}, + {gitName: git.NewRefName("main")}, + {gitName: git.NewRefNameWithBranch("main", "dev")}, + }, + expectedClones: 3, + }, + { + name: "different depth is separate clones", + reads: []gitTestRead{ + {gitName: git.NewBranchName("main"), depth: 1}, + {gitName: git.NewBranchName("main"), depth: 50}, + }, + expectedClones: 2, + }, + { + name: "different recurse submodules is separate clones", + reads: []gitTestRead{ + {gitName: git.NewBranchName("main")}, + {gitName: git.NewBranchName("main"), recurseSubmodules: true}, + }, + expectedClones: 2, + }, + { + // With a filter the subdirectory is a sparse checkout. + name: "different subdirs with a filter are separate clones", + reads: []gitTestRead{ + {gitName: git.NewBranchName("main"), subDirPath: "svc-a", filter: "blob:none"}, + {gitName: git.NewBranchName("main"), subDirPath: "svc-b", filter: "blob:none"}, + }, + expectedClones: 2, + }, + } { + t.Run(testCase.name, func(t *testing.T) { + t.Parallel() + ctx := t.Context() + cloner := &testCloner{} + reader := NewReader( + slogtestext.NewLogger(t), + storageos.NewProvider(), + WithReaderGit(cloner), + ) + for _, read := range testCase.reads { + depth := read.depth + if depth == 0 { + depth = 1 + } + gitRef, err := newGitRef( + "git", + "https://github.com/foo/bar.git", + read.gitName, + depth, + read.recurseSubmodules, + read.subDirPath, + read.filter, + ) + require.NoError(t, err) + readBucketCloser, _, err := reader.GetReadBucketCloser( + ctx, + newTestStdinContainer(), + gitRef, + ) + require.NoError(t, err) + require.NoError(t, readBucketCloser.Close()) + } + require.Equal(t, testCase.expectedClones, cloner.cloneCount) + }) + } +} + +type archiveTestRead struct { + subDirPath string + stripComponents uint32 +} + +type gitTestRead struct { + gitName git.Name + depth uint32 + recurseSubmodules bool + subDirPath string + filter string +} + +// testCloner writes fixed contents on every clone and counts calls. +type testCloner struct { + cloneCount int +} + +func (c *testCloner) CloneToBucket( + ctx context.Context, + _ app.EnvContainer, + _ string, + _ uint32, + writeBucket storage.WriteBucket, + _ git.CloneToBucketOptions, +) error { + c.cloneCount++ + return putTestArchiveFiles(ctx, writeBucket) +} + +// newTestArchiveServer serves a gzipped tarball of putTestArchiveFiles, and +// counts requests. +func newTestArchiveServer(t *testing.T) (*httptest.Server, *atomic.Int64) { + t.Helper() + ctx := t.Context() + readWriteBucket := storagemem.NewReadWriteBucket() + require.NoError(t, putTestArchiveFiles(ctx, readWriteBucket)) + buffer := bytes.NewBuffer(nil) + gzipWriter := gzip.NewWriter(buffer) + require.NoError(t, storagearchive.Tar(ctx, readWriteBucket, gzipWriter)) + require.NoError(t, gzipWriter.Close()) + archiveData := buffer.Bytes() + var requestCount atomic.Int64 + server := httptest.NewServer( + http.HandlerFunc(func(responseWriter http.ResponseWriter, request *http.Request) { + requestCount.Add(1) + _, err := responseWriter.Write(archiveData) + require.NoError(t, err) + }), + ) + t.Cleanup(server.Close) + return server, &requestCount +} + +// putTestArchiveFiles writes a workspace with one module per subdirectory. Each +// file's contents are the name of the subdirectory containing it. +func putTestArchiveFiles(ctx context.Context, writeBucket storage.WriteBucket) error { + if err := storage.PutPath(ctx, writeBucket, "buf.yaml", []byte("version: v2\n")); err != nil { + return err + } + for _, subDirPath := range []string{"svc-a", "svc-b"} { + if err := storage.PutPath( + ctx, + writeBucket, + filepath.ToSlash(filepath.Join(subDirPath, "test.proto")), + []byte(subDirPath), + ); err != nil { + return err + } + } + return nil +} + +func newTestHTTPReader(t *testing.T) Reader { + t.Helper() + return NewReader( + slogtestext.NewLogger(t), + storageos.NewProvider(), + WithReaderHTTP(http.DefaultClient, httpauth.NewNopAuthenticator()), + ) +} + +func newTestStdinContainer() app.EnvStdinContainer { + return app.NewContainer(nil, nil, nil, nil) +} diff --git a/private/pkg/git/git.go b/private/pkg/git/git.go index 3e9807adad..32d08275da 100644 --- a/private/pkg/git/git.go +++ b/private/pkg/git/git.go @@ -63,6 +63,19 @@ type Name interface { checkout() string } +// FetchIdentity returns the values that determine what a clone of name fetches: +// the branch to clone, if any, and the ref to check out afterwards, if any. A +// nil name has an empty identity. +// +// Names with equal identities produce equal clones. String is not sufficient +// for this: a branch, a tag, and a ref with the same value all share a String. +func FetchIdentity(name Name) (cloneBranch string, checkout string) { + if name == nil { + return "", "" + } + return name.cloneBranch(), name.checkout() +} + // NewBranchName returns a new Name for the branch. func NewBranchName(branch string) Name { return newBranch(branch) From 3cdae6e25982925d41392802950336b4f255f60e Mon Sep 17 00:00:00 2001 From: Edward McFarlane Date: Wed, 16 Sep 2026 16:50:22 +0100 Subject: [PATCH 2/2] Configure cache settings --- cmd/buf/internal/command/generate/generate.go | 2 + private/buf/bufctl/controller.go | 6 +++ private/buf/bufctl/option.go | 10 ++++ private/buf/buffetch/buffetch.go | 20 ++++++-- private/buf/buffetch/internal/internal.go | 7 +++ private/buf/buffetch/internal/reader.go | 31 +++++++++---- .../buffetch/internal/reader_cache_test.go | 46 +++++++++++++++++-- private/buf/buffetch/reader.go | 39 +++++++++++----- 8 files changed, 132 insertions(+), 29 deletions(-) diff --git a/cmd/buf/internal/command/generate/generate.go b/cmd/buf/internal/command/generate/generate.go index 706f4bde26..7b7dd31681 100644 --- a/cmd/buf/internal/command/generate/generate.go +++ b/cmd/buf/internal/command/generate/generate.go @@ -534,6 +534,8 @@ func run( container, bufctl.WithDisableSymlinks(flags.DisableSymlinks), bufctl.WithFileAnnotationErrorFormat(flags.ErrorFormat), + // Multiple inputs may resolve to the same remote, fetch each one once. + bufctl.WithReaderFetchCache(), ) if err != nil { return err diff --git a/private/buf/bufctl/controller.go b/private/buf/bufctl/controller.go index d37e5498d4..232d431131 100644 --- a/private/buf/bufctl/controller.go +++ b/private/buf/bufctl/controller.go @@ -217,6 +217,7 @@ type controller struct { fileAnnotationErrorFormat string fileAnnotationsToStdout bool copyToInMemory bool + readerFetchCacheEnabled bool storageosProvider storageos.Provider buffetchRefParser buffetch.RefParser @@ -263,6 +264,10 @@ func newController( } controller.storageosProvider = newStorageosProvider(controller.disableSymlinks) controller.buffetchRefParser = buffetch.NewRefParser(logger) + var buffetchReaderOptions []buffetch.ReaderOption + if controller.readerFetchCacheEnabled { + buffetchReaderOptions = append(buffetchReaderOptions, buffetch.WithReaderFetchCache()) + } controller.buffetchReader = buffetch.NewReader( logger, controller.storageosProvider, @@ -274,6 +279,7 @@ func newController( gitClonerOptions, ), moduleKeyProvider, + buffetchReaderOptions..., ) controller.buffetchWriter = buffetch.NewWriter(logger) controller.workspaceProvider = bufworkspace.NewWorkspaceProvider( diff --git a/private/buf/bufctl/option.go b/private/buf/bufctl/option.go index 74760444df..250682a1bc 100644 --- a/private/buf/bufctl/option.go +++ b/private/buf/bufctl/option.go @@ -42,6 +42,16 @@ func WithFileAnnotationsToStdout() ControllerOption { } } +// WithReaderFetchCache returns a new ControllerOption that fetches a given +// remote input at most once for the lifetime of the Controller. +// +// See buffetch.WithReaderFetchCache. +func WithReaderFetchCache() ControllerOption { + return func(controller *controller) { + controller.readerFetchCacheEnabled = true + } +} + // WithCopyToInMemory returns a new ControllerOption that copies to memory. func WithCopyToInMemory() ControllerOption { return func(controller *controller) { diff --git a/private/buf/buffetch/buffetch.go b/private/buf/buffetch/buffetch.go index d8cdcf6334..bef22ed94b 100644 --- a/private/buf/buffetch/buffetch.go +++ b/private/buf/buffetch/buffetch.go @@ -414,9 +414,6 @@ type ModuleFetcher interface { } // Reader is a reader for Buf. -// -// A Reader fetches a given remote file or git repository at most once, for the -// lifetime of the Reader. type Reader interface { MessageReader SourceReader @@ -424,6 +421,21 @@ type Reader interface { ModuleFetcher } +// ReaderOption is a Reader option. +type ReaderOption func(*readerOptions) + +// WithReaderFetchCache returns a ReaderOption that fetches a given remote file +// or git repository at most once, so that Refs resolving to the same remote +// share a fetch. +// +// Fetches are held in memory for the lifetime of the Reader, so only use this +// for a Reader that does not outlive the Refs it is created for. +func WithReaderFetchCache() ReaderOption { + return func(readerOptions *readerOptions) { + readerOptions.fetchCacheEnabled = true + } +} + // NewReader returns a new Reader. func NewReader( logger *slog.Logger, @@ -432,6 +444,7 @@ func NewReader( httpAuthenticator httpauth.Authenticator, gitCloner git.Cloner, moduleKeyProvider bufmodule.ModuleKeyProvider, + options ...ReaderOption, ) Reader { return newReader( logger, @@ -440,6 +453,7 @@ func NewReader( httpAuthenticator, gitCloner, moduleKeyProvider, + options..., ) } diff --git a/private/buf/buffetch/internal/internal.go b/private/buf/buffetch/internal/internal.go index 4c75233b48..c8e5c3e920 100644 --- a/private/buf/buffetch/internal/internal.go +++ b/private/buf/buffetch/internal/internal.go @@ -778,6 +778,13 @@ func WithReaderStdio() ReaderOption { } } +// WithReaderFetchCache enables caching of remote fetches. +func WithReaderFetchCache() ReaderOption { + return func(reader *reader) { + reader.fetchCacheEnabled = true + } +} + // WriterOption is an Writer option. type WriterOption func(*writer) diff --git a/private/buf/buffetch/internal/reader.go b/private/buf/buffetch/internal/reader.go index 2020696317..d778920891 100644 --- a/private/buf/buffetch/internal/reader.go +++ b/private/buf/buffetch/internal/reader.go @@ -63,7 +63,9 @@ type reader struct { moduleEnabled bool moduleKeyProvider bufmodule.ModuleKeyProvider - // Keyed so that Refs sharing a fetch share an entry, see reader_cache.go. + // Caches are keyed so that Refs sharing a fetch share an entry, see + // reader_cache.go. Entries are held for the lifetime of the reader. + fetchCacheEnabled bool fileDataCache cache.Cache[fileDataCacheKey, []byte] archiveBucketCache cache.Cache[archiveBucketCacheKey, storage.ReadBucket] gitBucketCache cache.Cache[gitBucketCacheKey, storage.ReadBucket] @@ -271,7 +273,7 @@ func (r *reader) getArchiveReadBucket( container app.EnvStdinContainer, archiveRef ArchiveRef, ) (storage.ReadBucket, error) { - if !isRemoteFileScheme(archiveRef.FileScheme()) { + if !r.fetchCacheEnabled || !isRemoteFileScheme(archiveRef.FileScheme()) { return r.unarchive(ctx, container, archiveRef) } return r.archiveBucketCache.GetOrAdd( @@ -393,12 +395,7 @@ func (r *reader) getGitBucket( if r.gitCloner == nil { return nil, nil, errors.New("git cloner is nil") } - readBucket, err := r.gitBucketCache.GetOrAdd( - newGitBucketCacheKey(gitRef), - func() (storage.ReadBucket, error) { - return r.clone(ctx, container, gitRef) - }, - ) + readBucket, err := r.getGitReadBucket(ctx, container, gitRef) if err != nil { return nil, nil, err } @@ -413,6 +410,22 @@ func (r *reader) getGitBucket( ) } +func (r *reader) getGitReadBucket( + ctx context.Context, + container app.EnvStdinContainer, + gitRef GitRef, +) (storage.ReadBucket, error) { + if !r.fetchCacheEnabled { + return r.clone(ctx, container, gitRef) + } + return r.gitBucketCache.GetOrAdd( + newGitBucketCacheKey(gitRef), + func() (storage.ReadBucket, error) { + return r.clone(ctx, container, gitRef) + }, + ) +} + func (r *reader) clone( ctx context.Context, container app.EnvStdinContainer, @@ -472,7 +485,7 @@ func (r *reader) getFileReadCloserAndSize( fileRef FileRef, keepFileCompression bool, ) (io.ReadCloser, int64, error) { - if !isRemoteFileScheme(fileRef.FileScheme()) { + if !r.fetchCacheEnabled || !isRemoteFileScheme(fileRef.FileScheme()) { return r.getFileReadCloserAndSizeUncached(ctx, container, fileRef, keepFileCompression) } data, err := r.fileDataCache.GetOrAdd( diff --git a/private/buf/buffetch/internal/reader_cache_test.go b/private/buf/buffetch/internal/reader_cache_test.go index c55081b2a5..d76bd1aedd 100644 --- a/private/buf/buffetch/internal/reader_cache_test.go +++ b/private/buf/buffetch/internal/reader_cache_test.go @@ -41,7 +41,7 @@ func TestReaderArchiveFetchDeduplication(t *testing.T) { t.Parallel() ctx := t.Context() server, requestCount := newTestArchiveServer(t) - reader := newTestHTTPReader(t) + reader := newTestHTTPReader(t, WithReaderFetchCache()) // Reads differing only in what is applied after the fetch share one fetch. // StripComponents is applied while unarchiving, so it is not one of those: // the last read fetches again. @@ -78,6 +78,33 @@ func TestReaderArchiveFetchDeduplication(t *testing.T) { require.Equal(t, int64(2), requestCount.Load()) } +func TestReaderArchiveFetchWithoutCache(t *testing.T) { + t.Parallel() + ctx := t.Context() + server, requestCount := newTestArchiveServer(t) + reader := newTestHTTPReader(t) + for range 2 { + archiveRef, err := newArchiveRef( + "targz", + server.URL+"/archive.tar.gz", + ArchiveTypeTar, + CompressionTypeGzip, + 0, + "svc-a", + ) + require.NoError(t, err) + readBucketCloser, _, err := reader.GetReadBucketCloser( + ctx, + newTestStdinContainer(), + archiveRef, + ) + require.NoError(t, err) + require.NoError(t, readBucketCloser.Close()) + } + // Without WithReaderFetchCache, every read fetches. + require.Equal(t, int64(2), requestCount.Load()) +} + func TestReaderFileFetchDeduplication(t *testing.T) { t.Parallel() ctx := t.Context() @@ -90,7 +117,7 @@ func TestReaderFileFetchDeduplication(t *testing.T) { }), ) t.Cleanup(server.Close) - reader := newTestHTTPReader(t) + reader := newTestHTTPReader(t, WithReaderFetchCache()) for range 3 { singleRef, err := newSingleRef("binpb", server.URL+"/image.binpb", CompressionTypeNone, nil) require.NoError(t, err) @@ -107,7 +134,12 @@ func TestReaderFileFetchDeduplication(t *testing.T) { func TestReaderFileFetchNoDeduplicationForStdin(t *testing.T) { t.Parallel() ctx := t.Context() - reader := NewReader(slogtestext.NewLogger(t), storageos.NewProvider(), WithReaderStdio()) + reader := NewReader( + slogtestext.NewLogger(t), + storageos.NewProvider(), + WithReaderStdio(), + WithReaderFetchCache(), + ) container := app.NewContainer(nil, strings.NewReader("image"), nil, nil) singleRef, err := newSingleRef("binpb", "-", CompressionTypeNone, nil) require.NoError(t, err) @@ -193,6 +225,7 @@ func TestReaderGitCloneDeduplication(t *testing.T) { slogtestext.NewLogger(t), storageos.NewProvider(), WithReaderGit(cloner), + WithReaderFetchCache(), ) for _, read := range testCase.reads { depth := read.depth @@ -295,12 +328,15 @@ func putTestArchiveFiles(ctx context.Context, writeBucket storage.WriteBucket) e return nil } -func newTestHTTPReader(t *testing.T) Reader { +func newTestHTTPReader(t *testing.T, options ...ReaderOption) Reader { t.Helper() return NewReader( slogtestext.NewLogger(t), storageos.NewProvider(), - WithReaderHTTP(http.DefaultClient, httpauth.NewNopAuthenticator()), + append( + []ReaderOption{WithReaderHTTP(http.DefaultClient, httpauth.NewNopAuthenticator())}, + options..., + )..., ) } diff --git a/private/buf/buffetch/reader.go b/private/buf/buffetch/reader.go index 195ede9196..3ae3a21d85 100644 --- a/private/buf/buffetch/reader.go +++ b/private/buf/buffetch/reader.go @@ -33,6 +33,10 @@ type reader struct { internalReader internal.Reader } +type readerOptions struct { + fetchCacheEnabled bool +} + func newReader( logger *slog.Logger, storageosProvider storageos.Provider, @@ -40,23 +44,34 @@ func newReader( httpAuthenticator httpauth.Authenticator, gitCloner git.Cloner, moduleKeyProvider bufmodule.ModuleKeyProvider, + options ...ReaderOption, ) *reader { + readerOptions := &readerOptions{} + for _, option := range options { + option(readerOptions) + } + internalReaderOptions := []internal.ReaderOption{ + internal.WithReaderHTTP( + httpClient, + httpAuthenticator, + ), + internal.WithReaderGit( + gitCloner, + ), + internal.WithReaderLocal(), + internal.WithReaderStdio(), + internal.WithReaderModule( + moduleKeyProvider, + ), + } + if readerOptions.fetchCacheEnabled { + internalReaderOptions = append(internalReaderOptions, internal.WithReaderFetchCache()) + } return &reader{ internalReader: internal.NewReader( logger, storageosProvider, - internal.WithReaderHTTP( - httpClient, - httpAuthenticator, - ), - internal.WithReaderGit( - gitCloner, - ), - internal.WithReaderLocal(), - internal.WithReaderStdio(), - internal.WithReaderModule( - moduleKeyProvider, - ), + internalReaderOptions..., ), } }