Skip to content
Open
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
12 changes: 5 additions & 7 deletions sources/host-ctr/cmd/host-ctr/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -983,6 +983,9 @@ func pullImage(ctx context.Context, source string, client *containerd.Client, re
pullOpts := []containerd.RemoteOpt{
withDynamicResolver(ctx, source, registryConfig),
containerd.WithSchema1Conversion,
// Lease cached layers while unpacking so GC cannot remove a parent before container creation.
containerd.WithPullUnpack,
containerd.WithPullSnapshotter(containerd.DefaultSnapshotter),
}

if len(labels) != 0 {
Expand All @@ -992,15 +995,15 @@ func pullImage(ctx context.Context, source string, client *containerd.Client, re
img, err = client.Pull(ctx, source, pullOpts...)

if err == nil {
log.G(ctx).WithField("img", img.Name()).Info("pulled image successfully")
log.G(ctx).WithField("img", img.Name()).Info("pulled and unpacked image successfully")
break
}
if retryAttempts >= maxRetryAttempts {
return nil, errors.Wrap(err, "retries exhausted")
}
// Add a random jitter between 2 - 6 seconds to the retry interval
retryIntervalWithJitter := retryInterval + time.Duration(rand.Int31n(jitterPeakAmplitude))*time.Millisecond + jitterLowerBound*time.Millisecond
log.G(ctx).WithError(err).Warnf("failed to pull image. waiting %s before retrying...", retryIntervalWithJitter)
log.G(ctx).WithError(err).Warnf("failed to pull or unpack image. waiting %s before retrying...", retryIntervalWithJitter)
timer := time.NewTimer(retryIntervalWithJitter)
select {
case <-timer.C:
Expand All @@ -1014,11 +1017,6 @@ func pullImage(ctx context.Context, source string, client *containerd.Client, re
}
}

log.G(ctx).WithField("img", img.Name()).Info("unpacking image...")
if err := img.Unpack(ctx, containerd.DefaultSnapshotter); err != nil {
return nil, errors.Wrap(err, "failed to unpack image")
}

return img, nil
}

Expand Down
149 changes: 149 additions & 0 deletions sources/host-ctr/cmd/host-ctr/pull_integration_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
package main

import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"runtime"
"strings"
"sync"
"testing"
"time"

"github.com/containerd/containerd"
"github.com/containerd/containerd/content"
"github.com/containerd/containerd/errdefs"
"github.com/containerd/containerd/images"
"github.com/containerd/containerd/mount"
"github.com/containerd/containerd/namespaces"
"github.com/containerd/containerd/snapshots"
digest "github.com/opencontainers/go-digest"
spec "github.com/opencontainers/image-spec/specs-go"
oci "github.com/opencontainers/image-spec/specs-go/v1"
"github.com/stretchr/testify/require"
)

// gcBoundary removes the old image at the cached-layer boundary of either unpacker.
// Synchronous GC makes the otherwise timing-dependent loss reproducible.
type gcBoundary struct {
snapshots.Snapshotter
client *containerd.Client
parent string
once sync.Once
fired bool
err error
}

func (b *gcBoundary) collect(ctx context.Context) {
b.once.Do(func() {
b.fired = true
b.err = b.client.ImageService().Delete(ctx, "example.invalid/test:old", images.SynchronousDelete())
})
}

func (b *gcBoundary) Stat(ctx context.Context, key string) (snapshots.Info, error) {
info, err := b.Snapshotter.Stat(ctx, key)
if err == nil && key == b.parent {
b.collect(ctx)
}
return info, err
}

func (b *gcBoundary) Prepare(ctx context.Context, key, parent string, opts ...snapshots.Opt) ([]mount.Mount, error) {
mounts, err := b.Snapshotter.Prepare(ctx, key, parent, opts...)
// Pull's parallel unpacker leases the cached snapshot through Prepare first.
if errdefs.IsAlreadyExists(err) {
b.collect(ctx)
}
return mounts, err
}

func TestPullImageCachedSnapshotSurvivesGC(t *testing.T) {
// Opt in against a disposable Linux containerd with the default snapshotter.
// No registry, credentials, image execution or external network is needed.
socket := os.Getenv("HOST_CTR_TEST_SOCKET")
if socket == "" {
t.Skip("set HOST_CTR_TEST_SOCKET to a disposable containerd socket")
}
ctx, cancel := context.WithTimeout(namespaces.WithNamespace(context.Background(), fmt.Sprintf("host-ctr-gc-%d", time.Now().UnixNano())), 30*time.Second)
defer cancel()
client, err := containerd.New(socket)
require.NoError(t, err)
defer client.Close()
seed, release, err := client.WithLease(ctx)
require.NoError(t, err)
t.Cleanup(func() { _ = release(seed) })
snapshotter := containerd.DefaultSnapshotter
sn := client.SnapshotService(snapshotter)

write := func(data []byte, media string, labels map[string]string) oci.Descriptor {
d := oci.Descriptor{MediaType: media, Digest: digest.FromBytes(data), Size: int64(len(data))}
require.NoError(t, content.WriteBlob(seed, client.ContentStore(), d.Digest.String(), bytes.NewReader(data), d, content.WithLabels(labels)))
return d
}
// A valid empty tar layer gives both manifests identical cached filesystem data.
layer := write(make([]byte, 1024), oci.MediaTypeImageLayer, nil)
_, err = sn.Prepare(seed, "seed", "")
require.NoError(t, err)
require.NoError(t, sn.Commit(seed, layer.Digest.String(), "seed"))
makeImage := func(role string) images.Image {
labels := map[string]string{}
if role == "old" {
labels["containerd.io/gc.ref.snapshot."+snapshotter] = layer.Digest.String()
}
cfg, err := json.Marshal(oci.Image{Platform: oci.Platform{Architecture: runtime.GOARCH, OS: "linux"},
Config: oci.ImageConfig{Labels: map[string]string{"role": role}}, RootFS: oci.RootFS{Type: "layers", DiffIDs: []digest.Digest{layer.Digest}}})
require.NoError(t, err)
config := write(cfg, oci.MediaTypeImageConfig, labels)
body, err := json.Marshal(oci.Manifest{Versioned: spec.Versioned{SchemaVersion: 2}, MediaType: oci.MediaTypeImageManifest, Config: config, Layers: []oci.Descriptor{layer}})
require.NoError(t, err)
root := write(body, oci.MediaTypeImageManifest, map[string]string{"containerd.io/gc.ref.content.config": config.Digest.String(), "containerd.io/gc.ref.content.l.0": layer.Digest.String()})
img, err := client.ImageService().Create(seed, images.Image{Name: "example.invalid/test:" + role, Target: root})
require.NoError(t, err)
return img
}
makeImage("old")
next := makeImage("new")
require.NoError(t, release(seed))
body, err := content.ReadBlob(ctx, client.ContentStore(), next.Target)
require.NoError(t, err)
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/v2/" {
w.WriteHeader(http.StatusOK)
return
}
if strings.HasPrefix(r.URL.Path, "/v2/test/manifests/") {
w.Header().Set("Content-Type", next.Target.MediaType)
w.Header().Set("Docker-Content-Digest", next.Target.Digest.String())
w.Header().Set("Content-Length", fmt.Sprint(len(body)))
w.WriteHeader(http.StatusOK)
if r.Method != "HEAD" {
_, _ = w.Write(body)
}
return
}
http.NotFound(w, r)
}))
defer server.Close()
boundary := &gcBoundary{Snapshotter: sn, client: client, parent: layer.Digest.String()}
observed, err := containerd.New(socket, containerd.WithServices(containerd.WithSnapshotters(map[string]snapshots.Snapshotter{snapshotter: boundary})))
require.NoError(t, err)
defer observed.Close()
// Explicitly configure the synthetic loopback mirror; production TLS defaults stay intact.
registry := filepath.Join(t.TempDir(), "registry.toml")
host := strings.TrimPrefix(server.URL, "http://")
require.NoError(t, os.WriteFile(registry, []byte(fmt.Sprintf("[mirrors.%q]\nendpoints = [%q]\n", host, server.URL)), 0600))
_, err = pullImage(ctx, host+"/test:latest", observed, registry, nil)
require.NoError(t, err)
require.True(t, boundary.fired, "test must cross the cached-layer GC boundary")
require.NoError(t, boundary.err)
// Pull's lease has ended; a second collection must also preserve the new image.
require.NoError(t, client.ImageService().Delete(ctx, next.Name, images.SynchronousDelete()))
_, err = sn.Prepare(ctx, "container", layer.Digest.String())
require.NoError(t, err, "the returned image must remain usable for container creation")
}
20 changes: 20 additions & 0 deletions sources/host-ctr/cmd/host-ctr/pull_integration_test.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
# Cached-layer snapshot regression

`TestPullImageCachedSnapshotSurvivesGC` requires an isolated Linux containerd with
its default overlayfs snapshotter. Set `HOST_CTR_TEST_SOCKET` to that disposable
daemon's Unix socket and run:

```sh
HOST_CTR_TEST_SOCKET=/work/containerd.sock go test ./cmd/host-ctr -run TestPullImageCachedSnapshotSurvivesGC -count=1 -v
```

The test seeds two synthetic image manifests sharing an empty filesystem layer,
serves the new manifest from a loopback HTTP registry, and calls the real
`pullImage`. A snapshotter wrapper synchronously deletes the old image at the
cached-layer boundary. The old separate unpack path loses the parent snapshot;
the integrated pull/unpack path retains it. Another synchronous collection after
`pullImage` returns checks that container snapshot creation remains possible.

No image is executed and no registry credentials or external network are used.
Run against a disposable daemon: synthetic namespace contents are intentionally
left available for inspection. Without the environment variable the test skips.