Merge pull request #2057 from ktock/export-compression

exporter: Enable to specify the compression type for all layers of the finally exported image
This commit is contained in:
Tõnis Tiigi
2021-07-06 21:52:36 -07:00
committed by GitHub
17 changed files with 1087 additions and 31 deletions

View File

@@ -231,8 +231,8 @@ Keys supported by image output:
* `unpack=true`: unpack image after creation (for use with containerd)
* `dangling-name-prefix=[value]`: name image with `prefix@<digest>` , used for anonymous images
* `name-canonical=true`: add additional canonical name `name@<digest>`
* `compression=[uncompressed,gzip]`: choose compression type for layer, gzip is default value
* `compression=[uncompressed,gzip]`: choose compression type for layers newly created and cached, gzip is default value
* `force-compression=true`: forcefully apply `compression` option to all layers (including already existing layers).
If credentials are required, `buildctl` will attempt to read Docker configuration file `$DOCKER_CONFIG/config.json`.
`$DOCKER_CONFIG` defaults to `~/.docker`.

60
cache/blobs.go vendored
View File

@@ -26,7 +26,8 @@ var ErrNoBlobs = errors.Errorf("no blobs for snapshot")
// computeBlobChain ensures every ref in a parent chain has an associated blob in the content store. If
// a blob is missing and createIfNeeded is true, then the blob will be created, otherwise ErrNoBlobs will
// be returned. Caller must hold a lease when calling this function.
func (sr *immutableRef) computeBlobChain(ctx context.Context, createIfNeeded bool, compressionType compression.Type, s session.Group) error {
// If forceCompression is specified but the blob of compressionType doesn't exist, this function creates it.
func (sr *immutableRef) computeBlobChain(ctx context.Context, createIfNeeded bool, compressionType compression.Type, forceCompression bool, s session.Group) error {
if _, ok := leases.FromContext(ctx); !ok {
return errors.Errorf("missing lease requirement for computeBlobChain")
}
@@ -39,22 +40,31 @@ func (sr *immutableRef) computeBlobChain(ctx context.Context, createIfNeeded boo
ctx = winlayers.UseWindowsLayerMode(ctx)
}
return computeBlobChain(ctx, sr, createIfNeeded, compressionType, s)
return computeBlobChain(ctx, sr, createIfNeeded, compressionType, forceCompression, s)
}
func computeBlobChain(ctx context.Context, sr *immutableRef, createIfNeeded bool, compressionType compression.Type, s session.Group) error {
func computeBlobChain(ctx context.Context, sr *immutableRef, createIfNeeded bool, compressionType compression.Type, forceCompression bool, s session.Group) error {
baseCtx := ctx
eg, ctx := errgroup.WithContext(ctx)
var currentDescr ocispec.Descriptor
if sr.parent != nil {
eg.Go(func() error {
return computeBlobChain(ctx, sr.parent, createIfNeeded, compressionType, s)
return computeBlobChain(ctx, sr.parent, createIfNeeded, compressionType, forceCompression, s)
})
}
eg.Go(func() error {
dp, err := g.Do(ctx, sr.ID(), func(ctx context.Context) (interface{}, error) {
refInfo := sr.Info()
if refInfo.Blob != "" {
if forceCompression {
desc, err := sr.ociDesc()
if err != nil {
return nil, err
}
if err := ensureCompression(ctx, sr, desc, compressionType, s); err != nil {
return nil, err
}
}
return nil, nil
} else if !createIfNeeded {
return nil, errors.WithStack(ErrNoBlobs)
@@ -127,6 +137,12 @@ func computeBlobChain(ctx context.Context, sr *immutableRef, createIfNeeded bool
return nil, errors.Errorf("unknown layer compression type")
}
if forceCompression {
if err := ensureCompression(ctx, sr, descr, compressionType, s); err != nil {
return nil, err
}
}
return descr, nil
})
@@ -224,3 +240,39 @@ func isTypeWindows(sr *immutableRef) bool {
}
return false
}
// ensureCompression ensures the specified ref has the blob of the specified compression Type.
func ensureCompression(ctx context.Context, ref *immutableRef, desc ocispec.Descriptor, compressionType compression.Type, s session.Group) error {
// Resolve converters
layerConvertFunc, _, err := getConverters(desc, compressionType)
if err != nil {
return err
} else if layerConvertFunc == nil {
return nil // no need to convert
}
// First, lookup local content store
if _, err := ref.getCompressionBlob(ctx, compressionType); err == nil {
return nil // found the compression variant. no need to convert.
}
// Convert layer compression type
if err := (lazyRefProvider{
ref: ref,
desc: desc,
dh: ref.descHandlers[desc.Digest],
session: s,
}).Unlazy(ctx); err != nil {
return err
}
newDesc, err := layerConvertFunc(ctx, ref.cm.ContentStore, desc)
if err != nil {
return err
}
// Start to track converted layer
if err := ref.addCompressionBlob(ctx, newDesc.Digest, compressionType); err != nil {
return err
}
return nil
}

138
cache/converter.go vendored Normal file
View File

@@ -0,0 +1,138 @@
package cache
import (
"compress/gzip"
"context"
"fmt"
"io"
"github.com/containerd/containerd/content"
"github.com/containerd/containerd/errdefs"
"github.com/containerd/containerd/images"
"github.com/containerd/containerd/images/converter"
"github.com/containerd/containerd/images/converter/uncompress"
"github.com/containerd/containerd/labels"
"github.com/moby/buildkit/util/compression"
"github.com/opencontainers/go-digest"
ocispec "github.com/opencontainers/image-spec/specs-go/v1"
)
// getConverters returns converter functions according to the specified compression type.
// If no conversion is needed, this returns nil without error.
func getConverters(desc ocispec.Descriptor, compressionType compression.Type) (converter.ConvertFunc, func(string) string, error) {
switch compressionType {
case compression.Uncompressed:
if !images.IsLayerType(desc.MediaType) || uncompress.IsUncompressedType(desc.MediaType) {
// No conversion. No need to return an error here.
return nil, nil, nil
}
return uncompress.LayerConvertFunc, convertMediaTypeToUncompress, nil
case compression.Gzip:
if !images.IsLayerType(desc.MediaType) || isGzipCompressedType(desc.MediaType) {
// No conversion. No need to return an error here.
return nil, nil, nil
}
return gzipLayerConvertFunc, convertMediaTypeToGzip, nil
default:
return nil, nil, fmt.Errorf("unknown compression type during conversion: %q", compressionType)
}
}
func gzipLayerConvertFunc(ctx context.Context, cs content.Store, desc ocispec.Descriptor) (*ocispec.Descriptor, error) {
if !images.IsLayerType(desc.MediaType) || isGzipCompressedType(desc.MediaType) {
// No conversion. No need to return an error here.
return nil, nil
}
// prepare the source and destination
info, err := cs.Info(ctx, desc.Digest)
if err != nil {
return nil, err
}
labelz := info.Labels
if labelz == nil {
labelz = make(map[string]string)
}
ra, err := cs.ReaderAt(ctx, desc)
if err != nil {
return nil, err
}
defer ra.Close()
ref := fmt.Sprintf("convert-gzip-from-%s", desc.Digest)
w, err := cs.Writer(ctx, content.WithRef(ref))
if err != nil {
return nil, err
}
defer w.Close()
if err := w.Truncate(0); err != nil { // Old written data possibly remains
return nil, err
}
zw := gzip.NewWriter(w)
defer zw.Close()
// convert this layer
diffID := digest.Canonical.Digester()
if _, err := io.Copy(zw, io.TeeReader(io.NewSectionReader(ra, 0, ra.Size()), diffID.Hash())); err != nil {
return nil, err
}
if err := zw.Close(); err != nil { // Flush the writer
return nil, err
}
labelz[labels.LabelUncompressed] = diffID.Digest().String() // update diffID label
if err = w.Commit(ctx, 0, "", content.WithLabels(labelz)); err != nil && !errdefs.IsAlreadyExists(err) {
return nil, err
}
if err := w.Close(); err != nil {
return nil, err
}
info, err = cs.Info(ctx, w.Digest())
if err != nil {
return nil, err
}
newDesc := desc
newDesc.MediaType = convertMediaTypeToGzip(newDesc.MediaType)
newDesc.Digest = info.Digest
newDesc.Size = info.Size
return &newDesc, nil
}
func isGzipCompressedType(mt string) bool {
switch mt {
case
images.MediaTypeDockerSchema2LayerGzip,
images.MediaTypeDockerSchema2LayerForeignGzip,
ocispec.MediaTypeImageLayerGzip,
ocispec.MediaTypeImageLayerNonDistributableGzip:
return true
default:
return false
}
}
func convertMediaTypeToUncompress(mt string) string {
switch mt {
case images.MediaTypeDockerSchema2LayerGzip:
return images.MediaTypeDockerSchema2Layer
case images.MediaTypeDockerSchema2LayerForeignGzip:
return images.MediaTypeDockerSchema2LayerForeign
case ocispec.MediaTypeImageLayerGzip:
return ocispec.MediaTypeImageLayer
case ocispec.MediaTypeImageLayerNonDistributableGzip:
return ocispec.MediaTypeImageLayerNonDistributable
default:
return mt
}
}
func convertMediaTypeToGzip(mt string) string {
if uncompress.IsUncompressedType(mt) {
if images.IsDockerType(mt) {
mt += ".gzip"
} else {
mt += "+gzip"
}
return mt
}
return mt
}

63
cache/refs.go vendored
View File

@@ -7,6 +7,7 @@ import (
"sync"
"time"
"github.com/containerd/containerd/content"
"github.com/containerd/containerd/errdefs"
"github.com/containerd/containerd/leases"
"github.com/containerd/containerd/mount"
@@ -45,7 +46,7 @@ type ImmutableRef interface {
Info() RefInfo
Extract(ctx context.Context, s session.Group) error // +progress
GetRemote(ctx context.Context, createIfNeeded bool, compressionType compression.Type, s session.Group) (*solver.Remote, error)
GetRemote(ctx context.Context, createIfNeeded bool, compressionType compression.Type, forceCompression bool, s session.Group) (*solver.Remote, error)
}
type RefInfo struct {
@@ -207,6 +208,20 @@ func (cr *cacheRecord) Size(ctx context.Context) (int64, error) {
if err == nil {
usage.Size += info.Size
}
for k, v := range info.Labels {
// accumulate size of compression variant blobs
if strings.HasPrefix(k, compressionVariantDigestLabelPrefix) {
if cdgst, err := digest.Parse(v); err == nil {
if digest.Digest(dgst) == cdgst {
// do not double count if the label points to this content itself.
continue
}
if info, err := cr.cm.ContentStore.Info(ctx, cdgst); err == nil {
usage.Size += info.Size
}
}
}
}
}
cr.mu.Lock()
setSize(cr.md, usage.Size)
@@ -370,6 +385,52 @@ func (sr *immutableRef) ociDesc() (ocispec.Descriptor, error) {
return desc, nil
}
const compressionVariantDigestLabelPrefix = "buildkit.io/compression/digest."
func compressionVariantDigestLabel(compressionType compression.Type) string {
return compressionVariantDigestLabelPrefix + compressionType.String()
}
func (sr *immutableRef) getCompressionBlob(ctx context.Context, compressionType compression.Type) (content.Info, error) {
cs := sr.cm.ContentStore
info, err := cs.Info(ctx, digest.Digest(getBlob(sr.md)))
if err != nil {
return content.Info{}, err
}
dgstS, ok := info.Labels[compressionVariantDigestLabel(compressionType)]
if ok {
dgst, err := digest.Parse(dgstS)
if err != nil {
return content.Info{}, err
}
return cs.Info(ctx, dgst)
}
return content.Info{}, errdefs.ErrNotFound
}
func (sr *immutableRef) addCompressionBlob(ctx context.Context, dgst digest.Digest, compressionType compression.Type) error {
cs := sr.cm.ContentStore
if err := sr.cm.ManagerOpt.LeaseManager.AddResource(ctx, leases.Lease{ID: sr.ID()}, leases.Resource{
ID: dgst.String(),
Type: "content",
}); err != nil {
return err
}
info, err := cs.Info(ctx, digest.Digest(getBlob(sr.md)))
if err != nil {
return err
}
if info.Labels == nil {
info.Labels = make(map[string]string)
}
cachedVariantLabel := compressionVariantDigestLabel(compressionType)
info.Labels[cachedVariantLabel] = dgst.String()
if _, err := cs.Update(ctx, info, "labels."+cachedVariantLabel); err != nil {
return err
}
return nil
}
// order is from parent->child, sr will be at end of slice
func (sr *immutableRef) parentRefChain() []*immutableRef {
var count int

35
cache/remote.go vendored
View File

@@ -27,24 +27,25 @@ type Unlazier interface {
// GetRemote gets a *solver.Remote from content store for this ref (potentially pulling lazily).
// Note: Use WorkerRef.GetRemote instead as moby integration requires custom GetRemote implementation.
func (sr *immutableRef) GetRemote(ctx context.Context, createIfNeeded bool, compressionType compression.Type, s session.Group) (*solver.Remote, error) {
func (sr *immutableRef) GetRemote(ctx context.Context, createIfNeeded bool, compressionType compression.Type, forceCompression bool, s session.Group) (*solver.Remote, error) {
ctx, done, err := leaseutil.WithLease(ctx, sr.cm.LeaseManager, leaseutil.MakeTemporary)
if err != nil {
return nil, err
}
defer done(ctx)
err = sr.computeBlobChain(ctx, createIfNeeded, compressionType, s)
err = sr.computeBlobChain(ctx, createIfNeeded, compressionType, forceCompression, s)
if err != nil {
return nil, err
}
mprovider := &lazyMultiProvider{mprovider: contentutil.NewMultiProvider(nil)}
chain := sr.parentRefChain()
mproviderBase := contentutil.NewMultiProvider(nil)
mprovider := &lazyMultiProvider{mprovider: mproviderBase}
remote := &solver.Remote{
Provider: mprovider,
}
for _, ref := range sr.parentRefChain() {
for _, ref := range chain {
desc, err := ref.ociDesc()
if err != nil {
return nil, err
@@ -104,6 +105,30 @@ func (sr *immutableRef) GetRemote(ctx context.Context, createIfNeeded bool, comp
}
}
if forceCompression {
// ensure the compression type.
// compressed blob must be created and stored in the content store.
_, convertMediaTypeFunc, err := getConverters(desc, compressionType)
if err != nil {
return nil, err
}
if convertMediaTypeFunc != nil {
// needs conversion
info, err := ref.getCompressionBlob(ctx, compressionType)
if err != nil {
return nil, err
}
newDesc := desc
newDesc.MediaType = convertMediaTypeFunc(newDesc.MediaType)
newDesc.Digest = info.Digest
newDesc.Size = info.Size
if desc.Digest != newDesc.Digest {
mproviderBase.Add(newDesc.Digest, ref.cm.ContentStore)
}
desc = newDesc
}
}
remote.Descriptors = append(remote.Descriptors, desc)
mprovider.Add(lazyRefProvider{
ref: ref,

View File

@@ -2050,6 +2050,22 @@ func testBuildExportWithUncompressed(t *testing.T, sb integration.Sandbox) {
}, nil)
require.NoError(t, err)
allCompressedTarget := registry + "/buildkit/build/exporter:withallcompressed"
_, err = c.Solve(context.TODO(), def, SolveOpt{
Exports: []ExportEntry{
{
Type: ExporterImage,
Attrs: map[string]string{
"name": allCompressedTarget,
"push": "true",
"compression": "gzip",
"force-compression": "true",
},
},
},
}, nil)
require.NoError(t, err)
if cdAddress == "" {
t.Skip("rest of test requires containerd worker")
}
@@ -2058,9 +2074,12 @@ func testBuildExportWithUncompressed(t *testing.T, sb integration.Sandbox) {
require.NoError(t, err)
err = client.ImageService().Delete(ctx, compressedTarget, images.SynchronousDelete())
require.NoError(t, err)
err = client.ImageService().Delete(ctx, allCompressedTarget, images.SynchronousDelete())
require.NoError(t, err)
checkAllReleasable(t, c, sb, true)
// check if the new layer is compressed with compression option
img, err := client.Pull(ctx, compressedTarget)
require.NoError(t, err)
@@ -2099,6 +2118,51 @@ func testBuildExportWithUncompressed(t *testing.T, sb integration.Sandbox) {
require.True(t, ok)
require.Equal(t, int32(item.Header.Typeflag), tar.TypeReg)
require.Equal(t, []byte("gzip"), item.Data)
err = client.ImageService().Delete(ctx, compressedTarget, images.SynchronousDelete())
require.NoError(t, err)
checkAllReleasable(t, c, sb, true)
// check if all layers are compressed with force-compressoin option
img, err = client.Pull(ctx, allCompressedTarget)
require.NoError(t, err)
dt, err = content.ReadBlob(ctx, img.ContentStore(), img.Target())
require.NoError(t, err)
mfst = struct {
MediaType string `json:"mediaType,omitempty"`
ocispec.Manifest
}{}
err = json.Unmarshal(dt, &mfst)
require.NoError(t, err)
require.Equal(t, 2, len(mfst.Layers))
require.Equal(t, images.MediaTypeDockerSchema2LayerGzip, mfst.Layers[0].MediaType)
require.Equal(t, images.MediaTypeDockerSchema2LayerGzip, mfst.Layers[1].MediaType)
dt, err = content.ReadBlob(ctx, img.ContentStore(), ocispec.Descriptor{Digest: mfst.Layers[0].Digest})
require.NoError(t, err)
m, err = testutil.ReadTarToMap(dt, true)
require.NoError(t, err)
item, ok = m["data"]
require.True(t, ok)
require.Equal(t, int32(item.Header.Typeflag), tar.TypeReg)
require.Equal(t, []byte("uncompressed"), item.Data)
dt, err = content.ReadBlob(ctx, img.ContentStore(), ocispec.Descriptor{Digest: mfst.Layers[1].Digest})
require.NoError(t, err)
m, err = testutil.ReadTarToMap(dt, true)
require.NoError(t, err)
item, ok = m["data"]
require.True(t, ok)
require.Equal(t, int32(item.Header.Typeflag), tar.TypeReg)
require.Equal(t, []byte("gzip"), item.Data)
}
func testBuildPushAndValidate(t *testing.T, sb integration.Sandbox) {

View File

@@ -37,6 +37,7 @@ const (
keyDanglingPrefix = "dangling-name-prefix"
keyNameCanonical = "name-canonical"
keyLayerCompression = "compression"
keyForceCompression = "force-compression"
ociTypes = "oci-mediatypes"
)
@@ -142,6 +143,16 @@ func (e *imageExporter) Resolve(ctx context.Context, opt map[string]string) (exp
default:
return nil, errors.Errorf("unsupported layer compression type: %v", v)
}
case keyForceCompression:
if v == "" {
i.forceCompression = true
continue
}
b, err := strconv.ParseBool(v)
if err != nil {
return nil, errors.Wrapf(err, "non-bool value specified for %s", k)
}
i.forceCompression = b
default:
if i.meta == nil {
i.meta = make(map[string][]byte)
@@ -163,6 +174,7 @@ type imageExporterInstance struct {
nameCanonical bool
danglingPrefix string
layerCompression compression.Type
forceCompression bool
meta map[string][]byte
}
@@ -184,7 +196,7 @@ func (e *imageExporterInstance) Export(ctx context.Context, src exporter.Source,
}
defer done(context.TODO())
desc, err := e.opt.ImageWriter.Commit(ctx, src, e.ociTypes, e.layerCompression, sessionID)
desc, err := e.opt.ImageWriter.Commit(ctx, src, e.ociTypes, e.layerCompression, e.forceCompression, sessionID)
if err != nil {
return nil, err
}
@@ -242,7 +254,7 @@ func (e *imageExporterInstance) Export(ctx context.Context, src exporter.Source,
annotations := map[digest.Digest]map[string]string{}
mprovider := contentutil.NewMultiProvider(e.opt.ImageWriter.ContentStore())
if src.Ref != nil {
remote, err := src.Ref.GetRemote(ctx, false, e.layerCompression, session.NewGroup(sessionID))
remote, err := src.Ref.GetRemote(ctx, false, e.layerCompression, e.forceCompression, session.NewGroup(sessionID))
if err != nil {
return nil, err
}
@@ -253,7 +265,7 @@ func (e *imageExporterInstance) Export(ctx context.Context, src exporter.Source,
}
if len(src.Refs) > 0 {
for _, r := range src.Refs {
remote, err := r.GetRemote(ctx, false, e.layerCompression, session.NewGroup(sessionID))
remote, err := r.GetRemote(ctx, false, e.layerCompression, e.forceCompression, session.NewGroup(sessionID))
if err != nil {
return nil, err
}
@@ -306,7 +318,7 @@ func (e *imageExporterInstance) unpackImage(ctx context.Context, img images.Imag
}
}
remote, err := topLayerRef.GetRemote(ctx, true, e.layerCompression, s)
remote, err := topLayerRef.GetRemote(ctx, true, e.layerCompression, e.forceCompression, s)
if err != nil {
return err
}

View File

@@ -44,7 +44,7 @@ type ImageWriter struct {
opt WriterOpt
}
func (ic *ImageWriter) Commit(ctx context.Context, inp exporter.Source, oci bool, compressionType compression.Type, sessionID string) (*ocispec.Descriptor, error) {
func (ic *ImageWriter) Commit(ctx context.Context, inp exporter.Source, oci bool, compressionType compression.Type, forceCompression bool, sessionID string) (*ocispec.Descriptor, error) {
platformsBytes, ok := inp.Metadata[exptypes.ExporterPlatformsKey]
if len(inp.Refs) > 0 && !ok {
@@ -52,7 +52,7 @@ func (ic *ImageWriter) Commit(ctx context.Context, inp exporter.Source, oci bool
}
if len(inp.Refs) == 0 {
remotes, err := ic.exportLayers(ctx, compressionType, session.NewGroup(sessionID), inp.Ref)
remotes, err := ic.exportLayers(ctx, compressionType, forceCompression, session.NewGroup(sessionID), inp.Ref)
if err != nil {
return nil, err
}
@@ -83,7 +83,7 @@ func (ic *ImageWriter) Commit(ctx context.Context, inp exporter.Source, oci bool
refs = append(refs, r)
}
remotes, err := ic.exportLayers(ctx, compressionType, session.NewGroup(sessionID), refs...)
remotes, err := ic.exportLayers(ctx, compressionType, forceCompression, session.NewGroup(sessionID), refs...)
if err != nil {
return nil, err
}
@@ -148,7 +148,7 @@ func (ic *ImageWriter) Commit(ctx context.Context, inp exporter.Source, oci bool
return &idxDesc, nil
}
func (ic *ImageWriter) exportLayers(ctx context.Context, compressionType compression.Type, s session.Group, refs ...cache.ImmutableRef) ([]solver.Remote, error) {
func (ic *ImageWriter) exportLayers(ctx context.Context, compressionType compression.Type, forceCompression bool, s session.Group, refs ...cache.ImmutableRef) ([]solver.Remote, error) {
eg, ctx := errgroup.WithContext(ctx)
layersDone := oneOffProgress(ctx, "exporting layers")
@@ -160,7 +160,7 @@ func (ic *ImageWriter) exportLayers(ctx context.Context, compressionType compres
return
}
eg.Go(func() error {
remote, err := ref.GetRemote(ctx, true, compressionType, s)
remote, err := ref.GetRemote(ctx, true, compressionType, forceCompression, s)
if err != nil {
return err
}

View File

@@ -32,6 +32,7 @@ const (
VariantOCI = "oci"
VariantDocker = "docker"
ociTypes = "oci-mediatypes"
keyForceCompression = "force-compression"
)
type Opt struct {
@@ -69,6 +70,16 @@ func (e *imageExporter) Resolve(ctx context.Context, opt map[string]string) (exp
default:
return nil, errors.Errorf("unsupported layer compression type: %v", v)
}
case keyForceCompression:
if v == "" {
i.forceCompression = true
continue
}
b, err := strconv.ParseBool(v)
if err != nil {
return nil, errors.Wrapf(err, "non-bool value specified for %s", k)
}
i.forceCompression = b
case ociTypes:
ot = new(bool)
if v == "" {
@@ -101,6 +112,7 @@ type imageExporterInstance struct {
name string
ociTypes bool
layerCompression compression.Type
forceCompression bool
}
func (e *imageExporterInstance) Name() string {
@@ -125,7 +137,7 @@ func (e *imageExporterInstance) Export(ctx context.Context, src exporter.Source,
}
defer done(context.TODO())
desc, err := e.opt.ImageWriter.Commit(ctx, src, e.ociTypes, e.layerCompression, sessionID)
desc, err := e.opt.ImageWriter.Commit(ctx, src, e.ociTypes, e.layerCompression, e.forceCompression, sessionID)
if err != nil {
return nil, err
}
@@ -181,7 +193,7 @@ func (e *imageExporterInstance) Export(ctx context.Context, src exporter.Source,
mprovider := contentutil.NewMultiProvider(e.opt.ImageWriter.ContentStore())
if src.Ref != nil {
remote, err := src.Ref.GetRemote(ctx, false, e.layerCompression, session.NewGroup(sessionID))
remote, err := src.Ref.GetRemote(ctx, false, e.layerCompression, e.forceCompression, session.NewGroup(sessionID))
if err != nil {
return nil, err
}
@@ -198,7 +210,7 @@ func (e *imageExporterInstance) Export(ctx context.Context, src exporter.Source,
}
if len(src.Refs) > 0 {
for _, r := range src.Refs {
remote, err := r.GetRemote(ctx, false, e.layerCompression, session.NewGroup(sessionID))
remote, err := r.GetRemote(ctx, false, e.layerCompression, e.forceCompression, session.NewGroup(sessionID))
if err != nil {
return nil, err
}

View File

@@ -89,6 +89,6 @@ func workerRefConverter(g session.Group) func(ctx context.Context, res solver.Re
return nil, errors.Errorf("invalid result: %T", res.Sys())
}
return ref.GetRemote(ctx, true, compression.Default, g)
return ref.GetRemote(ctx, true, compression.Default, false, g)
}
}

View File

@@ -274,7 +274,7 @@ func inlineCache(ctx context.Context, e remotecache.Exporter, res solver.CachedR
return nil, errors.Errorf("invalid reference: %T", res.Sys())
}
remote, err := workerRef.GetRemote(ctx, true, compression.Default, g)
remote, err := workerRef.GetRemote(ctx, true, compression.Default, false, g)
if err != nil || remote == nil {
return nil, nil
}

View File

@@ -0,0 +1,126 @@
/*
Copyright The containerd Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
// Package converter provides image converter
package converter
import (
"context"
"github.com/containerd/containerd/content"
"github.com/containerd/containerd/images"
"github.com/containerd/containerd/leases"
"github.com/containerd/containerd/platforms"
)
type convertOpts struct {
layerConvertFunc ConvertFunc
docker2oci bool
indexConvertFunc ConvertFunc
platformMC platforms.MatchComparer
}
// Opt is an option for Convert()
type Opt func(*convertOpts) error
// WithLayerConvertFunc specifies the function that converts layers.
func WithLayerConvertFunc(fn ConvertFunc) Opt {
return func(copts *convertOpts) error {
copts.layerConvertFunc = fn
return nil
}
}
// WithDockerToOCI converts Docker media types into OCI ones.
func WithDockerToOCI(v bool) Opt {
return func(copts *convertOpts) error {
copts.docker2oci = true
return nil
}
}
// WithPlatform specifies the platform.
// Defaults to all platforms.
func WithPlatform(p platforms.MatchComparer) Opt {
return func(copts *convertOpts) error {
copts.platformMC = p
return nil
}
}
// WithIndexConvertFunc specifies the function that converts manifests and index (manifest lists).
// Defaults to DefaultIndexConvertFunc.
func WithIndexConvertFunc(fn ConvertFunc) Opt {
return func(copts *convertOpts) error {
copts.indexConvertFunc = fn
return nil
}
}
// Client is implemented by *containerd.Client .
type Client interface {
WithLease(ctx context.Context, opts ...leases.Opt) (context.Context, func(context.Context) error, error)
ContentStore() content.Store
ImageService() images.Store
}
// Convert converts an image.
func Convert(ctx context.Context, client Client, dstRef, srcRef string, opts ...Opt) (*images.Image, error) {
var copts convertOpts
for _, o := range opts {
if err := o(&copts); err != nil {
return nil, err
}
}
if copts.platformMC == nil {
copts.platformMC = platforms.All
}
if copts.indexConvertFunc == nil {
copts.indexConvertFunc = DefaultIndexConvertFunc(copts.layerConvertFunc, copts.docker2oci, copts.platformMC)
}
ctx, done, err := client.WithLease(ctx)
if err != nil {
return nil, err
}
defer done(ctx)
cs := client.ContentStore()
is := client.ImageService()
srcImg, err := is.Get(ctx, srcRef)
if err != nil {
return nil, err
}
dstDesc, err := copts.indexConvertFunc(ctx, cs, srcImg.Target)
if err != nil {
return nil, err
}
dstImg := srcImg
dstImg.Name = dstRef
if dstDesc != nil {
dstImg.Target = *dstDesc
}
var res images.Image
if dstRef != srcRef {
_ = is.Delete(ctx, dstRef)
res, err = is.Create(ctx, dstImg)
} else {
res, err = is.Update(ctx, dstImg)
}
return &res, err
}

View File

@@ -0,0 +1,442 @@
/*
Copyright The containerd Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package converter
import (
"bytes"
"context"
"encoding/json"
"fmt"
"strings"
"sync"
"github.com/containerd/containerd/content"
"github.com/containerd/containerd/images"
"github.com/containerd/containerd/platforms"
"github.com/opencontainers/go-digest"
ocispec "github.com/opencontainers/image-spec/specs-go/v1"
"github.com/sirupsen/logrus"
"golang.org/x/sync/errgroup"
)
// ConvertFunc returns a converted content descriptor.
// When the content was not converted, ConvertFunc returns nil.
type ConvertFunc func(ctx context.Context, cs content.Store, desc ocispec.Descriptor) (*ocispec.Descriptor, error)
// DefaultIndexConvertFunc is the default convert func used by Convert.
func DefaultIndexConvertFunc(layerConvertFunc ConvertFunc, docker2oci bool, platformMC platforms.MatchComparer) ConvertFunc {
c := &defaultConverter{
layerConvertFunc: layerConvertFunc,
docker2oci: docker2oci,
platformMC: platformMC,
diffIDMap: make(map[digest.Digest]digest.Digest),
}
return c.convert
}
type defaultConverter struct {
layerConvertFunc ConvertFunc
docker2oci bool
platformMC platforms.MatchComparer
diffIDMap map[digest.Digest]digest.Digest // key: old diffID, value: new diffID
diffIDMapMu sync.RWMutex
}
// convert dispatches desc.MediaType and calls c.convert{Layer,Manifest,Index,Config}.
//
// Also converts media type if c.docker2oci is set.
func (c *defaultConverter) convert(ctx context.Context, cs content.Store, desc ocispec.Descriptor) (*ocispec.Descriptor, error) {
var (
newDesc *ocispec.Descriptor
err error
)
if images.IsLayerType(desc.MediaType) {
newDesc, err = c.convertLayer(ctx, cs, desc)
} else if images.IsManifestType(desc.MediaType) {
newDesc, err = c.convertManifest(ctx, cs, desc)
} else if images.IsIndexType(desc.MediaType) {
newDesc, err = c.convertIndex(ctx, cs, desc)
} else if images.IsConfigType(desc.MediaType) {
newDesc, err = c.convertConfig(ctx, cs, desc)
}
if err != nil {
return nil, err
}
if images.IsDockerType(desc.MediaType) {
if c.docker2oci {
if newDesc == nil {
newDesc = copyDesc(desc)
}
newDesc.MediaType = ConvertDockerMediaTypeToOCI(newDesc.MediaType)
} else if (newDesc == nil && len(desc.Annotations) != 0) || (newDesc != nil && len(newDesc.Annotations) != 0) {
// Annotations is supported only on OCI manifest.
// We need to remove annotations for Docker media types.
if newDesc == nil {
newDesc = copyDesc(desc)
}
newDesc.Annotations = nil
}
}
logrus.WithField("old", desc).WithField("new", newDesc).Debugf("converted")
return newDesc, nil
}
func copyDesc(desc ocispec.Descriptor) *ocispec.Descriptor {
descCopy := desc
return &descCopy
}
// convertLayer converts image image layers if c.layerConvertFunc is set.
//
// c.layerConvertFunc can be nil, e.g., for converting Docker media types to OCI ones.
func (c *defaultConverter) convertLayer(ctx context.Context, cs content.Store, desc ocispec.Descriptor) (*ocispec.Descriptor, error) {
if c.layerConvertFunc != nil {
return c.layerConvertFunc(ctx, cs, desc)
}
return nil, nil
}
// convertManifest converts image manifests.
//
// - clears `.mediaType` if the target format is OCI
//
// - records diff ID changes in c.diffIDMap
func (c *defaultConverter) convertManifest(ctx context.Context, cs content.Store, desc ocispec.Descriptor) (*ocispec.Descriptor, error) {
var (
manifest DualManifest
modified bool
)
labels, err := readJSON(ctx, cs, &manifest, desc)
if err != nil {
return nil, err
}
if labels == nil {
labels = make(map[string]string)
}
if images.IsDockerType(manifest.MediaType) && c.docker2oci {
manifest.MediaType = ""
modified = true
}
var mu sync.Mutex
eg, ctx2 := errgroup.WithContext(ctx)
for i, l := range manifest.Layers {
i := i
l := l
oldDiffID, err := images.GetDiffID(ctx, cs, l)
if err != nil {
return nil, err
}
eg.Go(func() error {
newL, err := c.convert(ctx2, cs, l)
if err != nil {
return err
}
if newL != nil {
mu.Lock()
// update GC labels
ClearGCLabels(labels, l.Digest)
labelKey := fmt.Sprintf("containerd.io/gc.ref.content.l.%d", i)
labels[labelKey] = newL.Digest.String()
manifest.Layers[i] = *newL
modified = true
mu.Unlock()
// diffID changes if the tar entries were modified.
// diffID stays same if only the compression type was changed.
// When diffID changed, add a map entry so that we can update image config.
newDiffID, err := images.GetDiffID(ctx, cs, *newL)
if err != nil {
return err
}
if newDiffID != oldDiffID {
c.diffIDMapMu.Lock()
c.diffIDMap[oldDiffID] = newDiffID
c.diffIDMapMu.Unlock()
}
}
return nil
})
}
if err := eg.Wait(); err != nil {
return nil, err
}
newConfig, err := c.convert(ctx, cs, manifest.Config)
if err != nil {
return nil, err
}
if newConfig != nil {
ClearGCLabels(labels, manifest.Config.Digest)
labels["containerd.io/gc.ref.content.config"] = newConfig.Digest.String()
manifest.Config = *newConfig
modified = true
}
if modified {
return writeJSON(ctx, cs, &manifest, desc, labels)
}
return nil, nil
}
// convertIndex converts image index.
//
// - clears `.mediaType` if the target format is OCI
//
// - clears manifest entries that do not match c.platformMC
func (c *defaultConverter) convertIndex(ctx context.Context, cs content.Store, desc ocispec.Descriptor) (*ocispec.Descriptor, error) {
var (
index DualIndex
modified bool
)
labels, err := readJSON(ctx, cs, &index, desc)
if err != nil {
return nil, err
}
if labels == nil {
labels = make(map[string]string)
}
if images.IsDockerType(index.MediaType) && c.docker2oci {
index.MediaType = ""
modified = true
}
newManifests := make([]ocispec.Descriptor, len(index.Manifests))
newManifestsToBeRemoved := make(map[int]struct{}) // slice index
var mu sync.Mutex
eg, ctx2 := errgroup.WithContext(ctx)
for i, mani := range index.Manifests {
i := i
mani := mani
labelKey := fmt.Sprintf("containerd.io/gc.ref.content.m.%d", i)
eg.Go(func() error {
if mani.Platform != nil && !c.platformMC.Match(*mani.Platform) {
mu.Lock()
ClearGCLabels(labels, mani.Digest)
newManifestsToBeRemoved[i] = struct{}{}
modified = true
mu.Unlock()
return nil
}
newMani, err := c.convert(ctx2, cs, mani)
if err != nil {
return err
}
mu.Lock()
if newMani != nil {
ClearGCLabels(labels, mani.Digest)
labels[labelKey] = newMani.Digest.String()
// NOTE: for keeping manifest order, we specify `i` index explicitly
newManifests[i] = *newMani
modified = true
} else {
newManifests[i] = mani
}
mu.Unlock()
return nil
})
}
if err := eg.Wait(); err != nil {
return nil, err
}
if modified {
var newManifestsClean []ocispec.Descriptor
for i, m := range newManifests {
if _, ok := newManifestsToBeRemoved[i]; !ok {
newManifestsClean = append(newManifestsClean, m)
}
}
index.Manifests = newManifestsClean
return writeJSON(ctx, cs, &index, desc, labels)
}
return nil, nil
}
// convertConfig converts image config contents.
//
// - updates `.rootfs.diff_ids` using c.diffIDMap .
//
// - clears legacy `.config.Image` and `.container_config.Image` fields if `.rootfs.diff_ids` was updated.
func (c *defaultConverter) convertConfig(ctx context.Context, cs content.Store, desc ocispec.Descriptor) (*ocispec.Descriptor, error) {
var (
cfg DualConfig
cfgAsOCI ocispec.Image // read only, used for parsing cfg
modified bool
)
labels, err := readJSON(ctx, cs, &cfg, desc)
if err != nil {
return nil, err
}
if labels == nil {
labels = make(map[string]string)
}
if _, err := readJSON(ctx, cs, &cfgAsOCI, desc); err != nil {
return nil, err
}
if rootfs := cfgAsOCI.RootFS; rootfs.Type == "layers" {
rootfsModified := false
c.diffIDMapMu.RLock()
for i, oldDiffID := range rootfs.DiffIDs {
if newDiffID, ok := c.diffIDMap[oldDiffID]; ok && newDiffID != oldDiffID {
rootfs.DiffIDs[i] = newDiffID
rootfsModified = true
}
}
c.diffIDMapMu.RUnlock()
if rootfsModified {
rootfsB, err := json.Marshal(rootfs)
if err != nil {
return nil, err
}
cfg["rootfs"] = (*json.RawMessage)(&rootfsB)
modified = true
}
}
if modified {
// cfg may have dummy value for legacy `.config.Image` and `.container_config.Image`
// We should clear the ID if we changed the diff IDs.
if _, err := clearDockerV1DummyID(cfg); err != nil {
return nil, err
}
return writeJSON(ctx, cs, &cfg, desc, labels)
}
return nil, nil
}
// clearDockerV1DummyID clears the dummy values for legacy `.config.Image` and `.container_config.Image`.
// Returns true if the cfg was modified.
func clearDockerV1DummyID(cfg DualConfig) (bool, error) {
var modified bool
f := func(k string) error {
if configX, ok := cfg[k]; ok && configX != nil {
var configField map[string]*json.RawMessage
if err := json.Unmarshal(*configX, &configField); err != nil {
return err
}
delete(configField, "Image")
b, err := json.Marshal(configField)
if err != nil {
return err
}
cfg[k] = (*json.RawMessage)(&b)
modified = true
}
return nil
}
if err := f("config"); err != nil {
return modified, err
}
if err := f("container_config"); err != nil {
return modified, err
}
return modified, nil
}
// ObjectWithMediaType represents an object with a MediaType field
type ObjectWithMediaType struct {
// MediaType appears on Docker manifests and manifest lists.
// MediaType does not appear on OCI manifests and index
MediaType string `json:"mediaType,omitempty"`
}
// DualManifest covers Docker manifest and OCI manifest
type DualManifest struct {
ocispec.Manifest
ObjectWithMediaType
}
// DualIndex covers Docker manifest list and OCI index
type DualIndex struct {
ocispec.Index
ObjectWithMediaType
}
// DualConfig covers Docker config (v1.0, v1.1, v1.2) and OCI config.
// Unmarshalled as map[string]*json.RawMessage to retain unknown fields on remarshalling.
type DualConfig map[string]*json.RawMessage
func readJSON(ctx context.Context, cs content.Store, x interface{}, desc ocispec.Descriptor) (map[string]string, error) {
info, err := cs.Info(ctx, desc.Digest)
if err != nil {
return nil, err
}
labels := info.Labels
b, err := content.ReadBlob(ctx, cs, desc)
if err != nil {
return nil, err
}
if err := json.Unmarshal(b, x); err != nil {
return nil, err
}
return labels, nil
}
func writeJSON(ctx context.Context, cs content.Store, x interface{}, oldDesc ocispec.Descriptor, labels map[string]string) (*ocispec.Descriptor, error) {
b, err := json.Marshal(x)
if err != nil {
return nil, err
}
dgst := digest.SHA256.FromBytes(b)
ref := fmt.Sprintf("converter-write-json-%s", dgst.String())
w, err := content.OpenWriter(ctx, cs, content.WithRef(ref))
if err != nil {
return nil, err
}
if err := content.Copy(ctx, w, bytes.NewReader(b), int64(len(b)), dgst, content.WithLabels(labels)); err != nil {
return nil, err
}
if err := w.Close(); err != nil {
return nil, err
}
newDesc := oldDesc
newDesc.Size = int64(len(b))
newDesc.Digest = dgst
return &newDesc, nil
}
// ConvertDockerMediaTypeToOCI converts a media type string
func ConvertDockerMediaTypeToOCI(mt string) string {
switch mt {
case images.MediaTypeDockerSchema2ManifestList:
return ocispec.MediaTypeImageIndex
case images.MediaTypeDockerSchema2Manifest:
return ocispec.MediaTypeImageManifest
case images.MediaTypeDockerSchema2LayerGzip:
return ocispec.MediaTypeImageLayerGzip
case images.MediaTypeDockerSchema2LayerForeignGzip:
return ocispec.MediaTypeImageLayerNonDistributableGzip
case images.MediaTypeDockerSchema2Layer:
return ocispec.MediaTypeImageLayer
case images.MediaTypeDockerSchema2LayerForeign:
return ocispec.MediaTypeImageLayerNonDistributable
case images.MediaTypeDockerSchema2Config:
return ocispec.MediaTypeImageConfig
default:
return mt
}
}
// ClearGCLabels clears GC labels for the given digest.
func ClearGCLabels(labels map[string]string, dgst digest.Digest) {
for k, v := range labels {
if v == dgst.String() && strings.HasPrefix(k, "containerd.io/gc.ref.content") {
delete(labels, k)
}
}
}

View File

@@ -0,0 +1,122 @@
/*
Copyright The containerd Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package uncompress
import (
"compress/gzip"
"context"
"fmt"
"io"
"github.com/containerd/containerd/content"
"github.com/containerd/containerd/errdefs"
"github.com/containerd/containerd/images"
"github.com/containerd/containerd/images/converter"
"github.com/containerd/containerd/labels"
ocispec "github.com/opencontainers/image-spec/specs-go/v1"
)
var _ converter.ConvertFunc = LayerConvertFunc
// LayerConvertFunc converts tar.gz layers into uncompressed tar layers.
// Media type is changed, e.g., "application/vnd.oci.image.layer.v1.tar+gzip" -> "application/vnd.oci.image.layer.v1.tar"
func LayerConvertFunc(ctx context.Context, cs content.Store, desc ocispec.Descriptor) (*ocispec.Descriptor, error) {
if !images.IsLayerType(desc.MediaType) || IsUncompressedType(desc.MediaType) {
// No conversion. No need to return an error here.
return nil, nil
}
info, err := cs.Info(ctx, desc.Digest)
if err != nil {
return nil, err
}
readerAt, err := cs.ReaderAt(ctx, desc)
if err != nil {
return nil, err
}
defer readerAt.Close()
sr := io.NewSectionReader(readerAt, 0, desc.Size)
newR, err := gzip.NewReader(sr)
if err != nil {
return nil, err
}
defer newR.Close()
ref := fmt.Sprintf("convert-uncompress-from-%s", desc.Digest)
w, err := content.OpenWriter(ctx, cs, content.WithRef(ref))
if err != nil {
return nil, err
}
defer w.Close()
// Reset the writing position
// Old writer possibly remains without aborted
// (e.g. conversion interrupted by a signal)
if err := w.Truncate(0); err != nil {
return nil, err
}
n, err := io.Copy(w, newR)
if err != nil {
return nil, err
}
if err := newR.Close(); err != nil {
return nil, err
}
// no need to retain "containerd.io/uncompressed" label, but retain other labels ("containerd.io/distribution.source.*")
labelsMap := info.Labels
delete(labelsMap, labels.LabelUncompressed)
if err = w.Commit(ctx, 0, "", content.WithLabels(labelsMap)); err != nil && !errdefs.IsAlreadyExists(err) {
return nil, err
}
if err := w.Close(); err != nil {
return nil, err
}
newDesc := desc
newDesc.Digest = w.Digest()
newDesc.Size = n
newDesc.MediaType = convertMediaType(newDesc.MediaType)
return &newDesc, nil
}
// IsUncompressedType returns whether the provided media type is considered
// an uncompressed layer type
func IsUncompressedType(mt string) bool {
switch mt {
case
images.MediaTypeDockerSchema2Layer,
images.MediaTypeDockerSchema2LayerForeign,
ocispec.MediaTypeImageLayer,
ocispec.MediaTypeImageLayerNonDistributable:
return true
default:
return false
}
}
func convertMediaType(mt string) string {
switch mt {
case images.MediaTypeDockerSchema2LayerGzip:
return images.MediaTypeDockerSchema2Layer
case images.MediaTypeDockerSchema2LayerForeignGzip:
return images.MediaTypeDockerSchema2LayerForeign
case ocispec.MediaTypeImageLayerGzip:
return ocispec.MediaTypeImageLayer
case ocispec.MediaTypeImageLayerNonDistributableGzip:
return ocispec.MediaTypeImageLayerNonDistributable
default:
return mt
}
}

2
vendor/modules.txt vendored
View File

@@ -79,6 +79,8 @@ github.com/containerd/containerd/gc
github.com/containerd/containerd/identifiers
github.com/containerd/containerd/images
github.com/containerd/containerd/images/archive
github.com/containerd/containerd/images/converter
github.com/containerd/containerd/images/converter/uncompress
github.com/containerd/containerd/labels
github.com/containerd/containerd/leases
github.com/containerd/containerd/leases/proxy

View File

@@ -79,7 +79,7 @@ func (s *cacheResultStorage) LoadRemote(ctx context.Context, res solver.CacheRes
}
defer ref.Release(context.TODO())
wref := WorkerRef{ref, w}
remote, err := wref.GetRemote(ctx, false, compression.Default, g)
remote, err := wref.GetRemote(ctx, false, compression.Default, false, g)
if err != nil {
return nil, nil // ignore error. loadRemote is best effort
}

View File

@@ -29,13 +29,13 @@ func (wr *WorkerRef) ID() string {
// GetRemote method abstracts ImmutableRef's GetRemote to allow a Worker to override.
// This is needed for moby integration.
// Use this method instead of calling ImmutableRef.GetRemote() directly.
func (wr *WorkerRef) GetRemote(ctx context.Context, createIfNeeded bool, compressionType compression.Type, g session.Group) (*solver.Remote, error) {
func (wr *WorkerRef) GetRemote(ctx context.Context, createIfNeeded bool, compressionType compression.Type, forceCompression bool, g session.Group) (*solver.Remote, error) {
if w, ok := wr.Worker.(interface {
GetRemote(context.Context, cache.ImmutableRef, bool, compression.Type, session.Group) (*solver.Remote, error)
GetRemote(context.Context, cache.ImmutableRef, bool, compression.Type, bool, session.Group) (*solver.Remote, error)
}); ok {
return w.GetRemote(ctx, wr.ImmutableRef, createIfNeeded, compressionType, g)
return w.GetRemote(ctx, wr.ImmutableRef, createIfNeeded, compressionType, forceCompression, g)
}
return wr.ImmutableRef.GetRemote(ctx, createIfNeeded, compressionType, g)
return wr.ImmutableRef.GetRemote(ctx, createIfNeeded, compressionType, forceCompression, g)
}
type workerRefResult struct {