From d171c2c4410fb483b5b85dec44d55b9b2868638f Mon Sep 17 00:00:00 2001 From: Anthony DiSanti Date: Sat, 26 Sep 2026 18:15:23 +0400 Subject: [PATCH 1/2] host-ctr: protect cached snapshots during image unpack --- sources/host-ctr/cmd/host-ctr/main.go | 12 +- .../cmd/host-ctr/pull_integration_test.go | 148 ++++++++++++++++++ .../cmd/host-ctr/pull_integration_test.md | 20 +++ 3 files changed, 173 insertions(+), 7 deletions(-) create mode 100644 sources/host-ctr/cmd/host-ctr/pull_integration_test.go create mode 100644 sources/host-ctr/cmd/host-ctr/pull_integration_test.md diff --git a/sources/host-ctr/cmd/host-ctr/main.go b/sources/host-ctr/cmd/host-ctr/main.go index 1f90d2628..2ec6b2da3 100644 --- a/sources/host-ctr/cmd/host-ctr/main.go +++ b/sources/host-ctr/cmd/host-ctr/main.go @@ -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 { @@ -992,7 +995,7 @@ 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 { @@ -1000,7 +1003,7 @@ func pullImage(ctx context.Context, source string, client *containerd.Client, re } // 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: @@ -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 } diff --git a/sources/host-ctr/cmd/host-ctr/pull_integration_test.go b/sources/host-ctr/cmd/host-ctr/pull_integration_test.go new file mode 100644 index 000000000..147c2ee75 --- /dev/null +++ b/sources/host-ctr/cmd/host-ctr/pull_integration_test.go @@ -0,0 +1,148 @@ +package main + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "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: "arm64", 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") +} diff --git a/sources/host-ctr/cmd/host-ctr/pull_integration_test.md b/sources/host-ctr/cmd/host-ctr/pull_integration_test.md new file mode 100644 index 000000000..b803f3ffa --- /dev/null +++ b/sources/host-ctr/cmd/host-ctr/pull_integration_test.md @@ -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. From 6a4c053219b4741a1dddc0738e7a91d099deff83 Mon Sep 17 00:00:00 2001 From: Anthony DiSanti Date: Sat, 26 Sep 2026 18:26:40 +0400 Subject: [PATCH 2/2] test: use the native architecture for synthetic image manifests --- sources/host-ctr/cmd/host-ctr/pull_integration_test.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/sources/host-ctr/cmd/host-ctr/pull_integration_test.go b/sources/host-ctr/cmd/host-ctr/pull_integration_test.go index 147c2ee75..02e285fdc 100644 --- a/sources/host-ctr/cmd/host-ctr/pull_integration_test.go +++ b/sources/host-ctr/cmd/host-ctr/pull_integration_test.go @@ -9,6 +9,7 @@ import ( "net/http/httptest" "os" "path/filepath" + "runtime" "strings" "sync" "testing" @@ -95,7 +96,7 @@ func TestPullImageCachedSnapshotSurvivesGC(t *testing.T) { if role == "old" { labels["containerd.io/gc.ref.snapshot."+snapshotter] = layer.Digest.String() } - cfg, err := json.Marshal(oci.Image{Platform: oci.Platform{Architecture: "arm64", OS: "linux"}, + 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)