From 9fa72212974fc0fad43137db7ed14fa7b2b80c88 Mon Sep 17 00:00:00 2001 From: Mike Date: Sun, 17 Mar 2024 14:49:59 -0700 Subject: [PATCH 1/4] Remove edge code --- solver/jobs.go | 104 ++++++++++++++++++++++++++++--------------------- 1 file changed, 60 insertions(+), 44 deletions(-) diff --git a/solver/jobs.go b/solver/jobs.go index 53afa515d..747d3f830 100644 --- a/solver/jobs.go +++ b/solver/jobs.go @@ -529,33 +529,48 @@ func (jl *Solver) deleteIfUnreferenced(k digest.Digest, st *state) { } } -type dumbBuilder struct { +type simpleSolver struct { resolveOpFunc ResolveOpFunc solver *Solver job *Job } -var cache = map[string]CachedResult{} +var cache = map[string]Result{} +var inFlight = map[string]struct{}{} var mu = sync.Mutex{} -func (b *dumbBuilder) build(ctx context.Context, e Edge) (CachedResult, error) { +func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) { // Ordered list of vertices to build. - digests, vertices := b.exploreVertices(e) + digests, vertices := s.exploreVertices(e) - var ret CachedResult + var ret Result for _, d := range digests { - mu.Lock() vertex, ok := vertices[d] if !ok { return nil, errors.Errorf("digest %s not found", d) } + for { + mu.Lock() + _, ok := inFlight[d.String()] + mu.Unlock() + if ok { + time.Sleep(100 * time.Millisecond) + } else { + break + } + } + + mu.Lock() + inFlight[d.String()] = struct{}{} + mu.Unlock() + defaultCache := NewInMemoryCacheManager() st := &state{ - opts: SolverOpt{DefaultCache: defaultCache, ResolveOpFunc: b.resolveOpFunc}, + opts: SolverOpt{DefaultCache: defaultCache, ResolveOpFunc: s.resolveOpFunc}, parents: map[digest.Digest]struct{}{}, childVtx: map[digest.Digest]struct{}{}, allPw: map[progress.Writer]struct{}{}, @@ -564,71 +579,72 @@ func (b *dumbBuilder) build(ctx context.Context, e Edge) (CachedResult, error) { vtx: vertex, clientVertex: initClientVertex(vertex), edges: map[Index]*edge{}, - index: b.solver.index, + index: s.solver.index, mainCache: defaultCache, cache: map[string]CacheManager{}, - solver: b.solver, + solver: s.solver, origDigest: vertex.Digest(), } st.jobs = map[*Job]struct{}{ - b.job: {}, + s.job: {}, } - st.mpw.Add(b.job.pw) + st.mpw.Add(s.job.pw) - //fmt.Println("Processing vertex", vertex.Name(), d.String()) + notifyCompleted := notifyStarted(ctx, &st.clientVertex, true) - edge := st.getEdge(e.Index) - - edge.deps = make([]*dep, 0, len(vertex.Inputs())) - inputs := vertex.Inputs() - for i := range inputs { - dep := newDep(Index(i)) - if v, ok := cache[inputs[i].Vertex.Digest().String()]; ok { - dep.result = NewSharedCachedResult(v) - } - edge.deps = append(edge.deps, dep) + mu.Lock() + if v, ok := cache[d.String()]; ok { + delete(inFlight, d.String()) + mu.Unlock() + notifyCompleted(nil, true) + ret = v + continue } + mu.Unlock() - cm, err := edge.op.CacheMap(ctx, int(e.Index)) + op := newSharedOp(st.opts.ResolveOpFunc, st.opts.DefaultCache, st) + + // CacheMap populates required fields in SourceOp. + _, err := op.CacheMap(ctx, int(e.Index)) if err != nil { return nil, err } - edge.cacheMap = cm.CacheMap + var inDigests []string + for _, in := range vertex.Inputs() { + inDigests = append(inDigests, in.Vertex.Digest().String()) + } - if r, ok := cache[d.String()]; ok { - st.clientVertex.Cached = true - st.mpw.Write(identity.NewID(), st.clientVertex) - ret = r - mu.Unlock() - continue + var inputs []Result + for _, in := range vertex.Inputs() { + if v, ok := cache[in.Vertex.Digest().String()]; ok { + inputs = append(inputs, v) + } } - res, err := edge.execOp(ctx) + results, _, err := op.Exec(ctx, inputs) if err != nil { + notifyCompleted(err, false) return nil, err } - cachedResult := res.(CachedResult) - - st.mpw.Write(identity.NewID(), st.clientVertex) + notifyCompleted(nil, false) - edge.result = NewSharedCachedResult(cachedResult) - edge.state = edgeStatusComplete + res := results[int(e.Index)] + ret = res - ret = cachedResult - - fmt.Println("Adding cache!", d.String()) - cache[d.String()] = cachedResult + mu.Lock() + delete(inFlight, d.String()) + cache[d.String()] = res mu.Unlock() } - return ret, nil + return NewCachedResult(ret, []ExportableCacheKey{}), nil } -func (b *dumbBuilder) exploreVertices(e Edge) ([]digest.Digest, map[digest.Digest]Vertex) { +func (s *simpleSolver) exploreVertices(e Edge) ([]digest.Digest, map[digest.Digest]Vertex) { digests := []digest.Digest{e.Vertex.Digest()} vertices := map[digest.Digest]Vertex{ @@ -636,7 +652,7 @@ func (b *dumbBuilder) exploreVertices(e Edge) ([]digest.Digest, map[digest.Diges } for _, edge := range e.Vertex.Inputs() { - d, v := b.exploreVertices(edge) + d, v := s.exploreVertices(edge) digests = append(d, digests...) for key, value := range v { vertices[key] = value @@ -660,7 +676,7 @@ func (j *Job) Build(ctx context.Context, e Edge) (CachedResultWithProvenance, er j.span = span } - b := &dumbBuilder{ + b := &simpleSolver{ resolveOpFunc: j.list.opts.ResolveOpFunc, solver: j.list, job: j, From 8a63a039abf96ad6e425cddbd7d55337c4acdf03 Mon Sep 17 00:00:00 2001 From: Mike Date: Mon, 18 Mar 2024 09:27:39 -0700 Subject: [PATCH 2/4] Cache --- solver/jobs.go | 48 +++++++++++++++++++++++++++++++++--------------- 1 file changed, 33 insertions(+), 15 deletions(-) diff --git a/solver/jobs.go b/solver/jobs.go index 747d3f830..14841b503 100644 --- a/solver/jobs.go +++ b/solver/jobs.go @@ -536,6 +536,8 @@ type simpleSolver struct { } var cache = map[string]Result{} +var cacheDigests = map[string]string{} + var inFlight = map[string]struct{}{} var mu = sync.Mutex{} @@ -594,32 +596,48 @@ func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) notifyCompleted := notifyStarted(ctx, &st.clientVertex, true) - mu.Lock() - if v, ok := cache[d.String()]; ok { - delete(inFlight, d.String()) - mu.Unlock() - notifyCompleted(nil, true) - ret = v - continue - } - mu.Unlock() - op := newSharedOp(st.opts.ResolveOpFunc, st.opts.DefaultCache, st) // CacheMap populates required fields in SourceOp. - _, err := op.CacheMap(ctx, int(e.Index)) + cm, err := op.CacheMap(ctx, int(e.Index)) if err != nil { return nil, err } - var inDigests []string + fmt.Println(vertex.Name()) + fmt.Println("CacheMap Digest:", cm.Digest) + fmt.Println("LLB Digest:", d) + for _, in := range vertex.Inputs() { - inDigests = append(inDigests, in.Vertex.Digest().String()) + fmt.Println("> input name", in.Vertex.Name()) + fmt.Println("> input digest", in.Vertex.Digest()) + } + + for _, in := range cm.Deps { + fmt.Println("> CacheMap selector:", in.Selector) } + fmt.Println("----") + + cacheDigest := cm.Digest.String() + mu.Lock() + cacheDigests[d.String()] = cacheDigest + mu.Unlock() + + mu.Lock() + if v, ok := cache[cacheDigest]; ok { + delete(inFlight, d.String()) + mu.Unlock() + notifyCompleted(nil, true) + ret = v + continue + } + mu.Unlock() + var inputs []Result for _, in := range vertex.Inputs() { - if v, ok := cache[in.Vertex.Digest().String()]; ok { + key := cacheDigests[in.Vertex.Digest().String()] + if v, ok := cache[key]; ok { inputs = append(inputs, v) } } @@ -637,7 +655,7 @@ func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) mu.Lock() delete(inFlight, d.String()) - cache[d.String()] = res + cache[cacheDigest] = res mu.Unlock() } From 676524c66a37508d3dffd6ee5305766ee0920b27 Mon Sep 17 00:00:00 2001 From: Mike Date: Tue, 19 Mar 2024 10:21:53 -0700 Subject: [PATCH 3/4] Generate consistent cache keys --- solver/jobs.go | 137 +++++++++++++++++++++++++++++++++++++++---------- 1 file changed, 111 insertions(+), 26 deletions(-) diff --git a/solver/jobs.go b/solver/jobs.go index 14841b503..ddc37961b 100644 --- a/solver/jobs.go +++ b/solver/jobs.go @@ -2,10 +2,14 @@ package solver import ( "context" + "crypto/sha256" "fmt" + "hash" + "io" "sync" "time" + "github.com/davecgh/go-spew/spew" "github.com/moby/buildkit/client" "github.com/moby/buildkit/identity" "github.com/moby/buildkit/session" @@ -536,11 +540,21 @@ type simpleSolver struct { } var cache = map[string]Result{} -var cacheDigests = map[string]string{} - +var cacheMaps = map[string]*simpleCacheMap{} var inFlight = map[string]struct{}{} var mu = sync.Mutex{} +type cacheMapDep struct { + selector string + computed string +} + +type simpleCacheMap struct { + digest string + inputs []string + deps []cacheMapDep +} + func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) { // Ordered list of vertices to build. @@ -605,27 +619,61 @@ func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) } fmt.Println(vertex.Name()) - fmt.Println("CacheMap Digest:", cm.Digest) - fmt.Println("LLB Digest:", d) - - for _, in := range vertex.Inputs() { - fmt.Println("> input name", in.Vertex.Name()) - fmt.Println("> input digest", in.Vertex.Digest()) + fmt.Println("LLB digest:", d) + fmt.Println("CacheMap digest:", cm.Digest) + + // TODO: handle cm.Opts (CacheOpts) + scm := simpleCacheMap{ + digest: cm.Digest.String(), + deps: make([]cacheMapDep, len(cm.Deps)), + inputs: make([]string, len(cm.Deps)), } - - for _, in := range cm.Deps { - fmt.Println("> CacheMap selector:", in.Selector) + var inputs []Result + for i, in := range vertex.Inputs() { + cacheKey, err := s.cacheKey(ctx, in.Vertex.Digest().String()) + if err != nil { + return nil, err + } + mu.Lock() + res, ok := cache[cacheKey] + mu.Unlock() + if !ok { + return nil, errors.Errorf("cache key not found: %s", cacheKey) + } + dep := cm.Deps[i] + spew.Dump(dep) + if dep.PreprocessFunc != nil { + err = dep.PreprocessFunc(ctx, res, st) + if err != nil { + return nil, err + } + } + scm.deps[i] = cacheMapDep{ + selector: dep.Selector.String(), + } + if dep.ComputeDigestFunc != nil { + compDigest, err := dep.ComputeDigestFunc(ctx, res, st) + if err != nil { + return nil, err + } + scm.deps[i].computed = compDigest.String() + } + scm.inputs[i] = in.Vertex.Digest().String() + inputs = append(inputs, res) } - - fmt.Println("----") - - cacheDigest := cm.Digest.String() mu.Lock() - cacheDigests[d.String()] = cacheDigest + cacheMaps[d.String()] = &scm mu.Unlock() + cacheKey, err := s.cacheKey(ctx, d.String()) + if err != nil { + return nil, err + } + + fmt.Println("Cache key", cacheKey) + mu.Lock() - if v, ok := cache[cacheDigest]; ok { + if v, ok := cache[cacheKey]; ok && v != nil { delete(inFlight, d.String()) mu.Unlock() notifyCompleted(nil, true) @@ -634,14 +682,6 @@ func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) } mu.Unlock() - var inputs []Result - for _, in := range vertex.Inputs() { - key := cacheDigests[in.Vertex.Digest().String()] - if v, ok := cache[key]; ok { - inputs = append(inputs, v) - } - } - results, _, err := op.Exec(ctx, inputs) if err != nil { notifyCompleted(err, false) @@ -650,18 +690,63 @@ func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) notifyCompleted(nil, false) + fmt.Println("Result length:", len(results)) + res := results[int(e.Index)] ret = res mu.Lock() delete(inFlight, d.String()) - cache[cacheDigest] = res + cache[cacheKey] = res mu.Unlock() } return NewCachedResult(ret, []ExportableCacheKey{}), nil } +func (s *simpleSolver) cacheKey(ctx context.Context, d string) (string, error) { + h := sha256.New() + + fmt.Println("Computing cache key") + + err := s.calcCacheKey(ctx, d, h) + if err != nil { + return "", err + } + + return fmt.Sprintf("%x", h.Sum(nil)), nil +} + +func (s *simpleSolver) calcCacheKey(ctx context.Context, d string, h hash.Hash) error { + mu.Lock() + c, ok := cacheMaps[d] + mu.Unlock() + if !ok { + return errors.New("missing cache map key") + } + + spew.Dump(c) + + for _, in := range c.inputs { + err := s.calcCacheKey(ctx, in, h) + if err != nil { + return err + } + } + + io.WriteString(h, c.digest) + for _, dep := range c.deps { + if dep.selector != "" { + io.WriteString(h, dep.selector) + } + if dep.computed != "" { + io.WriteString(h, dep.computed) + } + } + + return nil +} + func (s *simpleSolver) exploreVertices(e Edge) ([]digest.Digest, map[digest.Digest]Vertex) { digests := []digest.Digest{e.Vertex.Digest()} From 392d0ad7142c40b04a17b7fd1cedb5047cc8f2d2 Mon Sep 17 00:00:00 2001 From: Mike Date: Tue, 19 Mar 2024 11:16:32 -0700 Subject: [PATCH 4/4] Refactoring --- solver/jobs.go | 207 ++++++++++++++++++++++++++++--------------------- 1 file changed, 120 insertions(+), 87 deletions(-) diff --git a/solver/jobs.go b/solver/jobs.go index ddc37961b..d6f53f0af 100644 --- a/solver/jobs.go +++ b/solver/jobs.go @@ -9,7 +9,6 @@ import ( "sync" "time" - "github.com/davecgh/go-spew/spew" "github.com/moby/buildkit/client" "github.com/moby/buildkit/identity" "github.com/moby/buildkit/session" @@ -541,7 +540,7 @@ type simpleSolver struct { var cache = map[string]Result{} var cacheMaps = map[string]*simpleCacheMap{} -var inFlight = map[string]struct{}{} +var execInProgress = map[string]struct{}{} var mu = sync.Mutex{} type cacheMapDep struct { @@ -563,50 +562,27 @@ func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) var ret Result for _, d := range digests { + fmt.Println() + vertex, ok := vertices[d] if !ok { return nil, errors.Errorf("digest %s not found", d) } + mu.Lock() + // TODO: replace busy-wait loop with a wait-for-channel-to-close approach for { - mu.Lock() - _, ok := inFlight[d.String()] - mu.Unlock() - if ok { - time.Sleep(100 * time.Millisecond) - } else { + if _, shouldWait := execInProgress[d.String()]; !shouldWait { + execInProgress[d.String()] = struct{}{} + mu.Unlock() break } + mu.Unlock() + time.Sleep(time.Millisecond * 10) + mu.Lock() } - mu.Lock() - inFlight[d.String()] = struct{}{} - mu.Unlock() - - defaultCache := NewInMemoryCacheManager() - - st := &state{ - opts: SolverOpt{DefaultCache: defaultCache, ResolveOpFunc: s.resolveOpFunc}, - parents: map[digest.Digest]struct{}{}, - childVtx: map[digest.Digest]struct{}{}, - allPw: map[progress.Writer]struct{}{}, - mpw: progress.NewMultiWriter(progress.WithMetadata("vertex", d)), - mspan: tracing.NewMultiSpan(), - vtx: vertex, - clientVertex: initClientVertex(vertex), - edges: map[Index]*edge{}, - index: s.solver.index, - mainCache: defaultCache, - cache: map[string]CacheManager{}, - solver: s.solver, - origDigest: vertex.Digest(), - } - - st.jobs = map[*Job]struct{}{ - s.job: {}, - } - - st.mpw.Add(s.job.pw) + st := s.createState(vertex) notifyCompleted := notifyStarted(ctx, &st.clientVertex, true) @@ -619,62 +595,25 @@ func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) } fmt.Println(vertex.Name()) - fmt.Println("LLB digest:", d) + fmt.Println("LLB digest:", d.String()) fmt.Println("CacheMap digest:", cm.Digest) - // TODO: handle cm.Opts (CacheOpts) - scm := simpleCacheMap{ - digest: cm.Digest.String(), - deps: make([]cacheMapDep, len(cm.Deps)), - inputs: make([]string, len(cm.Deps)), - } - var inputs []Result - for i, in := range vertex.Inputs() { - cacheKey, err := s.cacheKey(ctx, in.Vertex.Digest().String()) - if err != nil { - return nil, err - } - mu.Lock() - res, ok := cache[cacheKey] - mu.Unlock() - if !ok { - return nil, errors.Errorf("cache key not found: %s", cacheKey) - } - dep := cm.Deps[i] - spew.Dump(dep) - if dep.PreprocessFunc != nil { - err = dep.PreprocessFunc(ctx, res, st) - if err != nil { - return nil, err - } - } - scm.deps[i] = cacheMapDep{ - selector: dep.Selector.String(), - } - if dep.ComputeDigestFunc != nil { - compDigest, err := dep.ComputeDigestFunc(ctx, res, st) - if err != nil { - return nil, err - } - scm.deps[i].computed = compDigest.String() - } - scm.inputs[i] = in.Vertex.Digest().String() - inputs = append(inputs, res) + inputs, err := s.preprocessInputs(ctx, st, vertex, cm.CacheMap) + if err != nil { + return nil, err } - mu.Lock() - cacheMaps[d.String()] = &scm - mu.Unlock() cacheKey, err := s.cacheKey(ctx, d.String()) if err != nil { return nil, err } - fmt.Println("Cache key", cacheKey) + fmt.Println("Computed cache key:", cacheKey) mu.Lock() if v, ok := cache[cacheKey]; ok && v != nil { - delete(inFlight, d.String()) + fmt.Println("Cache hit!") + delete(execInProgress, d.String()) mu.Unlock() notifyCompleted(nil, true) ret = v @@ -684,19 +623,20 @@ func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) results, _, err := op.Exec(ctx, inputs) if err != nil { + mu.Lock() + delete(execInProgress, d.String()) + mu.Unlock() notifyCompleted(err, false) return nil, err } notifyCompleted(nil, false) - fmt.Println("Result length:", len(results)) - res := results[int(e.Index)] ret = res mu.Lock() - delete(inFlight, d.String()) + delete(execInProgress, d.String()) cache[cacheKey] = res mu.Unlock() } @@ -704,11 +644,106 @@ func (s *simpleSolver) build(ctx context.Context, e Edge) (CachedResult, error) return NewCachedResult(ret, []ExportableCacheKey{}), nil } +func (s *simpleSolver) createState(vertex Vertex) *state { + defaultCache := NewInMemoryCacheManager() + + st := &state{ + opts: SolverOpt{DefaultCache: defaultCache, ResolveOpFunc: s.resolveOpFunc}, + parents: map[digest.Digest]struct{}{}, + childVtx: map[digest.Digest]struct{}{}, + allPw: map[progress.Writer]struct{}{}, + mpw: progress.NewMultiWriter(progress.WithMetadata("vertex", vertex.Digest())), + mspan: tracing.NewMultiSpan(), + vtx: vertex, + clientVertex: initClientVertex(vertex), + edges: map[Index]*edge{}, + index: s.solver.index, + mainCache: defaultCache, + cache: map[string]CacheManager{}, + solver: s.solver, + origDigest: vertex.Digest(), + } + + st.jobs = map[*Job]struct{}{ + s.job: {}, + } + + st.mpw.Add(s.job.pw) + + return st +} + +func (s *simpleSolver) preprocessInputs(ctx context.Context, st *state, vertex Vertex, cm *CacheMap) ([]Result, error) { + // This struct is used to reconstruct a cache key from an LLB digest & all + // parents using consistent digests that depend on the full dependency chain. + // TODO: handle cm.Opts (CacheOpts)? + scm := simpleCacheMap{ + digest: cm.Digest.String(), + deps: make([]cacheMapDep, len(cm.Deps)), + inputs: make([]string, len(cm.Deps)), + } + + var inputs []Result + + for i, in := range vertex.Inputs() { + // Compute a cache key given the LLB digest value. + cacheKey, err := s.cacheKey(ctx, in.Vertex.Digest().String()) + if err != nil { + return nil, err + } + + // Lookup the result for that cache key. + mu.Lock() + res, ok := cache[cacheKey] + mu.Unlock() + if !ok { + return nil, errors.Errorf("cache key not found: %s", cacheKey) + } + + dep := cm.Deps[i] + + // Unlazy the result. + if dep.PreprocessFunc != nil { + err = dep.PreprocessFunc(ctx, res, st) + if err != nil { + return nil, err + } + } + + // Add selectors (usually file references) to the struct. + scm.deps[i] = cacheMapDep{ + selector: dep.Selector.String(), + } + + // ComputeDigestFunc will usually checksum files. This is then used as + // part of the cache key to ensure it's consistent & distinct for this + // operation. + if dep.ComputeDigestFunc != nil { + compDigest, err := dep.ComputeDigestFunc(ctx, res, st) + if err != nil { + return nil, err + } + scm.deps[i].computed = compDigest.String() + } + + // Add input references to the struct as to link dependencies. + scm.inputs[i] = in.Vertex.Digest().String() + + // Add the cached result to the input set. These inputs are used to + // reconstruct dependencies (mounts, etc.) for a new container run. + inputs = append(inputs, res) + } + + mu.Lock() + cacheMaps[vertex.Digest().String()] = &scm + mu.Unlock() + + return inputs, nil +} + func (s *simpleSolver) cacheKey(ctx context.Context, d string) (string, error) { h := sha256.New() - fmt.Println("Computing cache key") - err := s.calcCacheKey(ctx, d, h) if err != nil { return "", err @@ -725,8 +760,6 @@ func (s *simpleSolver) calcCacheKey(ctx context.Context, d string, h hash.Hash) return errors.New("missing cache map key") } - spew.Dump(c) - for _, in := range c.inputs { err := s.calcCacheKey(ctx, in, h) if err != nil {