From 3bfb9cfc4d95c836e11c94fdc19bdb92e882c2a6 Mon Sep 17 00:00:00 2001 From: Erik Sipsma Date: Fri, 4 Mar 2022 19:48:08 -0800 Subject: [PATCH] Fully initialize progress controller in FromRemote Before this, the worker's FromRemote method only partially intialized progress controllers, leaving out the vertex digest. This meant that only status updates would be sent and vertex start/stops would not be sent. Signed-off-by: Erik Sipsma --- solver/llbsolver/ops/exec.go | 12 ++++++++++++ solver/llbsolver/ops/file.go | 12 ++++++++++++ worker/base/worker.go | 19 ++++++++++++++++--- 3 files changed, 40 insertions(+), 3 deletions(-) diff --git a/solver/llbsolver/ops/exec.go b/solver/llbsolver/ops/exec.go index 51df9ff37..6cca733c0 100644 --- a/solver/llbsolver/ops/exec.go +++ b/solver/llbsolver/ops/exec.go @@ -21,6 +21,8 @@ import ( "github.com/moby/buildkit/solver/llbsolver/errdefs" "github.com/moby/buildkit/solver/llbsolver/mounts" "github.com/moby/buildkit/solver/pb" + "github.com/moby/buildkit/util/progress" + "github.com/moby/buildkit/util/progress/controller" "github.com/moby/buildkit/util/progress/logs" utilsystem "github.com/moby/buildkit/util/system" "github.com/moby/buildkit/worker" @@ -43,6 +45,7 @@ type execOp struct { platform *pb.Platform numInputs int parallelism *semaphore.Weighted + vtx solver.Vertex } func NewExecOp(v solver.Vertex, op *pb.Op_Exec, platform *pb.Platform, cm cache.Manager, parallelism *semaphore.Weighted, sm *session.Manager, exec executor.Executor, w worker.Worker) (solver.Op, error) { @@ -60,6 +63,7 @@ func NewExecOp(v solver.Vertex, op *pb.Op_Exec, platform *pb.Platform, cm cache. w: w, platform: platform, parallelism: parallelism, + vtx: v, }, nil } @@ -141,6 +145,14 @@ func (e *execOp) CacheMap(ctx context.Context, g session.Group, index int) (*sol ComputeDigestFunc solver.ResultBasedCacheFunc PreprocessFunc solver.PreprocessFunc }, e.numInputs), + Opts: solver.CacheOpts(map[interface{}]interface{}{ + cache.ProgressKey{}: &controller.Controller{ + WriterFactory: progress.FromContext(ctx), + Digest: e.vtx.Digest(), + Name: e.vtx.Name(), + ProgressGroup: e.vtx.Options().ProgressGroup, + }, + }), } deps, err := e.getMountDeps() diff --git a/solver/llbsolver/ops/file.go b/solver/llbsolver/ops/file.go index 27deec5ca..012ef4cc1 100644 --- a/solver/llbsolver/ops/file.go +++ b/solver/llbsolver/ops/file.go @@ -19,6 +19,8 @@ import ( "github.com/moby/buildkit/solver/llbsolver/ops/fileoptypes" "github.com/moby/buildkit/solver/pb" "github.com/moby/buildkit/util/flightcontrol" + "github.com/moby/buildkit/util/progress" + "github.com/moby/buildkit/util/progress/controller" "github.com/moby/buildkit/worker" digest "github.com/opencontainers/go-digest" "github.com/pkg/errors" @@ -35,6 +37,7 @@ type fileOp struct { solver *FileOpSolver numInputs int parallelism *semaphore.Weighted + vtx solver.Vertex } func NewFileOp(v solver.Vertex, op *pb.Op_File, cm cache.Manager, parallelism *semaphore.Weighted, w worker.Worker) (solver.Op, error) { @@ -48,6 +51,7 @@ func NewFileOp(v solver.Vertex, op *pb.Op_File, cm cache.Manager, parallelism *s w: w, solver: NewFileOpSolver(w, &file.Backend{}, file.NewRefManager(cm)), parallelism: parallelism, + vtx: v, }, nil } @@ -134,6 +138,14 @@ func (f *fileOp) CacheMap(ctx context.Context, g session.Group, index int) (*sol ComputeDigestFunc solver.ResultBasedCacheFunc PreprocessFunc solver.PreprocessFunc }, f.numInputs), + Opts: solver.CacheOpts(map[interface{}]interface{}{ + cache.ProgressKey{}: &controller.Controller{ + WriterFactory: progress.FromContext(ctx), + Digest: f.vtx.Digest(), + Name: f.vtx.Name(), + ProgressGroup: f.vtx.Options().ProgressGroup, + }, + }), } for idx, m := range selectors { diff --git a/worker/base/worker.go b/worker/base/worker.go index 7762f9850..fa8b7692d 100644 --- a/worker/base/worker.go +++ b/worker/base/worker.go @@ -392,11 +392,24 @@ func (w *Worker) FromRemote(ctx context.Context, remote *solver.Remote) (ref cac } } + var pg progress.Controller + optGetter := solver.CacheOptGetterOf(ctx) + if optGetter != nil { + if kv := optGetter(false, cache.ProgressKey{}); kv != nil { + if v, ok := kv[cache.ProgressKey{}].(progress.Controller); ok { + pg = v + } + } + } + if pg == nil { + pg = &controller.Controller{ + WriterFactory: progress.FromContext(ctx), + } + } + descHandler := &cache.DescHandler{ Provider: func(session.Group) content.Provider { return remote.Provider }, - Progress: &controller.Controller{ - WriterFactory: progress.FromContext(ctx), - }, + Progress: pg, } snapshotLabels := func([]ocispecs.Descriptor, int) map[string]string { return nil } if cd, ok := remote.Provider.(interface {