Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions core/remotes/docker/fetcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -231,6 +231,13 @@ func (r dockerFetcher) Fetch(ctx context.Context, desc ocispec.Descriptor) (io.R
return nil, err
}

if r.warningHandler != nil {
ctx = context.WithValue(ctx, warningSourceKey{}, WarningSource{
Desc: &desc,
Digest: &desc.Digest,
})
}

return newHTTPReadSeeker(desc.Size, func(offset int64) (io.ReadCloser, error) {
// firstly try fetch via external urls
for _, us := range desc.URLs {
Expand All @@ -243,6 +250,7 @@ func (r dockerFetcher) Fetch(ctx context.Context, desc ocispec.Descriptor) (io.R
log.G(ctx).Debug("non-http(s) alternative url is unsupported")
continue
}

ctx = log.WithLogger(ctx, log.G(ctx).WithField("url", u))
log.G(ctx).Info("request")

Expand Down Expand Up @@ -379,6 +387,12 @@ func (r dockerFetcher) FetchByDigest(ctx context.Context, dgst digest.Digest, op
return nil, desc, err
}

if r.warningHandler != nil {
ctx = context.WithValue(ctx, warningSourceKey{}, WarningSource{
Digest: &dgst,
})
}

var (
getReq *request
sz int64
Expand Down
6 changes: 6 additions & 0 deletions core/remotes/docker/pusher.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,12 @@ func (p dockerPusher) push(ctx context.Context, desc ocispec.Descriptor, ref str
if err != nil {
return nil, err
}
if p.dockerBase.warningHandler != nil {
ctx = context.WithValue(ctx, warningSourceKey{}, WarningSource{
Desc: &desc,
Digest: &desc.Digest,
})
}
status, err := p.tracker.GetStatus(ref)
if err == nil {
if status.Committed && status.Offset == status.Total {
Expand Down
104 changes: 73 additions & 31 deletions core/remotes/docker/resolver.go
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,13 @@ type ResolverOptions struct {
// mechanism for getting blob upload status is expensive.
Tracker StatusTracker

// WarningHandler is called for each warning received from the registry.
// Warnings are reported via HTTP Warning headers with warn-code 299.
// It may be called concurrently from multiple goroutines, so
// implementations must be safe for concurrent use.
// If nil, warnings are ignored.
WarningHandler WarningHandler

// Authorizer is used to authorize registry requests
//
// Deprecated: use Hosts.
Expand Down Expand Up @@ -139,11 +146,12 @@ func DefaultHost(ns string) (string, error) {
}

type dockerResolver struct {
hosts RegistryHosts
header http.Header
resolveHeader http.Header
tracker StatusTracker
config transfer.ImageResolverOptions
hosts RegistryHosts
header http.Header
resolveHeader http.Header
tracker StatusTracker
config transfer.ImageResolverOptions
warningHandler WarningHandler
}

// NewResolver returns a new resolver to a Docker registry
Expand Down Expand Up @@ -198,10 +206,11 @@ func NewResolver(options ResolverOptions) remotes.Resolver {
options.Hosts = ConfigureDefaultRegistries(opts...)
}
return &dockerResolver{
hosts: options.Hosts,
header: options.Headers,
resolveHeader: resolveHeader,
tracker: options.Tracker,
hosts: options.Hosts,
header: options.Headers,
resolveHeader: resolveHeader,
tracker: options.Tracker,
warningHandler: options.WarningHandler,
}
}

Expand Down Expand Up @@ -316,6 +325,10 @@ func (r *dockerResolver) Resolve(ctx context.Context, ref string) (string, ocisp
"method": req.method,
"url": req.sanitizedURL(),
}))

// Don't report warnings during resolution; will report after descriptor construction
req.warningHandler = nil

log.G(ctx).Debug("resolving")
resp, err := req.doWithRetries(ctx, i == len(hosts)-1)
if err != nil {
Expand All @@ -329,6 +342,10 @@ func (r *dockerResolver) Resolve(ctx context.Context, ref string) (string, ocisp
log.G(ctx).WithError(err).Info(nextHostOrFail(i))
continue // try another host
}
var respHeaders http.Header
if r.warningHandler != nil {
respHeaders = resp.Header.Clone()
}
resp.Body.Close() // don't care about body contents.

if resp.StatusCode > 299 {
Expand Down Expand Up @@ -398,6 +415,9 @@ func (r *dockerResolver) Resolve(ctx context.Context, ref string) (string, ocisp
req.header[key] = append(req.header[key], value...)
}

// Don't report warnings during resolution
req.warningHandler = nil

resp, err := req.doWithRetries(ctx, true)
if err != nil {
return "", ocispec.Descriptor{}, err
Expand All @@ -409,6 +429,11 @@ func (r *dockerResolver) Resolve(ctx context.Context, ref string) (string, ocisp
return "", ocispec.Descriptor{}, unexpectedResponseErr(resp)
}

// Save response headers for warning reporting later
if r.warningHandler != nil {
respHeaders = resp.Header.Clone()
}

bodyReader := countingReader{reader: resp.Body}

contentType = getManifestMediaType(resp)
Expand Down Expand Up @@ -447,6 +472,13 @@ func (r *dockerResolver) Resolve(ctx context.Context, ref string) (string, ocisp
Size: size,
}

// Report any warnings with the warning source
reportWarningsWithSource(ctx, respHeaders, r.warningHandler, WarningSource{
Ref: refspec,
Desc: &desc,
Digest: &dgst,
})

log.G(ctx).WithField("desc.digest", desc.Digest).Debug("resolved")
return ref, desc, nil
}
Expand Down Expand Up @@ -501,12 +533,13 @@ func (r *dockerResolver) resolveDockerBase(ref string) (*dockerBase, error) {
}

type dockerBase struct {
refspec reference.Spec
repository string
hosts []RegistryHost
header http.Header
performances transfer.ImageResolverPerformanceSettings
limiter *semaphore.Weighted
refspec reference.Spec
repository string
hosts []RegistryHost
header http.Header
performances transfer.ImageResolverPerformanceSettings
limiter *semaphore.Weighted
warningHandler WarningHandler
}

func (r *dockerBase) Acquire(ctx context.Context, weight int64) error {
Expand All @@ -529,12 +562,13 @@ func (r *dockerResolver) base(refspec reference.Spec) (*dockerBase, error) {
return nil, err
}
return &dockerBase{
refspec: refspec,
repository: strings.TrimPrefix(refspec.Locator, host+"/"),
hosts: hosts,
header: r.header,
performances: r.config.Performances,
limiter: r.config.DownloadLimiter,
refspec: refspec,
repository: strings.TrimPrefix(refspec.Locator, host+"/"),
hosts: hosts,
header: r.header,
performances: r.config.Performances,
limiter: r.config.DownloadLimiter,
warningHandler: r.warningHandler,
}, nil
}

Expand Down Expand Up @@ -568,10 +602,12 @@ func (r *dockerBase) request(host RegistryHost, method string, ps ...string) *re
p = p + "/"
}
return &request{
method: method,
path: p,
header: header,
host: host,
method: method,
path: p,
header: header,
host: host,
refspec: r.refspec,
warningHandler: r.warningHandler,
}
}

Expand Down Expand Up @@ -616,12 +652,14 @@ func (r *request) addNamespace(ns string) error {
}

type request struct {
method string
path string
header http.Header
host RegistryHost
body func() (io.ReadCloser, error)
size int64
method string
path string
header http.Header
host RegistryHost
body func() (io.ReadCloser, error)
size int64
refspec reference.Spec
warningHandler WarningHandler
}

func (r *request) clone() *request {
Expand Down Expand Up @@ -681,6 +719,7 @@ func (r *request) do(ctx context.Context) (*http.Response, error) {
return nil, fmt.Errorf("failed to do request: %w", err)
}
log.G(ctx).WithFields(responseFields(resp)).Debug("fetch response received")
reportWarnings(ctx, resp.Header, r.warningHandler)
return resp, nil
}

Expand Down Expand Up @@ -734,6 +773,9 @@ const maxAttempts = 5

func (r *request) doWithRetries(ctx context.Context, lastHost bool, checks ...doChecks) (resp *http.Response, err error) {
attempts := maxAttempts
if r.warningHandler != nil {
ctx = updateWarningSource(ctx, r.refspec)
}
resp, err = r.doWithRetriesInner(ctx, nil, &attempts, lastHost)
if err != nil {
return nil, err
Expand Down
Loading
Loading