From f41283caede23e32ed552fb015146aea696faf43 Mon Sep 17 00:00:00 2001 From: Tonis Tiigi Date: Mon, 25 Sep 2017 20:57:38 -0700 Subject: [PATCH] solver: fix shared request progress and cancellation Signed-off-by: Tonis Tiigi --- api/services/control/control.proto | 1 - client/graph.go | 11 --- client/solve.go | 1 - control/control.go | 92 +++++++++----------- solver/build.go | 16 +--- solver/jobs.go | 129 +++++++++++++++++++--------- solver/load.go | 80 +++++++++-------- solver/solver.go | 83 +++++++++--------- solver/vertex.go | 42 +++------ source/containerimage/pull.go | 9 +- util/bgfunc/bgfunc.go | 10 +-- util/progress/progressui/display.go | 7 -- 12 files changed, 236 insertions(+), 245 deletions(-) diff --git a/api/services/control/control.proto b/api/services/control/control.proto index ded70e3d4..08ca9e758 100644 --- a/api/services/control/control.proto +++ b/api/services/control/control.proto @@ -68,7 +68,6 @@ message Vertex { google.protobuf.Timestamp started = 5 [(gogoproto.stdtime) = true ]; google.protobuf.Timestamp completed = 6 [(gogoproto.stdtime) = true ]; string error = 7; // typed errors? - string parent = 8 [(gogoproto.customtype) = "github.com/opencontainers/go-digest.Digest", (gogoproto.nullable) = false]; } message VertexStatus { diff --git a/client/graph.go b/client/graph.go index d9fcd0656..1ea0843b5 100644 --- a/client/graph.go +++ b/client/graph.go @@ -14,7 +14,6 @@ type Vertex struct { Completed *time.Time Cached bool Error string - Parent digest.Digest } type VertexStatus struct { @@ -40,13 +39,3 @@ type SolveStatus struct { Statuses []*VertexStatus Logs []*VertexLog } - -// -// type VertexEvent struct { -// ID digest.Digest -// Vertex digest.Digest -// Name string -// Total int -// Current int -// Timestamp int64 -// } diff --git a/client/solve.go b/client/solve.go index 638d44072..3e072f2eb 100644 --- a/client/solve.go +++ b/client/solve.go @@ -132,7 +132,6 @@ func (c *Client) Solve(ctx context.Context, r io.Reader, opt SolveOpt, statusCha Completed: v.Completed, Error: v.Error, Cached: v.Cached, - Parent: v.Parent, }) } for _, v := range resp.Statuses { diff --git a/control/control.go b/control/control.go index 65619b4a5..acc831658 100644 --- a/control/control.go +++ b/control/control.go @@ -90,15 +90,6 @@ func (c *Controller) Solve(ctx context.Context, req *controlapi.SolveRequest) (* } } - var vertex solver.Vertex - if req.Frontend == "" { - v, err := solver.LoadLLB(req.Definition) - if err != nil { - return nil, errors.Wrap(err, "failed to load llb definition") - } - vertex = v - } - ctx = session.NewContext(ctx, req.Session) var expi exporter.ExporterInstance @@ -114,7 +105,7 @@ func (c *Controller) Solve(ctx context.Context, req *controlapi.SolveRequest) (* } } - if err := c.solver.Solve(ctx, req.Ref, frontend, vertex, expi, req.FrontendAttrs); err != nil { + if err := c.solver.Solve(ctx, req.Ref, frontend, req.Definition, expi, req.FrontendAttrs); err != nil { return nil, err } return &controlapi.SolveResponse{}, nil @@ -130,49 +121,44 @@ func (c *Controller) Status(req *controlapi.StatusRequest, stream controlapi.Con eg.Go(func() error { for { - select { - case <-ctx.Done(): - return ctx.Err() - case ss, ok := <-ch: - if !ok { - return nil - } - sr := controlapi.StatusResponse{} - for _, v := range ss.Vertexes { - sr.Vertexes = append(sr.Vertexes, &controlapi.Vertex{ - Digest: v.Digest, - Inputs: v.Inputs, - Name: v.Name, - Started: v.Started, - Completed: v.Completed, - Error: v.Error, - Cached: v.Cached, - Parent: v.Parent, - }) - } - for _, v := range ss.Statuses { - sr.Statuses = append(sr.Statuses, &controlapi.VertexStatus{ - ID: v.ID, - Vertex: v.Vertex, - Name: v.Name, - Current: v.Current, - Total: v.Total, - Timestamp: v.Timestamp, - Started: v.Started, - Completed: v.Completed, - }) - } - for _, v := range ss.Logs { - sr.Logs = append(sr.Logs, &controlapi.VertexLog{ - Vertex: v.Vertex, - Stream: int64(v.Stream), - Msg: v.Data, - Timestamp: v.Timestamp, - }) - } - if err := stream.SendMsg(&sr); err != nil { - return err - } + ss, ok := <-ch + if !ok { + return nil + } + sr := controlapi.StatusResponse{} + for _, v := range ss.Vertexes { + sr.Vertexes = append(sr.Vertexes, &controlapi.Vertex{ + Digest: v.Digest, + Inputs: v.Inputs, + Name: v.Name, + Started: v.Started, + Completed: v.Completed, + Error: v.Error, + Cached: v.Cached, + }) + } + for _, v := range ss.Statuses { + sr.Statuses = append(sr.Statuses, &controlapi.VertexStatus{ + ID: v.ID, + Vertex: v.Vertex, + Name: v.Name, + Current: v.Current, + Total: v.Total, + Timestamp: v.Timestamp, + Started: v.Started, + Completed: v.Completed, + }) + } + for _, v := range ss.Logs { + sr.Logs = append(sr.Logs, &controlapi.VertexLog{ + Vertex: v.Vertex, + Stream: int64(v.Stream), + Msg: v.Data, + Timestamp: v.Timestamp, + }) + } + if err := stream.SendMsg(&sr); err != nil { + return err } } }) diff --git a/solver/build.go b/solver/build.go index cf1c77687..5dd331e3d 100644 --- a/solver/build.go +++ b/solver/build.go @@ -109,21 +109,7 @@ func (b *buildOp) Run(ctx context.Context, inputs []Reference) (outputs []Refere lm.Unmount() lm = nil - v, err := LoadLLB(def) - if err != nil { - return nil, err - } - - if len(v.Inputs()) == 0 { - return nil, errors.New("required vertex needs to have inputs") - } - - index := v.Inputs()[0].Index - v = v.Inputs()[0].Vertex - - vv := toInternalVertex(v) - - newref, err := b.s.loadAndSolveChildVertex(ctx, b.v.Digest(), vv, index) + newref, err := b.s.loadAndSolve(ctx, b.v.Digest(), def) if err != nil { return nil, err } diff --git a/solver/jobs.go b/solver/jobs.go index fd3b80624..5d50b5ce5 100644 --- a/solver/jobs.go +++ b/solver/jobs.go @@ -7,6 +7,7 @@ import ( "github.com/moby/buildkit/client" "github.com/moby/buildkit/session" + "github.com/moby/buildkit/solver/pb" "github.com/moby/buildkit/util/progress" digest "github.com/opencontainers/go-digest" "github.com/pkg/errors" @@ -26,7 +27,7 @@ type jobList struct { } type state struct { - jobs map[*job]struct{} + jobs map[*job]*vertex solver VertexSolver mpw *progress.MultiWriter } @@ -90,7 +91,7 @@ func (jl *jobList) get(id string) (*job, error) { } } -func (jl *jobList) loadAndSolveChildVertex(ctx context.Context, dgst digest.Digest, vv *vertex, index Index, f ResolveOpFunc, cache InstructionCache) (Reference, error) { +func (jl *jobList) loadAndSolve(ctx context.Context, dgst digest.Digest, ops [][]byte, f ResolveOpFunc, cache InstructionCache) (Reference, error) { jl.mu.Lock() st, ok := jl.actives[dgst] @@ -99,18 +100,19 @@ func (jl *jobList) loadAndSolveChildVertex(ctx context.Context, dgst digest.Dige return nil, errors.Errorf("no such parent vertex: %v", dgst) } - var newst *state + var inp *Input for j := range st.jobs { var err error - newst, err = j.loadInternal(vv, f) + inp, err = j.loadInternal(ops, f) if err != nil { jl.mu.Unlock() return nil, err } } + st = jl.actives[inp.Vertex.Digest()] jl.mu.Unlock() - return getRef(newst.solver, ctx, vv, index, cache) + return getRef(st.solver, ctx, inp.Vertex.(*vertex), inp.Index, cache) // TODO: combine to pass single input } type job struct { @@ -121,49 +123,57 @@ type job struct { cache InstructionCache } -func (j *job) load(v *vertex, f ResolveOpFunc) error { +func (j *job) load(ops [][]byte, resolveOp ResolveOpFunc) (*Input, error) { j.l.mu.Lock() defer j.l.mu.Unlock() - _, err := j.loadInternal(v, f) - return err + return j.loadInternal(ops, resolveOp) } -func (j *job) loadInternal(v *vertex, f ResolveOpFunc) (*state, error) { - for _, inp := range v.inputs { - if _, err := j.loadInternal(inp.vertex, f); err != nil { - return nil, err +func (j *job) loadInternal(ops [][]byte, resolveOp ResolveOpFunc) (*Input, error) { + vtx, idx, err := loadLLB(ops, func(dgst digest.Digest, op *pb.Op, load func(digest.Digest) (interface{}, error)) (interface{}, error) { + if st, ok := j.l.actives[dgst]; ok { + if vtx, ok := st.jobs[j]; ok { + return vtx, nil + } } - } - - dgst := v.Digest() - st, ok := j.l.actives[dgst] - if !ok { - st = &state{ - jobs: map[*job]struct{}{}, - mpw: progress.NewMultiWriter(progress.WithMetadata("vertex", dgst)), - } - op, err := f(v) + vtx, err := newVertex(dgst, op, load) if err != nil { return nil, err } - ctx := progress.WithProgress(context.Background(), st.mpw) - ctx = session.NewContext(ctx, j.session) // TODO: support multiple - s, err := newVertexSolver(ctx, v, op, j.cache, j.getSolver) - if err != nil { - return nil, err + st, ok := j.l.actives[dgst] + if !ok { + st = &state{ + jobs: map[*job]*vertex{}, + mpw: progress.NewMultiWriter(progress.WithMetadata("vertex", dgst)), + } + op, err := resolveOp(vtx) + if err != nil { + return nil, err + } + ctx := progress.WithProgress(context.Background(), st.mpw) + ctx = session.NewContext(ctx, j.session) // TODO: support multiple + + s, err := newVertexSolver(ctx, vtx, op, j.cache, j.getSolver) + if err != nil { + return nil, err + } + st.solver = s + + j.l.actives[dgst] = st } - st.solver = s - - j.l.actives[dgst] = st + if _, ok := st.jobs[j]; !ok { + j.pw.Write(vtx.Digest().String(), vtx.clientVertex) + st.mpw.Add(j.pw) + st.jobs[j] = vtx + } + return vtx, nil + }) + if err != nil { + return nil, err } - if _, ok := st.jobs[j]; !ok { - j.pw.Write(v.Digest().String(), v.clientVertex) - st.mpw.Add(j.pw) - st.jobs[j] = struct{}{} - } - return st, nil + return &Input{Vertex: vtx.(*vertex), Index: idx}, nil } func (j *job) discard() { @@ -209,7 +219,7 @@ func getRef(s VertexSolver, ctx context.Context, v *vertex, index Index, cache I return nil, err } if ref != nil { - v.notifyCompleted(ctx, true, nil) + markCached(ctx, v.clientVertex) return ref.(Reference), nil } @@ -230,7 +240,7 @@ func getRef(s VertexSolver, ctx context.Context, v *vertex, index Index, cache I return nil, err } if ref != nil { - v.notifyCompleted(ctx, true, nil) + markCached(ctx, v.clientVertex) return ref.(Reference), nil } continue @@ -240,7 +250,13 @@ func getRef(s VertexSolver, ctx context.Context, v *vertex, index Index, cache I } func (j *job) pipe(ctx context.Context, ch chan *client.SolveStatus) error { + vs := &vertexStream{cache: map[digest.Digest]*client.Vertex{}} pr := j.pr.Reader(ctx) + defer func() { + if enc := vs.encore(); len(enc) > 0 { + ch <- &client.SolveStatus{Vertexes: enc} + } + }() for { p, err := pr.Read(ctx) if err != nil { @@ -253,7 +269,7 @@ func (j *job) pipe(ctx context.Context, ch chan *client.SolveStatus) error { for _, p := range p { switch v := p.Sys.(type) { case client.Vertex: - ss.Vertexes = append(ss.Vertexes, &v) + ss.Vertexes = append(ss.Vertexes, vs.append(v)...) case progress.Status: vtx, ok := p.Meta("vertex") @@ -290,3 +306,38 @@ func (j *job) pipe(ctx context.Context, ch chan *client.SolveStatus) error { } } } + +type vertexStream struct { + cache map[digest.Digest]*client.Vertex +} + +func (vs *vertexStream) append(v client.Vertex) []*client.Vertex { + var out []*client.Vertex + vs.cache[v.Digest] = &v + if v.Cached { + for _, inp := range v.Inputs { + if inpv, ok := vs.cache[inp]; ok { + if !inpv.Cached && inpv.Completed == nil { + inpv.Cached = true + inpv.Started = v.Completed + inpv.Completed = v.Completed + out = append(vs.append(*inpv), inpv) + } + } + } + } + return append(out, &v) +} + +func (vs *vertexStream) encore() []*client.Vertex { + var out []*client.Vertex + for _, v := range vs.cache { + if v.Started != nil && v.Completed == nil { + now := time.Now() + v.Completed = &now + v.Error = context.Canceled.Error() + out = append(out, v) + } + } + return out +} diff --git a/solver/load.go b/solver/load.go index b6929bb59..385d79ac0 100644 --- a/solver/load.go +++ b/solver/load.go @@ -8,32 +8,17 @@ import ( "github.com/pkg/errors" ) -func LoadLLB(ops [][]byte) (Vertex, error) { - if len(ops) == 0 { - return nil, errors.New("invalid empty definition") - } - - allOps := make(map[digest.Digest]*pb.Op) - - var lastOp *pb.Op - var lastDigest digest.Digest - - for _, dt := range ops { - var op pb.Op - if err := (&op).Unmarshal(dt); err != nil { - return nil, errors.Wrap(err, "failed to parse llb proto op") +func newVertex(dgst digest.Digest, op *pb.Op, load func(digest.Digest) (interface{}, error)) (*vertex, error) { + vtx := &vertex{sys: op.Op, digest: dgst, name: llbOpName(op)} + for _, in := range op.Inputs { + sub, err := load(in.Digest) + if err != nil { + return nil, err } - lastOp = &op - lastDigest = digest.FromBytes(dt) - allOps[lastDigest] = &op + vtx.inputs = append(vtx.inputs, &input{index: Index(in.Index), vertex: sub.(*vertex)}) } - - delete(allOps, lastDigest) // avoid loops - - cache := make(map[digest.Digest]*vertex) - - // TODO: validate the connections - return loadLLBVertexRecursive(lastDigest, lastOp, allOps, cache) + vtx.initClientVertex() + return vtx, nil } func toInternalVertex(v Vertex) *vertex { @@ -55,26 +40,45 @@ func loadInternalVertexHelper(v Vertex, cache map[digest.Digest]*vertex) *vertex return vtx } -func loadLLBVertexRecursive(dgst digest.Digest, op *pb.Op, all map[digest.Digest]*pb.Op, cache map[digest.Digest]*vertex) (*vertex, error) { - if v, ok := cache[dgst]; ok { - return v, nil +func loadLLB(ops [][]byte, fn func(digest.Digest, *pb.Op, func(digest.Digest) (interface{}, error)) (interface{}, error)) (interface{}, Index, error) { + if len(ops) == 0 { + return nil, 0, errors.New("invalid empty definition") } - vtx := &vertex{sys: op.Op, digest: dgst, name: llbOpName(op)} - for _, in := range op.Inputs { - dgst := digest.Digest(in.Digest) - op, ok := all[dgst] - if !ok { - return nil, errors.Errorf("failed to find %s", in) + + allOps := make(map[digest.Digest]*pb.Op) + + var dgst digest.Digest + + for _, dt := range ops { + var op pb.Op + if err := (&op).Unmarshal(dt); err != nil { + return nil, 0, errors.Wrap(err, "failed to parse llb proto op") } - sub, err := loadLLBVertexRecursive(dgst, op, all, cache) + dgst = digest.FromBytes(dt) + allOps[dgst] = &op + } + + lastOp := allOps[dgst] + delete(allOps, dgst) + dgst = lastOp.Inputs[0].Digest + + cache := make(map[digest.Digest]interface{}) + + var rec func(dgst digest.Digest) (interface{}, error) + rec = func(dgst digest.Digest) (interface{}, error) { + if v, ok := cache[dgst]; ok { + return v, nil + } + v, err := fn(dgst, allOps[dgst], rec) if err != nil { return nil, err } - vtx.inputs = append(vtx.inputs, &input{index: Index(in.Index), vertex: sub}) + cache[dgst] = v + return v, nil } - vtx.initClientVertex() - cache[dgst] = vtx - return vtx, nil + + v, err := rec(dgst) + return v, Index(lastOp.Inputs[0].Index), err } func llbOpName(op *pb.Op) string { diff --git a/solver/solver.go b/solver/solver.go index 142c6d3b9..2b0fa0aa6 100644 --- a/solver/solver.go +++ b/solver/solver.go @@ -4,11 +4,13 @@ import ( "encoding/json" "fmt" "sync" + "time" "github.com/moby/buildkit/cache" "github.com/moby/buildkit/client" "github.com/moby/buildkit/exporter" "github.com/moby/buildkit/frontend" + "github.com/moby/buildkit/identity" "github.com/moby/buildkit/solver/pb" "github.com/moby/buildkit/source" "github.com/moby/buildkit/util/bgfunc" @@ -80,7 +82,7 @@ func New(resolve ResolveOpFunc, cache InstructionCache, imageSource source.Sourc return &Solver{resolve: resolve, jobs: newJobList(), cache: cache, imageSource: imageSource} } -func (s *Solver) Solve(ctx context.Context, id string, f frontend.Frontend, v Vertex, exp exporter.ExporterInstance, frontendOpt map[string]string) error { +func (s *Solver) Solve(ctx context.Context, id string, f frontend.Frontend, dt [][]byte, exp exporter.ExporterInstance, frontendOpt map[string]string) error { ctx, cancel := context.WithCancel(ctx) defer cancel() @@ -88,19 +90,6 @@ func (s *Solver) Solve(ctx context.Context, id string, f frontend.Frontend, v Ve defer closeProgressWriter() - var vv *vertex - var index Index - if v != nil { - if len(v.Inputs()) == 0 { - return errors.New("required vertex needs to have inputs") - } - - index = v.Inputs()[0].Index - v = v.Inputs()[0].Vertex - - vv = toInternalVertex(v) - } - ctx, j, err := s.jobs.new(ctx, id, pr, s.cache) if err != nil { return err @@ -108,13 +97,14 @@ func (s *Solver) Solve(ctx context.Context, id string, f frontend.Frontend, v Ve var ref Reference var exporterOpt map[string]interface{} - // solver: s.getRef, - if vv != nil { - if err := j.load(vv, s.resolve); err != nil { + if dt != nil { + var inp *Input + inp, err = j.load(dt, s.resolve) + if err != nil { j.discard() return err } - ref, err = j.getRef(ctx, vv, index) + ref, err = j.getRef(ctx, inp.Vertex.(*vertex), inp.Index) } else { ref, exporterOpt, err = f.Solve(ctx, &llbBridge{ job: j, @@ -140,11 +130,15 @@ func (s *Solver) Solve(ctx context.Context, id string, f frontend.Frontend, v Ve } if exp != nil { - vv.notifyStarted(ctx) - pw, _, ctx := progress.FromContext(ctx, progress.WithMetadata("vertex", vv.Digest())) + v := client.Vertex{ + Digest: digest.FromBytes([]byte(identity.NewID())), + Name: exp.Name(), + } + notifyStarted(ctx, &v) + pw, _, ctx := progress.FromContext(ctx, progress.WithMetadata("vertex", v.Digest)) defer pw.Close() err := exp.Export(ctx, immutable, exporterOpt) - vv.notifyCompleted(ctx, false, err) + notifyCompleted(ctx, &v, err) if err != nil { return err } @@ -161,8 +155,8 @@ func (s *Solver) Status(ctx context.Context, id string, statusChan chan *client. return j.pipe(ctx, statusChan) } -func (s *Solver) loadAndSolveChildVertex(ctx context.Context, dgst digest.Digest, vv *vertex, index Index) (Reference, error) { - return s.jobs.loadAndSolveChildVertex(ctx, dgst, vv, index, s.resolve, s.cache) +func (s *Solver) loadAndSolve(ctx context.Context, dgst digest.Digest, def [][]byte) (Reference, error) { + return s.jobs.loadAndSolve(ctx, dgst, def, s.resolve, s.cache) } type VertexSolver interface { @@ -181,6 +175,7 @@ type vertexInput struct { type vertexSolver struct { inputs []*vertexInput v *vertex + cv client.Vertex op Op cache InstructionCache refs []*sharedRef @@ -196,7 +191,7 @@ type vertexSolver struct { type resolveF func(digest.Digest) (VertexSolver, error) -func newVertexSolver(ctx context.Context, v *vertex, op Op, c InstructionCache, resolve resolveF) (VertexSolver, error) { +func newVertexSolver(ctx context.Context, v *vertex, op Op, c InstructionCache, resolve resolveF) (*vertexSolver, error) { inputs := make([]*vertexInput, len(v.inputs)) for i, in := range v.inputs { s, err := resolve(in.vertex.digest) @@ -207,6 +202,7 @@ func newVertexSolver(ctx context.Context, v *vertex, op Op, c InstructionCache, if err != nil { return nil, err } + ev.Cancel() inputs[i] = &vertexInput{ solver: s, ev: ev, @@ -216,12 +212,26 @@ func newVertexSolver(ctx context.Context, v *vertex, op Op, c InstructionCache, ctx: ctx, inputs: inputs, v: v, + cv: v.clientVertex, op: op, cache: c, signal: newSignaller(), }, nil } +func markCached(ctx context.Context, cv client.Vertex) { + pw, _, _ := progress.FromContext(ctx) + defer pw.Close() + + if cv.Started == nil { + now := time.Now() + cv.Started = &now + cv.Completed = &now + cv.Cached = true + } + pw.Write(cv.Digest.String(), cv) +} + func (vs *vertexSolver) CacheKey(ctx context.Context, index Index) (digest.Digest, error) { vs.mu.Lock() defer vs.mu.Unlock() @@ -369,6 +379,7 @@ func (vs *vertexSolver) run(ctx context.Context, signal func()) (retErr error) { } if ref != nil { inp.ref = ref.(Reference) + markCached(ctx, inp.solver.(*vertexSolver).cv) return nil } } @@ -459,9 +470,9 @@ func (vs *vertexSolver) run(ctx context.Context, signal func()) (retErr error) { } // no cache hit. start evaluating the node - vs.v.notifyStarted(ctx) + notifyStarted(ctx, &vs.cv) defer func() { - vs.v.notifyCompleted(ctx, false, retErr) + notifyCompleted(ctx, &vs.cv, retErr) }() refs, err := vs.op.Run(ctx, inputRefs) @@ -472,11 +483,13 @@ func (vs *vertexSolver) run(ctx context.Context, signal func()) (retErr error) { for i, r := range refs { sr[i] = newSharedRef(r) } + vs.mu.Lock() vs.refs = sr + vs.mu.Unlock() // store the cacheKeys for current refs if vs.cache != nil { - cacheKey, err := vs.mainCacheKey() + cacheKey, err := vs.lastCacheKey() if err != nil { return err } @@ -556,23 +569,11 @@ type resolveImageConfig interface { } func (s *llbBridge) Solve(ctx context.Context, dt [][]byte) (cache.ImmutableRef, error) { - v, err := LoadLLB(dt) + inp, err := s.job.load(dt, s.resolveOp) if err != nil { return nil, err } - if len(v.Inputs()) == 0 { - return nil, errors.New("required vertex needs to have inputs") - } - - index := v.Inputs()[0].Index - v = v.Inputs()[0].Vertex - - vv := toInternalVertex(v) - - if err := s.job.load(vv, s.resolveOp); err != nil { - return nil, err - } - ref, err := s.job.getRef(ctx, vv, index) + ref, err := s.job.getRef(ctx, inp.Vertex.(*vertex), inp.Index) if err != nil { return nil, err } diff --git a/solver/vertex.go b/solver/vertex.go index 592e6c310..261a0a77a 100644 --- a/solver/vertex.go +++ b/solver/vertex.go @@ -1,13 +1,13 @@ package solver import ( + "context" "sync" "time" "github.com/moby/buildkit/client" "github.com/moby/buildkit/util/progress" digest "github.com/opencontainers/go-digest" - "golang.org/x/net/context" ) // Vertex is one node in the build graph @@ -78,44 +78,26 @@ func (v *vertex) Name() string { return v.name } -func (v *vertex) inputRequiresExport(i int) bool { - return true // TODO -} - -func (v *vertex) notifyStarted(ctx context.Context) { - v.recursiveMarkCached(ctx) +func notifyStarted(ctx context.Context, v *client.Vertex) { pw, _, _ := progress.FromContext(ctx) defer pw.Close() now := time.Now() - v.clientVertex.Started = &now - v.clientVertex.Completed = nil - pw.Write(v.Digest().String(), v.clientVertex) + v.Started = &now + v.Completed = nil + pw.Write(v.Digest.String(), *v) } -func (v *vertex) notifyCompleted(ctx context.Context, cached bool, err error) { +func notifyCompleted(ctx context.Context, v *client.Vertex, err error) { pw, _, _ := progress.FromContext(ctx) defer pw.Close() now := time.Now() - v.recursiveMarkCached(ctx) - if v.clientVertex.Started == nil { - v.clientVertex.Started = &now + if v.Started == nil { + v.Started = &now } - v.clientVertex.Completed = &now - v.clientVertex.Cached = cached + v.Completed = &now + v.Cached = false if err != nil { - v.clientVertex.Error = err.Error() + v.Error = err.Error() } - pw.Write(v.Digest().String(), v.clientVertex) -} - -func (v *vertex) recursiveMarkCached(ctx context.Context) { - for _, inp := range v.inputs { - inp.vertex.notifyMu.Lock() - if inp.vertex.clientVertex.Started == nil { - inp.vertex.recursiveMarkCached(ctx) - inp.vertex.notifyCompleted(ctx, true, nil) - } - inp.vertex.notifyMu.Unlock() - } - + pw.Write(v.Digest.String(), *v) } diff --git a/source/containerimage/pull.go b/source/containerimage/pull.go index cdafe2b2d..177e4d87e 100644 --- a/source/containerimage/pull.go +++ b/source/containerimage/pull.go @@ -295,11 +295,12 @@ func showProgress(ctx context.Context, ongoing *jobs, cs content.Store) { for _, j := range ongoing.jobs() { refKey := remotes.MakeRefKey(ctx, j.Descriptor) if a, ok := actives[refKey]; ok { + started := j.started pw.Write(j.Digest.String(), progress.Status{ Action: a.Status, Total: int(a.Total), Current: int(a.Offset), - Started: &j.started, + Started: &started, }) continue } @@ -318,12 +319,14 @@ func showProgress(ctx context.Context, ongoing *jobs, cs content.Store) { } if done || j.done { + started := j.started + createdAt := info.CreatedAt pw.Write(j.Digest.String(), progress.Status{ Action: "done", Current: int(info.Size), Total: int(info.Size), - Completed: &info.CreatedAt, - Started: &j.started, + Completed: &createdAt, + Started: &started, }) } } diff --git a/util/bgfunc/bgfunc.go b/util/bgfunc/bgfunc.go index a5c32b235..8621f56d7 100644 --- a/util/bgfunc/bgfunc.go +++ b/util/bgfunc/bgfunc.go @@ -41,10 +41,13 @@ func (f *F) run() { f.runMu.Lock() if !f.running && !f.done { f.running = true + ctx, cancel := context.WithCancel(f.mainCtx) + ctxErr := make(chan error, 1) + f.cancelCtx = cancel + f.ctxErr = ctxErr go func() { var err error var nodone bool - ctxErr := make(chan error, 1) defer func() { // release all cancellations f.runMu.Lock() @@ -59,11 +62,6 @@ func (f *F) run() { ctxErr <- err f.mu.Unlock() }() - ctx, cancel := context.WithCancel(f.mainCtx) - f.runMu.Lock() - f.cancelCtx = cancel - f.ctxErr = ctxErr - f.runMu.Unlock() err = f.f(ctx, func() { f.cond.Broadcast() }) diff --git a/util/progress/progressui/display.go b/util/progress/progressui/display.go index a472a2fe2..2908277bb 100644 --- a/util/progress/progressui/display.go +++ b/util/progress/progressui/display.go @@ -100,10 +100,6 @@ func (t *trace) update(s *client.SolveStatus) { t.byDigest[v.Digest] = &vertex{ byID: make(map[string]*status), } - } else { - if prev.Parent != v.Parent { // skip vertexes already in list for other parents - continue - } } if v.Started != nil && (prev == nil || prev.Started == nil) { if t.localTimeDiff == 0 { @@ -112,9 +108,6 @@ func (t *trace) update(s *client.SolveStatus) { t.vertexes = append(t.vertexes, t.byDigest[v.Digest]) } t.byDigest[v.Digest].Vertex = v - if v.Parent != "" { - t.byDigest[v.Digest].indent = t.byDigest[v.Parent].indent + "=> " - } } for _, s := range s.Statuses { v, ok := t.byDigest[s.Vertex]