diff --git a/server/download.go b/server/download.go index 5e3217f5..a50c0cd1 100644 --- a/server/download.go +++ b/server/download.go @@ -151,36 +151,9 @@ func (b *blobDownload) Run(ctx context.Context, requestURL *url.URL, opts *regis _ = file.Truncate(b.Total) - g, inner := NewLimitGroup(ctx, numDownloadParts) - - go func() { - ticker := time.NewTicker(time.Second) - var n int64 = 1 - var maxDelta float64 - var buckets []int64 - for { - select { - case <-ticker.C: - buckets = append(buckets, b.Completed.Load()) - if len(buckets) < 2 { - continue - } else if len(buckets) > 10 { - buckets = buckets[1:] - } - - delta := float64((buckets[len(buckets)-1] - buckets[0])) / float64(len(buckets)) - slog.Debug(fmt.Sprintf("delta: %s/s max_delta: %s/s", format.HumanBytes(int64(delta)), format.HumanBytes(int64(maxDelta)))) - if delta > maxDelta*1.5 { - maxDelta = delta - g.SetLimit(n) - n++ - } - - case <-ctx.Done(): - return - } - } - }() + var limit int64 = 2 + g, inner := NewLimitGroup(ctx, numDownloadParts, limit) + go watchDelta(inner, g, &b.Completed, limit) for i := range b.Parts { part := b.Parts[i] @@ -410,37 +383,84 @@ func downloadBlob(ctx context.Context, opts downloadOpts) error { type LimitGroup struct { *errgroup.Group - Semaphore *semaphore.Weighted - - weight, max_weight int64 + *semaphore.Weighted + size, limit int64 } -func NewLimitGroup(ctx context.Context, n int64) (*LimitGroup, context.Context) { +func NewLimitGroup(ctx context.Context, size, limit int64) (*LimitGroup, context.Context) { g, ctx := errgroup.WithContext(ctx) return &LimitGroup{ - Group: g, - Semaphore: semaphore.NewWeighted(n), - weight: n, - max_weight: n, + Group: g, + Weighted: semaphore.NewWeighted(size), + size: size, + limit: limit, }, ctx } func (g *LimitGroup) Go(ctx context.Context, fn func() error) { - weight := g.weight - _ = g.Semaphore.Acquire(ctx, weight) + var weight int64 = 1 + if g.limit > 0 { + weight = g.size / g.limit + } + + _ = g.Acquire(ctx, weight) if ctx.Err() != nil { return } g.Group.Go(func() error { - defer g.Semaphore.Release(weight) + defer g.Release(weight) return fn() }) } -func (g *LimitGroup) SetLimit(n int64) { - if n > 0 { - slog.Debug(fmt.Sprintf("setting limit to %d", n)) - g.weight = g.max_weight / n +func (g *LimitGroup) SetLimit(limit int64) { + if limit > g.limit { + g.limit = limit + } +} + +func watchDelta(ctx context.Context, g *LimitGroup, c *atomic.Int64, limit int64) { + var maxDelta float64 + var buckets []int64 + + // 5s ramp up period + nextUpdate := time.Now().Add(5 * time.Second) + + ticker := time.NewTicker(time.Second) + for { + select { + case <-ticker.C: + buckets = append(buckets, c.Load()) + if len(buckets) < 2 { + continue + } else if len(buckets) > 10 { + buckets = buckets[1:] + } + + delta := float64((buckets[len(buckets)-1] - buckets[0])) / float64(len(buckets)) + slog.Debug("", "limit", limit, "delta", format.HumanBytes(int64(delta)), "max_delta", format.HumanBytes(int64(maxDelta))) + + if time.Now().Before(nextUpdate) { + // quiet period; do not update ccy if recently updated + continue + } else if maxDelta > 0 { + x := delta / maxDelta + if x < 1.2 { + continue + } + + limit += int64(x) + slog.Debug("setting", "limit", limit) + g.SetLimit(limit) + } + + // 3s cooldown period + nextUpdate = time.Now().Add(3 * time.Second) + maxDelta = delta + + case <-ctx.Done(): + return + } } } diff --git a/server/upload.go b/server/upload.go index 590247d9..c4517268 100644 --- a/server/upload.go +++ b/server/upload.go @@ -136,41 +136,15 @@ func (b *blobUpload) Run(ctx context.Context, opts *registryOptions) { } defer b.file.Close() - g, inner := NewLimitGroup(ctx, numUploadParts) - - go func() { - ticker := time.NewTicker(time.Second) - var n int64 = 1 - var maxDelta float64 - var buckets []int64 - for { - select { - case <-ticker.C: - buckets = append(buckets, b.Completed.Load()) - if len(buckets) < 2 { - continue - } else if len(buckets) > 10 { - buckets = buckets[1:] - } - - delta := float64((buckets[len(buckets)-1] - buckets[0])) / float64(len(buckets)) - slog.Debug(fmt.Sprintf("delta: %s/s max_delta: %s/s", format.HumanBytes(int64(delta)), format.HumanBytes(int64(maxDelta)))) - if delta > maxDelta*1.5 { - maxDelta = delta - g.SetLimit(n) - n++ - } - - case <-ctx.Done(): - return - } - } - }() + var limit int64 = 2 + g, inner := NewLimitGroup(ctx, numUploadParts, limit) + go watchDelta(inner, g, &b.Completed, limit) for i := range b.Parts { part := &b.Parts[i] select { case <-inner.Done(): + break case requestURL := <-b.nextURL: g.Go(inner, func() error { var err error