diff --git a/examples/llbout/example.go b/examples/llbout/example.go index c2a0b9195..2806ab5cf 100644 --- a/examples/llbout/example.go +++ b/examples/llbout/example.go @@ -13,6 +13,7 @@ func main() { alpine := llb.Image("docker.io/library/alpine:latest") mod3 := mod2.Run(llb.Meta{Args: []string{"/bin/cp", "-a", "/alpine/etc/passwd", "baz"}, Cwd: "/"}) mod3.AddMount("/alpine", alpine) + mod3.AddMount("/redis", busybox) mod4 := mod3.Run(llb.Meta{Args: []string{"/bin/ls", "-l", "/"}, Cwd: "/"}) res := mod4 diff --git a/solver/solver.go b/solver/solver.go index e4100e4c6..9b17ae833 100644 --- a/solver/solver.go +++ b/solver/solver.go @@ -94,7 +94,7 @@ func (g *opVertex) release(ctx context.Context) (retErr error) { return retErr } -func (g *opVertex) getInputRef(i int) cache.ImmutableRef { +func (g *opVertex) getInputRefForIndex(i int) cache.ImmutableRef { input := g.op.Inputs[i] for _, v := range g.inputs { if v.dgst == digest.Digest(input.Digest) { @@ -143,86 +143,101 @@ func (g *opVertex) solve(ctx context.Context, opt Opt) (retErr error) { } } - g.notifyStarted(pw) - defer g.notifyComplete(pw) + g.notifyStarted(ctx) + defer g.notifyCompleted(ctx) switch op := g.op.Op.(type) { case *pb.Op_Source: - id, err := source.FromString(op.Source.Identifier) - if err != nil { + if err := g.runSourceOp(ctx, opt.SourceManager, op); err != nil { return err } - ref, err := opt.SourceManager.Pull(ctx, id) - if err != nil { - return err - } - g.refs = []cache.ImmutableRef{ref} case *pb.Op_Exec: - - mounts := make(map[string]cache.Mountable) - - var outputs []cache.MutableRef - - defer func() { - for _, o := range outputs { - if o != nil { - s, err := o.Freeze() // TODO: log error - if err == nil { - s.Release(ctx) - } - } - } - }() - - for _, m := range op.Exec.Mounts { - var mountable cache.Mountable - ref := g.getInputRef(int(m.Input)) - mountable = ref - if m.Output != -1 { - active, err := opt.CacheManager.New(ctx, ref) // TODO: should be method - if err != nil { - return err - } - outputs = append(outputs, active) - mountable = active - } - mounts[m.Dest] = mountable + if err := g.runExecOp(ctx, opt.CacheManager, opt.Worker, op); err != nil { + return err } - - meta := worker.Meta{ - Args: op.Exec.Meta.Args, - Env: op.Exec.Meta.Env, - Cwd: op.Exec.Meta.Cwd, - } - - if err := opt.Worker.Exec(ctx, meta, mounts, os.Stderr, os.Stderr); err != nil { - return errors.Wrapf(err, "worker failed running %v", meta.Args) - } - - g.refs = []cache.ImmutableRef{} - - for i, o := range outputs { - ref, err := o.ReleaseAndCommit(ctx) - if err != nil { - return errors.Wrapf(err, "error committing %s", ref.ID()) - } - g.refs = append(g.refs, ref) - outputs[i] = nil - } - default: return errors.Errorf("invalid op type") } return nil } -func (g *opVertex) notifyStarted(pw progress.Writer) { +func (g *opVertex) runSourceOp(ctx context.Context, sm *source.Manager, op *pb.Op_Source) error { + id, err := source.FromString(op.Source.Identifier) + if err != nil { + return err + } + ref, err := sm.Pull(ctx, id) + if err != nil { + return err + } + g.refs = []cache.ImmutableRef{ref} + return nil +} + +func (g *opVertex) runExecOp(ctx context.Context, cm cache.Manager, w worker.Worker, op *pb.Op_Exec) error { + mounts := make(map[string]cache.Mountable) + + var outputs []cache.MutableRef + + defer func() { + for _, o := range outputs { + if o != nil { + s, err := o.Freeze() // TODO: log error + if err == nil { + s.Release(ctx) + } + } + } + }() + + for _, m := range op.Exec.Mounts { + var mountable cache.Mountable + ref := g.getInputRefForIndex(int(m.Input)) + mountable = ref + if m.Output != -1 { + active, err := cm.New(ctx, ref) // TODO: should be method + if err != nil { + return err + } + outputs = append(outputs, active) + mountable = active + } + mounts[m.Dest] = mountable + } + + meta := worker.Meta{ + Args: op.Exec.Meta.Args, + Env: op.Exec.Meta.Env, + Cwd: op.Exec.Meta.Cwd, + } + + if err := w.Exec(ctx, meta, mounts, os.Stderr, os.Stderr); err != nil { + return errors.Wrapf(err, "worker failed running %v", meta.Args) + } + + g.refs = []cache.ImmutableRef{} + for i, o := range outputs { + ref, err := o.ReleaseAndCommit(ctx) + if err != nil { + return errors.Wrapf(err, "error committing %s", o.ID()) + } + g.refs = append(g.refs, ref) + outputs[i] = nil + } + return nil +} + +func (g *opVertex) notifyStarted(ctx context.Context) { + pw, _, _ := progress.FromContext(ctx) + defer pw.Close() now := time.Now() g.vtx.Started = &now pw.Write(g.dgst.String(), g.vtx) } -func (g *opVertex) notifyComplete(pw progress.Writer) { +func (g *opVertex) notifyCompleted(ctx context.Context) { + pw, _, _ := progress.FromContext(ctx) + defer pw.Close() now := time.Now() g.vtx.Completed = &now pw.Write(g.dgst.String(), g.vtx) diff --git a/util/flightcontrol/flightcontrol.go b/util/flightcontrol/flightcontrol.go index 7bfec5b36..83a120ff5 100644 --- a/util/flightcontrol/flightcontrol.go +++ b/util/flightcontrol/flightcontrol.go @@ -78,7 +78,7 @@ func newCall(fn func(ctx context.Context) (interface{}, error)) *call { c.ctx = ctx c.closeProgressWriter = closeProgressWriter - go c.progressState.run(pr) + go c.progressState.run(pr) // TODO: remove this, wrap writer instead return c }