diff --git a/client/client_cache_test.go b/client/client_cache_test.go index 7ec34307b65f..cda1872f9a4c 100644 --- a/client/client_cache_test.go +++ b/client/client_cache_test.go @@ -1,6 +1,7 @@ package client import ( + "context" "encoding/base64" "encoding/json" "fmt" @@ -15,13 +16,17 @@ import ( "time" ctd "github.com/containerd/containerd/v2/client" + "github.com/containerd/containerd/v2/core/content" "github.com/containerd/containerd/v2/core/images" "github.com/containerd/containerd/v2/pkg/namespaces" + contentlocal "github.com/containerd/containerd/v2/plugins/content/local" cerrdefs "github.com/containerd/errdefs" cacheimporttypes "github.com/moby/buildkit/cache/remotecache/v1/types" "github.com/moby/buildkit/client/llb" "github.com/moby/buildkit/exporter/containerimage/exptypes" "github.com/moby/buildkit/identity" + "github.com/moby/buildkit/session" + sessioncontent "github.com/moby/buildkit/session/content" "github.com/moby/buildkit/util/testutil" "github.com/moby/buildkit/util/testutil/helpers" "github.com/moby/buildkit/util/testutil/integration" @@ -30,8 +35,145 @@ import ( "github.com/pkg/errors" "github.com/stretchr/testify/require" "github.com/tonistiigi/fsutil" + "golang.org/x/sync/errgroup" ) +func testConcurrentLocalCacheExport(t *testing.T, sb integration.Sandbox) { + integration.SkipOnPlatform(t, "windows") + workers.CheckFeatureCompat(t, sb, workers.FeatureCacheExport, workers.FeatureCacheBackendLocal) + + ctx, cancel := context.WithTimeoutCause(sb.Context(), time.Minute, errors.New("concurrent local cache export timed out")) + defer cancel() + + c, err := New(ctx, sb.Address()) + require.NoError(t, err) + defer c.Close() + + st := llb.Image("busybox:latest").Run(llb.Shlex("sh -c 'echo shared > /out/payload'")).AddMount("/out", llb.Scratch()) + def, err := st.Marshal(ctx) + require.NoError(t, err) + + dir := t.TempDir() + opened := make(chan struct{}, 2) + blocked := make(chan struct{}, 1) + releases := [2]chan struct{}{make(chan struct{}), make(chan struct{})} + completed := make(chan error, 2) + eg, ctx := errgroup.WithContext(ctx) + for i := range 2 { + store, err := newLocalCacheStore(dir) + require.NoError(t, err) + s, err := session.NewSession(ctx, "") + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, s.Close()) }) + s.Allow(sessioncontent.NewAttachable(map[string]content.Store{ + "local:" + dir: &gatedCacheStore{Store: store, opened: opened, blocked: blocked, release: releases[i]}, + })) + eg.Go(func() error { return s.Run(ctx, c.Dialer()) }) + eg.Go(func() error { + defer s.Close() + _, err := c.Solve(ctx, def, SolveOpt{ + SharedSession: s, + SessionPreInitialized: true, + CacheExports: []CacheOptionsEntry{{Type: "local", Attrs: map[string]string{ + "dest": dir, + "tag": fmt.Sprintf("export-%d", i), + }}}, + }, nil) + completed <- err + return err + }) + if i == 0 { + select { + case <-opened: + case <-ctx.Done(): + t.Fatal(context.Cause(ctx)) + } + } + } + done := make(chan error, 1) + go func() { done <- eg.Wait() }() + + // The first export owns the ingest. The second must retry rather than + // opening and modifying the same temporary files. + select { + case <-blocked: + case <-opened: + t.Fatal("second export opened an active ingest") + case <-ctx.Done(): + t.Fatal(context.Cause(ctx)) + } + + entries, err := os.ReadDir(filepath.Join(dir, "ingest")) + require.NoError(t, err) + require.Len(t, entries, 1, "exports of the same reference must reuse its ingest") + // Index updates use a non-blocking lock. Delay the second export's retry + // until the first has published its index entry. + for _, release := range releases { + close(release) + require.NoError(t, <-completed) + } + require.NoError(t, <-done) + + dt, err := os.ReadFile(filepath.Join(dir, "index.json")) + require.NoError(t, err) + var index ocispecs.Index + require.NoError(t, json.Unmarshal(dt, &index)) + require.Len(t, index.Manifests, 2) + + store, err := contentlocal.NewStore(dir) + require.NoError(t, err) + + for _, desc := range index.Manifests { + dt, err := content.ReadBlob(sb.Context(), store, desc) + require.NoError(t, err) + var manifest ocispecs.Manifest + require.NoError(t, json.Unmarshal(dt, &manifest)) + for _, blob := range append(manifest.Layers, manifest.Config) { + dt, err := content.ReadBlob(sb.Context(), store, blob) + require.NoError(t, err) + require.Equal(t, blob.Digest, blob.Digest.Algorithm().FromBytes(dt)) + } + } +} + +// gatedCacheStore pauses the first writer so two exports can overlap without +// relying on scheduling or a large payload to keep an ingest active. +type gatedCacheStore struct { + content.Store + opened chan<- struct{} + blocked chan<- struct{} + release <-chan struct{} + paused atomic.Bool +} + +func (s *gatedCacheStore) Writer(ctx context.Context, opts ...content.WriterOpt) (content.Writer, error) { + w, err := s.Store.Writer(ctx, opts...) + if err != nil { + if errors.Is(err, cerrdefs.ErrUnavailable) { + select { + case s.blocked <- struct{}{}: + default: + } + select { + case <-s.release: + case <-ctx.Done(): + return nil, context.Cause(ctx) + } + } + return nil, err + } + if s.paused.CompareAndSwap(false, true) { + s.opened <- struct{}{} + select { + case <-s.release: + case <-ctx.Done(): + w.Close() + return nil, context.Cause(ctx) + } + } + return w, nil +} + func testBasicAzblobCacheImportExport(t *testing.T, sb integration.Sandbox) { integration.SkipOnPlatform(t, "windows") workers.CheckFeatureCompat(t, sb, diff --git a/client/client_test.go b/client/client_test.go index 7407d02541d3..1d6644705d94 100644 --- a/client/client_test.go +++ b/client/client_test.go @@ -21,6 +21,7 @@ var allTests = []func(t *testing.T, sb integration.Sandbox){ testBasicAzblobCacheImportExport, testBasicInlineCacheImportExport, testBasicLocalCacheImportExport, + testConcurrentLocalCacheExport, testBasicRegistryCacheImportExport, testBasicS3CacheImportExport, testCacheExportCacheDeletedContent, diff --git a/client/solve.go b/client/solve.go index 12d68b0c452c..b8e4e355872f 100644 --- a/client/solve.go +++ b/client/solve.go @@ -7,6 +7,7 @@ import ( "io" "maps" "os" + "path/filepath" "slices" "strconv" "strings" @@ -16,6 +17,8 @@ import ( "github.com/containerd/containerd/v2/core/content" "github.com/containerd/containerd/v2/core/images" contentlocal "github.com/containerd/containerd/v2/plugins/content/local" + cerrdefs "github.com/containerd/errdefs" + "github.com/gofrs/flock" controlapi "github.com/moby/buildkit/api/services/control" "github.com/moby/buildkit/client/llb" "github.com/moby/buildkit/client/ociindex" @@ -561,7 +564,7 @@ func parseCacheOptions(ctx context.Context, isGateway bool, opt SolveOpt) (*cach if err := os.MkdirAll(csDir, 0755); err != nil { return nil, err } - cs, err := contentlocal.NewStore(csDir) + cs, err := newLocalCacheStore(csDir) if err != nil { return nil, err } @@ -601,7 +604,7 @@ func parseCacheOptions(ctx context.Context, isGateway bool, opt SolveOpt) (*cach if csDir == "" { return nil, errors.New("local cache importer requires src") } - cs, err := contentlocal.NewStore(csDir) + cs, err := newLocalCacheStore(csDir) if err != nil { bklog.G(ctx).Warning("local cache import at " + csDir + " not found due to err: " + err.Error()) continue @@ -667,3 +670,90 @@ func parseCacheOptions(ctx context.Context, isGateway bool, opt SolveOpt) (*cach } return &res, nil } + +// localCacheStore coordinates ingests across clients sharing a cache directory. +// The underlying local store only locks references within one store instance. +type localCacheStore struct { + content.Store + dir string +} + +func newLocalCacheStore(dir string) (content.Store, error) { + store, err := contentlocal.NewStore(dir) + if err != nil { + return nil, err + } + return &localCacheStore{Store: store, dir: dir}, nil +} + +func (s *localCacheStore) lock(ref string) (*flock.Flock, error) { + dir := filepath.Join(s.dir, "ingest-locks") + if err := os.MkdirAll(dir, 0755); err != nil { + return nil, err + } + // Keep lock files after unlocking: deleting one could let another client + // lock a different inode while an existing waiter still uses the old one. + lock := flock.New(filepath.Join(dir, digest.FromString(ref).Encoded())) + locked, err := lock.TryLock() + if err != nil { + return nil, err + } + if !locked { + return nil, errors.Wrapf(cerrdefs.ErrUnavailable, "ingest ref %q is locked", ref) + } + return lock, nil +} + +func (s *localCacheStore) Writer(ctx context.Context, opts ...content.WriterOpt) (content.Writer, error) { + var wo content.WriterOpts + for _, opt := range opts { + if err := opt(&wo); err != nil { + return nil, err + } + } + if wo.Ref == "" { + return nil, errors.Wrap(cerrdefs.ErrInvalidArgument, "ref must not be empty") + } + lock, err := s.lock(wo.Ref) + if err != nil { + return nil, err + } + w, err := s.Store.Writer(ctx, content.WithRef(wo.Ref), content.WithDescriptor(wo.Desc)) + if err != nil { + lock.Close() + return nil, err + } + return &localCacheWriter{Writer: w, lock: lock}, nil +} + +func (s *localCacheStore) Abort(ctx context.Context, ref string) error { + lock, err := s.lock(ref) + if err != nil { + return err + } + defer lock.Close() + return s.Store.Abort(ctx, ref) +} + +type localCacheWriter struct { + content.Writer + lock *flock.Flock +} + +func (w *localCacheWriter) Commit(ctx context.Context, size int64, expected digest.Digest, opts ...content.Opt) error { + err := w.Writer.Commit(ctx, size, expected, opts...) + closeErr := w.Close() + if err != nil { + return err + } + return closeErr +} + +func (w *localCacheWriter) Close() error { + err := w.Writer.Close() + lockErr := w.lock.Close() + if err != nil { + return err + } + return lockErr +} diff --git a/client/solve_test.go b/client/solve_test.go index fd5258134b9a..1de07cac9e46 100644 --- a/client/solve_test.go +++ b/client/solve_test.go @@ -1,12 +1,100 @@ package client import ( + "context" + "os" + "path/filepath" "testing" + "github.com/containerd/containerd/v2/core/content" + cerrdefs "github.com/containerd/errdefs" "github.com/moby/buildkit/client/llb" + digest "github.com/opencontainers/go-digest" + ocispecs "github.com/opencontainers/image-spec/specs-go/v1" + "github.com/pkg/errors" "github.com/stretchr/testify/require" ) +func TestLocalCacheStoreLockAndResume(t *testing.T) { + ctx := t.Context() + dir := t.TempDir() + first, err := newLocalCacheStore(dir) + require.NoError(t, err) + second, err := newLocalCacheStore(dir) + require.NoError(t, err) + data := []byte("resumable cache layer") + desc := ocispecs.Descriptor{Digest: digest.FromBytes(data), Size: int64(len(data))} + opts := []content.WriterOpt{content.WithRef("layer"), content.WithDescriptor(desc)} + w, err := first.Writer(ctx, opts...) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, w.Close()) }) + _, err = w.Write(data[:5]) + require.NoError(t, err) + + _, err = second.Writer(ctx, opts...) + require.ErrorIs(t, err, cerrdefs.ErrUnavailable) + require.ErrorIs(t, second.Abort(ctx, "layer"), cerrdefs.ErrUnavailable) + // Cancellation of OpenWriter's retry must not disturb the active writer. + cancelled, cancel := context.WithCancelCause(ctx) + cancel(errors.New("cancel retry")) + _, err = content.OpenWriter(cancelled, second, opts...) + require.ErrorIs(t, err, cerrdefs.ErrUnavailable) + + // Independent references can still be written concurrently. + other, err := second.Writer(ctx, content.WithRef("other")) + require.NoError(t, err) + require.NoError(t, other.Close()) + require.NoError(t, second.Abort(ctx, "other")) + require.NoError(t, w.Close()) + + resumed, err := second.Writer(ctx, opts...) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, resumed.Close()) }) + status, err := resumed.Status() + require.NoError(t, err) + require.EqualValues(t, 5, status.Offset) + _, err = resumed.Write(data[5:]) + require.NoError(t, err) + require.NoError(t, resumed.Commit(ctx, desc.Size, desc.Digest)) + blob, err := content.ReadBlob(ctx, first, desc) + require.NoError(t, err) + require.Equal(t, data, blob) + entries, err := os.ReadDir(filepath.Join(dir, "ingest")) + require.NoError(t, err) + require.Empty(t, entries, "resuming an interrupted export must reclaim its ingest") + + // An already-existing blob must not leave its reference locked. + _, err = first.Writer(ctx, opts...) + require.ErrorIs(t, err, cerrdefs.ErrAlreadyExists) + require.NoError(t, second.Abort(ctx, "layer")) +} + +func TestLocalCacheStoreCommitError(t *testing.T) { + for _, failure := range []string{"size", "option"} { + t.Run(failure, func(t *testing.T) { + ctx := t.Context() + dir := t.TempDir() + first, err := newLocalCacheStore(dir) + require.NoError(t, err) + second, err := newLocalCacheStore(dir) + require.NoError(t, err) + w, err := first.Writer(ctx, content.WithRef("layer")) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, w.Close()) }) + if failure == "size" { + err = w.Commit(ctx, 1, "") + require.ErrorIs(t, err, cerrdefs.ErrFailedPrecondition) + } else { + sentinel := errors.New("invalid commit option") + err = w.Commit(ctx, 0, "", func(*content.Info) error { return sentinel }) + require.ErrorIs(t, err, sentinel) + } + // Commit closes the writer even on failure, including the file lock. + require.NoError(t, second.Abort(ctx, "layer")) + }) + } +} + func TestSolveRejectsInvalidLocalExporterMode(t *testing.T) { st := llb.Scratch().File( llb.Mkfile("fresh.txt", 0600, []byte("fresh")),