mirror of
https://github.com/containerd/containerd.git
synced 2026-08-09 17:39:22 +00:00
288 lines
8.8 KiB
Go
288 lines
8.8 KiB
Go
/*
|
|
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 transfer
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
|
|
"github.com/containerd/errdefs"
|
|
"github.com/containerd/log"
|
|
"github.com/containerd/platforms"
|
|
"github.com/containerd/plugin"
|
|
"github.com/containerd/plugin/registry"
|
|
|
|
"github.com/containerd/containerd/v2/core/diff"
|
|
"github.com/containerd/containerd/v2/core/leases"
|
|
"github.com/containerd/containerd/v2/core/metadata"
|
|
"github.com/containerd/containerd/v2/core/transfer/local"
|
|
"github.com/containerd/containerd/v2/core/unpack"
|
|
"github.com/containerd/containerd/v2/defaults"
|
|
"github.com/containerd/containerd/v2/internal/kmutex"
|
|
"github.com/containerd/containerd/v2/pkg/imageverifier"
|
|
"github.com/containerd/containerd/v2/plugins"
|
|
specs "github.com/opencontainers/image-spec/specs-go/v1"
|
|
|
|
// Load packages with type registrations
|
|
_ "github.com/containerd/containerd/v2/core/transfer/archive"
|
|
_ "github.com/containerd/containerd/v2/core/transfer/image"
|
|
_ "github.com/containerd/containerd/v2/core/transfer/registry"
|
|
)
|
|
|
|
// Register local transfer service plugin
|
|
func init() {
|
|
registry.Register(&plugin.Registration{
|
|
Type: plugins.TransferPlugin,
|
|
ID: "local",
|
|
Requires: []plugin.Type{
|
|
plugins.LeasePlugin,
|
|
plugins.MetadataPlugin,
|
|
plugins.DiffPlugin,
|
|
plugins.ImageVerifierPlugin,
|
|
plugins.SnapshotPlugin,
|
|
},
|
|
Config: defaultConfig(),
|
|
InitFn: func(ic *plugin.InitContext) (any, error) {
|
|
config := ic.Config.(*transferConfig)
|
|
m, err := ic.GetSingle(plugins.MetadataPlugin)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
ms := m.(*metadata.DB)
|
|
|
|
var lc local.TransferConfig
|
|
|
|
l, err := ic.GetSingle(plugins.LeasePlugin)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
lc.Leases = l.(leases.Manager)
|
|
|
|
vps, err := ic.GetByType(plugins.ImageVerifierPlugin)
|
|
if err != nil && !errors.Is(err, plugin.ErrPluginNotFound) {
|
|
return nil, err
|
|
}
|
|
if len(vps) > 0 {
|
|
lc.Verifiers = make(map[string]imageverifier.ImageVerifier)
|
|
for name, vp := range vps {
|
|
lc.Verifiers[name] = vp.(imageverifier.ImageVerifier)
|
|
}
|
|
}
|
|
|
|
// Set configuration based on default or user input
|
|
lc.MaxConcurrentDownloads = config.MaxConcurrentDownloads
|
|
lc.ConcurrentLayerFetchBuffer = config.ConcurrentLayerFetchBuffer
|
|
|
|
lc.MaxConcurrentUploadedLayers = config.MaxConcurrentUploadedLayers
|
|
lc.MaxConcurrentUnpacks = config.MaxConcurrentUnpacks
|
|
|
|
if err := configureUnpackPlatforms(ic, ms, config, &lc); err != nil {
|
|
return nil, err
|
|
}
|
|
lc.RegistryConfigPath = config.RegistryConfigPath
|
|
lc.DuplicationSuppressor = kmutex.New()
|
|
|
|
return local.NewTransferService(ms.ContentStore(), metadata.NewImageStore(ms), lc), nil
|
|
},
|
|
})
|
|
}
|
|
|
|
func configureUnpackPlatforms(ic *plugin.InitContext, ms *metadata.DB, config *transferConfig, lc *local.TransferConfig) error {
|
|
// If UnpackConfiguration is not defined, set the default.
|
|
// If UnpackConfiguration is defined and empty, ignore.
|
|
if config.UnpackConfiguration == nil {
|
|
config.UnpackConfiguration = defaultUnpackConfig()
|
|
}
|
|
for _, uc := range config.UnpackConfiguration {
|
|
p, err := platforms.Parse(uc.Platform)
|
|
if err != nil {
|
|
return fmt.Errorf("%s: platform configuration %v invalid", plugins.TransferPlugin, uc.Platform)
|
|
}
|
|
|
|
sn := ms.Snapshotter(uc.Snapshotter)
|
|
if sn == nil {
|
|
if uc.Optional {
|
|
continue
|
|
}
|
|
return fmt.Errorf("snapshotter %q not found: %w", uc.Snapshotter, errdefs.ErrNotFound)
|
|
}
|
|
var (
|
|
snExports map[string]string
|
|
snCapabilities []string
|
|
)
|
|
if p := ic.Plugins().Get(plugins.SnapshotPlugin, uc.Snapshotter); p != nil {
|
|
snExports = p.Meta.Exports
|
|
snCapabilities = p.Meta.Capabilities
|
|
}
|
|
|
|
applier, skip, err := getApplier(ic, uc, p)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if skip {
|
|
continue
|
|
}
|
|
if applier == nil {
|
|
if uc.Optional {
|
|
continue
|
|
}
|
|
return fmt.Errorf("no matching diff plugins: %w", errdefs.ErrNotFound)
|
|
}
|
|
|
|
target := platforms.Only(p)
|
|
// If CheckPlatformSupported is false, platforms.OnlyOS() is applied
|
|
if !config.CheckPlatformSupported {
|
|
target = platforms.OnlyOS(p)
|
|
}
|
|
|
|
up := unpack.Platform{
|
|
Platform: target,
|
|
SnapshotterKey: uc.Snapshotter,
|
|
Snapshotter: sn,
|
|
SnapshotterExports: snExports,
|
|
SnapshotterCapabilities: snCapabilities,
|
|
Applier: applier,
|
|
ConfigType: uc.ConfigType,
|
|
LayerTypes: uc.LayerTypes,
|
|
}
|
|
lc.UnpackPlatforms = append(lc.UnpackPlatforms, up)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func getApplier(ic *plugin.InitContext, uc unpackConfiguration, p specs.Platform) (diff.Applier, bool, error) {
|
|
if uc.Differ != "" {
|
|
inst, err := ic.GetByID(plugins.DiffPlugin, uc.Differ)
|
|
if err != nil {
|
|
if uc.Optional {
|
|
return nil, true, nil
|
|
}
|
|
return nil, false, fmt.Errorf("failed to get instance for diff plugin %q: %w", uc.Differ, err)
|
|
}
|
|
return inst.(diff.Applier), false, nil
|
|
}
|
|
|
|
var (
|
|
applier diff.Applier
|
|
applierID string
|
|
)
|
|
for _, candidate := range ic.GetAll() {
|
|
if candidate.Registration.Type != plugins.DiffPlugin {
|
|
continue
|
|
}
|
|
var matched bool
|
|
for _, pd := range candidate.Meta.Platforms {
|
|
// Note that we must use the platforms supported by the differ to
|
|
// match the platform in `UnpackConfiguration`.
|
|
//
|
|
// For example, a differ might only support "linux/amd64", while
|
|
// the platform in `UnpackConfiguration` is "linux(+erofs)/amd64".
|
|
// If we reverse this logic, this wrong differ will be applied.
|
|
if platforms.Only(pd).Match(p) {
|
|
matched = true
|
|
}
|
|
}
|
|
if !matched {
|
|
continue
|
|
}
|
|
if applier != nil {
|
|
skippedApplier := candidate.Registration.ID
|
|
|
|
// Prefer the default when multiple plugins match
|
|
if skippedApplier == defaults.DefaultDiffer {
|
|
skippedApplier = applierID
|
|
}
|
|
|
|
log.G(ic.Context).Warnf("multiple differs match for platform, set `differ` option to choose, skipping %q", skippedApplier)
|
|
|
|
if candidate.Registration.ID == skippedApplier {
|
|
continue
|
|
}
|
|
}
|
|
inst, err := candidate.Instance()
|
|
if err != nil {
|
|
if plugin.IsSkipPlugin(err) {
|
|
continue
|
|
}
|
|
if uc.Optional {
|
|
return nil, true, nil
|
|
}
|
|
return nil, false, fmt.Errorf("failed to get instance for diff plugin %q: %w", candidate.Registration.ID, err)
|
|
}
|
|
applier = inst.(diff.Applier)
|
|
applierID = candidate.Registration.ID
|
|
}
|
|
|
|
return applier, false, nil
|
|
}
|
|
|
|
type transferConfig struct {
|
|
// MaxConcurrentDownloads is the max concurrent content downloads for pull.
|
|
MaxConcurrentDownloads int `toml:"max_concurrent_downloads"`
|
|
|
|
// ConcurrentLayerFetchBuffer sets the maximum size in bytes for each chunk
|
|
// when downloading layers in parallel. Larger chunks reduce coordination
|
|
// overhead but use more memory. When ConcurrentLayerFetchBuffer is above
|
|
// 512 bytes, parallel layer fetch is enabled. It can accelerate pulls for
|
|
// big images.
|
|
ConcurrentLayerFetchBuffer int `toml:"concurrent_layer_fetch_buffer"`
|
|
|
|
// MaxConcurrentUploadedLayers is the max concurrent uploads for push
|
|
MaxConcurrentUploadedLayers int `toml:"max_concurrent_uploaded_layers"`
|
|
|
|
// CheckPlatformSupported enables platform check specified in UnpackConfiguration
|
|
CheckPlatformSupported bool `toml:"check_platform_supported"`
|
|
|
|
// UnpackConfiguration is used to read config from toml
|
|
UnpackConfiguration []unpackConfiguration `toml:"unpack_config,omitempty"`
|
|
|
|
// RegistryConfigPath is a path to the root directory containing registry-specific configurations
|
|
RegistryConfigPath string `toml:"config_path"`
|
|
|
|
// MaxConcurrentUnpacks controls the number of concurrent layer unpack operations.
|
|
MaxConcurrentUnpacks int `toml:"max_concurrent_unpacks"`
|
|
}
|
|
|
|
type unpackConfiguration struct {
|
|
// Platform is the target unpack platform to match
|
|
Platform string `toml:"platform"`
|
|
|
|
// Snapshotter is the snapshotter to use to unpack
|
|
Snapshotter string `toml:"snapshotter"`
|
|
|
|
// Differ is the diff plugin to be used for apply
|
|
Differ string `toml:"differ"`
|
|
|
|
// ConfigType is the config types for this unpack configuration
|
|
ConfigType string `toml:"config_type"`
|
|
|
|
// LayerTypes are the allowed layer types for this unpack configuration
|
|
LayerTypes []string `toml:"layer_types"`
|
|
|
|
// Optional skips the configuration when initialization fails
|
|
Optional bool `toml:"optional"`
|
|
}
|
|
|
|
func defaultConfig() *transferConfig {
|
|
return &transferConfig{
|
|
MaxConcurrentDownloads: 3,
|
|
MaxConcurrentUploadedLayers: 3,
|
|
MaxConcurrentUnpacks: 1,
|
|
CheckPlatformSupported: false,
|
|
}
|
|
}
|