diff --git a/internal/imagepush/mirror.go b/internal/imagepush/mirror.go index 2df2078..253a1d1 100644 --- a/internal/imagepush/mirror.go +++ b/internal/imagepush/mirror.go @@ -109,11 +109,15 @@ type Source struct { tokens map[string]string // host/repo → bearer token } +// sourceClient is shared by every Source without its own Client, so reads +// reuse keep-alive connections instead of opening one per request. +var sourceClient = &http.Client{Transport: &http.Transport{Proxy: http.ProxyFromEnvironment, ResponseHeaderTimeout: sourceRequestTimeout}} + func (s *Source) client() *http.Client { if s.Client != nil { return s.Client } - return &http.Client{Transport: &http.Transport{Proxy: http.ProxyFromEnvironment, ResponseHeaderTimeout: sourceRequestTimeout}} + return sourceClient } func (s *Source) url(r SourceRef, tail string) string { diff --git a/internal/imagepush/push.go b/internal/imagepush/push.go index 87cee52..f5e8916 100644 --- a/internal/imagepush/push.go +++ b/internal/imagepush/push.go @@ -240,10 +240,11 @@ func (p *Pusher) uploadBlob(ctx context.Context, r Ref, digest string, size int6 if err != nil { return err } - start.Body.Close() if start.StatusCode != http.StatusAccepted { + defer start.Body.Close() return statusError("start upload", start) } + drainClose(start) loc, err := start.Location() if err != nil { return fmt.Errorf("start upload: %w", err) @@ -261,10 +262,11 @@ func (p *Pusher) uploadBlob(ctx context.Context, r Ref, digest string, size int6 if err != nil { return err } - put.Body.Close() if put.StatusCode != http.StatusCreated { + defer put.Body.Close() return statusError("upload", put) } + drainClose(put) p.logf("pushed %s (%d bytes)", digest, size) return nil } @@ -311,13 +313,17 @@ func (p *Pusher) do(ctx context.Context, method, target string, body io.Reader, } c := p.Client if c == nil { - // No overall timeout: a modpack layer can take minutes on a slow disk, and - // the Job's activeDeadlineSeconds is the real bound. - c = &http.Client{Transport: &http.Transport{ResponseHeaderTimeout: 2 * time.Minute}} + c = pushClient } return c.Do(req) } +// pushClient is shared by every Pusher without its own Client, so a restore of +// many images reuses keep-alive connections instead of opening one per +// request. No overall timeout: a modpack layer can take minutes on a slow disk, +// and the caller's deadline is the real bound. +var pushClient = &http.Client{Transport: &http.Transport{ResponseHeaderTimeout: 2 * time.Minute}} + // retry runs fn up to Attempts times, backing off between tries. A refusal the // registry will repeat (401/403/4xx other than 408/429) is returned at once. A 503 // with Retry-After is waited out without using up an attempt, for as long as @@ -394,6 +400,13 @@ func (e *StatusError) Error() string { return fmt.Sprintf("%s: registry answered %d: %s", e.Op, e.Code, e.Body) } +// drainClose reads what is left of a small response body so its connection +// goes back to the pool. +func drainClose(resp *http.Response) { + io.Copy(io.Discard, io.LimitReader(resp.Body, 64<<10)) + resp.Body.Close() +} + func statusError(op string, resp *http.Response) error { b, _ := io.ReadAll(io.LimitReader(resp.Body, 2048)) se := &StatusError{Op: op, Code: resp.StatusCode, Body: strings.TrimSpace(string(b))}