diff --git a/pkg/build/anchor_prepass_test.go b/pkg/build/anchor_prepass_test.go new file mode 100644 index 0000000000..7a88c8b269 --- /dev/null +++ b/pkg/build/anchor_prepass_test.go @@ -0,0 +1,111 @@ +package build + +import ( + "context" + "fmt" + "time" + + "github.com/onsi/ginkgo/v2" + "github.com/onsi/gomega" + + "github.com/werf/werf/v3/pkg/build/image" + "github.com/werf/werf/v3/pkg/build/stage" + "github.com/werf/werf/v3/pkg/container_backend" + imagePkg "github.com/werf/werf/v3/pkg/image" +) + +type slowContentDependenciesStub struct { + *contentDependenciesStub + delay time.Duration +} + +var _ stage.Interface = (*slowContentDependenciesStub)(nil) + +func (s *slowContentDependenciesStub) GetContentDependencies(ctx context.Context, c stage.Conveyor, archive container_backend.BuildContextArchiver) (string, error) { + time.Sleep(s.delay) + + return s.contentDependenciesStub.GetContentDependencies(ctx, c, archive) +} + +// The graph is one dependent pair plus independent images: the dependency is +// the slowest image and comes first, its dependent comes last, so a prepass +// that ignores the graph hands them to different workers and leaves the +// dependent without a digest. +func newAnchorPrepassPhase(parallel bool) (*BuildPhase, []*image.Image) { + storageManager := &anchorLookupStorageManager{ + primaryStagesStorage: &anchorPrimaryStagesStorage{}, + secondaryStagesStorage: &fakeStagesStorage{}, + inPrimary: imagePkg.NewStageDescSet(), + inSecondary: imagePkg.NewStageDescSet(), + } + phase := newTestBuildPhase(storageManager, nil) + phase.Conveyor.Parallel = parallel + phase.Conveyor.ParallelTasksLimit = 4 + + var images []*image.Image + for i := range 6 { + name := fmt.Sprintf("image%d", i) + var dependencyNames []string + if i == 5 { + dependencyNames = append(dependencyNames, "image0") + } + img := newTestImage(name, true, dependencyNames...) + img.Conveyor = phase.Conveyor + anchor := &slowContentDependenciesStub{contentDependenciesStub: newContentDependenciesStub(stage.ImageSpec, name)} + if i == 0 { + anchor.delay = 300 * time.Millisecond + } + anchor.SetContentAnchor(true) + img.SetStages([]stage.Interface{anchor}) + images = append(images, img) + } + phase.Conveyor.imagesTree.SetImagesGraphForTests(newTestImagesGraph(images...)) + + return phase, images +} + +var _ = ginkgo.Describe("Anchor prepass", func() { + ginkgo.It("calculates the same digests whether or not the build is parallel", func(ctx ginkgo.SpecContext) { + serialPhase, serialImages := newAnchorPrepassPhase(false) + gomega.Expect(serialPhase.calculateAnchorDigests(ctx)).To(gomega.Succeed()) + + parallelPhase, parallelImages := newAnchorPrepassPhase(true) + gomega.Expect(parallelPhase.calculateAnchorDigests(ctx)).To(gomega.Succeed()) + + for i, img := range parallelImages { + gomega.Expect(img.GetAnchorDigest()).NotTo(gomega.BeEmpty(), img.Name) + gomega.Expect(img.GetAnchorDigest()).To(gomega.Equal(serialImages[i].GetAnchorDigest()), img.Name) + } + }) + + ginkgo.It("resolves every anchor exactly once in parallel", func(ctx ginkgo.SpecContext) { + phase, images := newAnchorPrepassPhase(true) + for i, img := range images { + img.SetAnchorDigest(fmt.Sprintf("anchor%d", i)) + } + storageManager := phase.Conveyor.StorageManager.(*anchorLookupStorageManager) + + gomega.Expect(phase.resolveAvailableContentAnchors(ctx)).To(gomega.Succeed()) + + gomega.Expect(storageManager.primaryLookups).To(gomega.BeZero()) + gomega.Expect(storageManager.secondaryLookups).To(gomega.BeZero()) + gomega.Expect(storageManager.cachedPrimaryLookups).To(gomega.Equal(len(images))) + gomega.Expect(storageManager.cachedSecondaryLookups).To(gomega.Equal(len(images))) + }) +}) + +var _ = ginkgo.Describe("Anchor prepass workers", func() { + ginkgo.DescribeTable("bounds concurrency by the build parallelism", func(parallel bool, limit int64, tasks, expected int) { + phase := newTestBuildPhase(nil, nil) + phase.Conveyor.Parallel = parallel + phase.Conveyor.ParallelTasksLimit = limit + + gomega.Expect(phase.prepassWorkers(tasks)).To(gomega.Equal(expected)) + }, + ginkgo.Entry("sequential build", false, int64(4), 10, 1), + ginkgo.Entry("single image", true, int64(4), 1, 1), + ginkgo.Entry("limit below the task count", true, int64(4), 10, 4), + ginkgo.Entry("limit above the task count", true, int64(16), 10, 10), + ginkgo.Entry("unlimited build", true, int64(0), 10, 10), + ) +}) diff --git a/pkg/build/build_phase.go b/pkg/build/build_phase.go index bafa3391f2..3508411115 100644 --- a/pkg/build/build_phase.go +++ b/pkg/build/build_phase.go @@ -176,18 +176,54 @@ func (phase *BuildPhase) calculateAnchorDigests(ctx context.Context) error { return nil } - for _, img := range graph.Nodes() { + nodes := graph.Nodes() + calculate := func(ctx context.Context, img *image.Image) error { dependencies := graph.Dependencies(img) if !canCalculateAnchorDigest(img, dependencies) { - continue + return nil + } + + return phase.calculateAnchorDigest(ctx, img, dependencies, false) + } + + workers := phase.prepassWorkers(len(nodes)) + if workers <= 1 { + for _, img := range nodes { + if err := calculate(ctx, img); err != nil { + return err + } } - if err := phase.calculateAnchorDigest(ctx, img, dependencies, false); err != nil { + return nil + } + + scheduler := newGraphScheduler(graph) + + return parallel.DoTasksDynamic(ctx, parallel.DoTasksOptions{MaxNumberOfWorkers: workers}, scheduler.next, func(ctx context.Context, taskID int) error { + img := nodes[taskID] + if err := calculate(ctx, img); err != nil { return err } + scheduler.complete(img) + + return nil + }) +} + +// prepassWorkers keeps the anchor prepass sequential unless the build itself is +// parallel, so a build that a user asked to run one image at a time does not +// start talking to the storage concurrently. +func (phase *BuildPhase) prepassWorkers(numberOfTasks int) int { + if !phase.Conveyor.Parallel || numberOfTasks <= 1 { + return 1 } - return nil + workers := int(phase.Conveyor.ParallelTasksLimit) + if workers <= 0 || workers > numberOfTasks { + workers = numberOfTasks + } + + return workers } func canCalculateAnchorDigest(img *image.Image, dependencies []*image.Image) bool { @@ -245,25 +281,42 @@ func (phase *BuildPhase) resolveAvailableContentAnchors(ctx context.Context) err return nil } - prepass := *phase - prepass.anchorPrepass = true - for _, img := range graph.Nodes() { + nodes := graph.Nodes() + resolve := func(ctx context.Context, img *image.Image) error { img.Requested = phase.isRequestedImage(img) if img.GetAnchorDigest() == "" { - continue + return nil } + prepass := *phase + prepass.anchorPrepass = true + prepass.StagesIterator = NewStagesIterator(phase.Conveyor) + var outBuf, errBuf bytes.Buffer resolveCtx := logboek.NewContext(ctx, logboek.Context(ctx).NewSubLogger(&outBuf, &errBuf)) - prepass.StagesIterator = NewStagesIterator(phase.Conveyor) if err := prepass.resolveContentAnchor(resolveCtx, img, false); err != nil { return fmt.Errorf("image %q: %w", img.Name, err) } img.ContentAnchorOutLog = bytes.Clone(outBuf.Bytes()) img.ContentAnchorErrLog = bytes.Clone(errBuf.Bytes()) + + return nil } - return nil + workers := phase.prepassWorkers(len(nodes)) + if workers <= 1 { + for _, img := range nodes { + if err := resolve(ctx, img); err != nil { + return err + } + } + + return nil + } + + return parallel.DoTasks(ctx, len(nodes), parallel.DoTasksOptions{MaxNumberOfWorkers: workers}, func(ctx context.Context, taskID int) error { + return resolve(ctx, nodes[taskID]) + }) } func (phase *BuildPhase) skipUnneededImages() { diff --git a/pkg/build/helpers_test.go b/pkg/build/helpers_test.go index 0f732d3f19..08d06fae36 100644 --- a/pkg/build/helpers_test.go +++ b/pkg/build/helpers_test.go @@ -45,6 +45,9 @@ func (m *anchorLookupStorageManager) GetStagesStorage() storage.PrimaryStagesSto } func (m *anchorLookupStorageManager) GetStageDescSetByDigestFromStagesStorageCached(_ context.Context, _, _ string, _ int64, stagesStorage storage.StagesStorage) (imagePkg.StageDescSet, error) { + m.mutex.Lock() + defer m.mutex.Unlock() + if stagesStorage == m.secondaryStagesStorage { m.cachedSecondaryLookups++ return m.inSecondary, nil diff --git a/pkg/build/unneeded_images_test.go b/pkg/build/unneeded_images_test.go index e5d3dc0195..91130d06c2 100644 --- a/pkg/build/unneeded_images_test.go +++ b/pkg/build/unneeded_images_test.go @@ -265,6 +265,7 @@ func (*nonLocalStorageManager) GetStagesStorage() storage.PrimaryStagesStorage { type anchorLookupStorageManager struct { manager.StorageManagerInterface + mutex sync.Mutex primaryStagesStorage storage.PrimaryStagesStorage secondaryStagesStorage storage.StagesStorage inPrimary imagePkg.StageDescSet @@ -276,6 +277,9 @@ type anchorLookupStorageManager struct { } func (m *anchorLookupStorageManager) GetStageDescSetByDigestWithCache(_ context.Context, _, _ string, _ int64) (imagePkg.StageDescSet, error) { + m.mutex.Lock() + defer m.mutex.Unlock() + m.primaryLookups++ return m.inPrimary, nil } @@ -285,6 +289,9 @@ func (m *anchorLookupStorageManager) GetSecondaryStagesStorageList() []storage.S } func (m *anchorLookupStorageManager) GetStageDescSetByDigestFromStagesStorageWithCache(_ context.Context, _, _ string, _ int64, stagesStorage storage.StagesStorage) (imagePkg.StageDescSet, error) { + m.mutex.Lock() + defer m.mutex.Unlock() + if stagesStorage == m.secondaryStagesStorage { m.secondaryLookups++ return m.inSecondary, nil