From 3dffed121b433347b2428dbad47dba0ba40639f1 Mon Sep 17 00:00:00 2001 From: CrazyMax <1951866+crazy-max@users.noreply.github.com> Date: Wed, 23 Sep 2026 11:23:40 +0200 Subject: [PATCH] solve: lock local cache ingests across clients Independent local content stores only coordinate writers within each store instance. Concurrent cache exports to the same directory can open the same ingest files, causing rename failures or incomplete writes. Acquire a filesystem lock per ingest reference before opening a writer or aborting an ingest. Release it when the writer closes or commits, including error paths, and report contention as unavailable so existing retry logic can wait. Preserve references so interrupted uploads remain resumable, and retain lock files to avoid races with replacement inodes. Signed-off-by: CrazyMax <1951866+crazy-max@users.noreply.github.com> --- client/client_cache_test.go | 142 ++++++++++++++++++++++++++++++++++++ client/client_test.go | 1 + client/solve.go | 94 +++++++++++++++++++++++- client/solve_test.go | 88 ++++++++++++++++++++++ 4 files changed, 323 insertions(+), 2 deletions(-) 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")),