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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,9 @@
only used to report parse errors and diffs.
- Fix lint comment ignores on proto2 `group` fields being ignored.
- Improve the `buf curl` error message for methods that accept a single request message.
- 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

Expand Down
2 changes: 2 additions & 0 deletions cmd/buf/internal/command/generate/generate.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions private/buf/bufctl/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -217,6 +217,7 @@ type controller struct {
fileAnnotationErrorFormat string
fileAnnotationsToStdout bool
copyToInMemory bool
readerFetchCacheEnabled bool

storageosProvider storageos.Provider
buffetchRefParser buffetch.RefParser
Expand Down Expand Up @@ -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,
Expand All @@ -274,6 +279,7 @@ func newController(
gitClonerOptions,
),
moduleKeyProvider,
buffetchReaderOptions...,
)
controller.buffetchWriter = buffetch.NewWriter(logger)
controller.workspaceProvider = bufworkspace.NewWorkspaceProvider(
Expand Down
10 changes: 10 additions & 0 deletions private/buf/bufctl/option.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
17 changes: 17 additions & 0 deletions private/buf/buffetch/buffetch.go
Original file line number Diff line number Diff line change
Expand Up @@ -421,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,
Expand All @@ -429,6 +444,7 @@ func NewReader(
httpAuthenticator httpauth.Authenticator,
gitCloner git.Cloner,
moduleKeyProvider bufmodule.ModuleKeyProvider,
options ...ReaderOption,
) Reader {
return newReader(
logger,
Expand All @@ -437,6 +453,7 @@ func NewReader(
httpAuthenticator,
gitCloner,
moduleKeyProvider,
options...,
)
}

Expand Down
7 changes: 7 additions & 0 deletions private/buf/buffetch/internal/internal.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
148 changes: 121 additions & 27 deletions private/buf/buffetch/internal/reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -61,6 +62,13 @@ type reader struct {

moduleEnabled bool
moduleKeyProvider bufmodule.ModuleKeyProvider

// 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]
}

func newReader(
Expand Down Expand Up @@ -244,11 +252,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 !r.fetchCacheEnabled || !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())
}()
Expand All @@ -263,21 +307,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(
Expand All @@ -289,20 +333,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(
Expand Down Expand Up @@ -359,10 +395,46 @@ func (r *reader) getGitBucket(
if r.gitCloner == nil {
return nil, nil, errors.New("git cloner is nil")
}
gitURL, err := getGitURL(gitRef)
readBucket, err := r.getGitReadBucket(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) 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,
gitRef GitRef,
) (storage.ReadBucket, error) {
gitURL, err := getGitURL(gitRef)
if err != nil {
return nil, err
}
readWriteBucket := storagemem.NewReadWriteBucket()
if err := r.gitCloner.CloneToBucket(
ctx,
Expand All @@ -377,17 +449,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(
Expand Down Expand Up @@ -420,6 +484,36 @@ func (r *reader) getFileReadCloserAndSize(
container app.EnvStdinContainer,
fileRef FileRef,
keepFileCompression bool,
) (io.ReadCloser, int64, error) {
if !r.fetchCacheEnabled || !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 {
Expand Down
Loading
Loading