diff options
Diffstat (limited to 'internal/pkg/pkg.go')
| -rw-r--r-- | internal/pkg/pkg.go | 96 |
1 files changed, 60 insertions, 36 deletions
diff --git a/internal/pkg/pkg.go b/internal/pkg/pkg.go index c8e36996..7f082d5a 100644 --- a/internal/pkg/pkg.go +++ b/internal/pkg/pkg.go @@ -413,9 +413,9 @@ type pendingArtifactDep struct { // Cache is a support layer that implementations of [Artifact] can use to store // cured [Artifact] data in a content addressed fashion. type Cache struct { - // Work for curing dependency [Artifact] is sent here and cured concurrently - // while subject to the cures limit. Invalid after the context is canceled. - cureDep chan<- *pendingArtifactDep + // Cures of any variant of [Artifact] sends to cures before entering the + // implementation and receives an equal amount of elements after. + cures chan struct{} // [context.WithCancel] over caller-supplied context, used by [Artifact] and // all dependency curing goroutines. @@ -1102,6 +1102,14 @@ func (c *Cache) Cure(a Artifact) ( checksum unique.Handle[Checksum], err error, ) { + select { + case <-c.ctx.Done(): + err = c.ctx.Err() + return + + default: + } + if c.threshold > 0 { var n uintptr for range Flood(a) { @@ -1114,7 +1122,7 @@ func (c *Cache) Cure(a Artifact) ( c.msg.Verbosef("visited %d artifacts", n) } - return c.cure(a) + return c.cure(a, true) } // CureError wraps a non-nil error returned attempting to cure an [Artifact]. @@ -1187,8 +1195,30 @@ func (e *DependencyCureError) Error() string { return buf.String() } +// enterCure must be called before entering an [Artifact] implementation. +func (c *Cache) enterCure(curesExempt bool) error { + if curesExempt { + return nil + } + + select { + case c.cures <- struct{}{}: + return nil + + case <-c.ctx.Done(): + return c.ctx.Err() + } +} + +// exitCure must be called after exiting an [Artifact] implementation. +func (c *Cache) exitCure(curesExempt bool) { + if !curesExempt { + <-c.cures + } +} + // cure implements Cure without checking the full dependency graph. -func (c *Cache) cure(a Artifact) ( +func (c *Cache) cure(a Artifact, curesExempt bool) ( pathname *check.Absolute, checksum unique.Handle[Checksum], err error, @@ -1294,7 +1324,11 @@ func (c *Cache) cure(a Artifact) ( } var data []byte + if err = c.enterCure(curesExempt); err != nil { + return + } data, err = f.Cure(c.ctx) + c.exitCure(curesExempt) if err != nil { return } @@ -1365,7 +1399,12 @@ func (c *Cache) cure(a Artifact) ( switch ca := a.(type) { case TrivialArtifact: defer t.destroy(&err) - if err = ca.Cure(&t); err != nil { + if err = c.enterCure(curesExempt); err != nil { + return + } + err = ca.Cure(&t) + c.exitCure(curesExempt) + if err != nil { return } break @@ -1381,14 +1420,7 @@ func (c *Cache) cure(a Artifact) ( var errsMu sync.Mutex for i, d := range deps { pending := pendingArtifactDep{d, &res[i], &errs, &errsMu, &wg} - select { - case c.cureDep <- &pending: - break - - case <-c.ctx.Done(): - err = c.ctx.Err() - return - } + go pending.cure(c) } wg.Wait() @@ -1401,7 +1433,12 @@ func (c *Cache) cure(a Artifact) ( } defer f.destroy(&err) - if err = ca.Cure(&f); err != nil { + if err = c.enterCure(curesExempt); err != nil { + return + } + err = ca.Cure(&f) + c.exitCure(curesExempt) + if err != nil { return } break @@ -1486,7 +1523,7 @@ func (pending *pendingArtifactDep) cure(c *Cache) { defer pending.Done() var err error - pending.resP.pathname, pending.resP.checksum, err = c.cure(pending.a) + pending.resP.pathname, pending.resP.checksum, err = c.cure(pending.a, false) if err == nil { return } @@ -1501,6 +1538,7 @@ func (c *Cache) Close() { c.closeOnce.Do(func() { c.cancel() c.wg.Wait() + close(c.cures) c.unlock() }) } @@ -1536,6 +1574,10 @@ func open( base *check.Absolute, lock bool, ) (*Cache, error) { + if cures < 1 { + cures = runtime.NumCPU() + } + for _, name := range []string{ dirIdentifier, dirChecksum, @@ -1548,6 +1590,8 @@ func open( } c := Cache{ + cures: make(chan struct{}, cures), + msg: msg, base: base, @@ -1556,8 +1600,6 @@ func open( identPending: make(map[unique.Handle[ID]]<-chan struct{}), } c.ctx, c.cancel = context.WithCancel(ctx) - cureDep := make(chan *pendingArtifactDep, cures) - c.cureDep = cureDep c.identPool.New = func() any { return new(extIdent) } if lock || !testing.Testing() { @@ -1572,23 +1614,5 @@ func open( c.unlock = func() {} } - if cures < 1 { - cures = runtime.NumCPU() - } - for i := 0; i < cures; i++ { - c.wg.Go(func() { - for { - select { - case <-c.ctx.Done(): - return - - case pending := <-cureDep: - pending.cure(&c) - break - } - } - }) - } - return &c, nil } |
