From 31cf97bda640c2f94c8c622f4b641f24875a6dda Mon Sep 17 00:00:00 2001 From: Harshitaakri Date: Sat, 29 Aug 2026 16:18:26 +0530 Subject: [PATCH 1/2] fix(state): resolve DirectDeliverer source by digest when present Closes the last gap in digest-preferred source resolution after #648 rewrote replication. store.go's sourceIdentifier() already handles both store paths; this applies the same semantics to the k3s tarball delivery path in direct_delivery.go. Includes a regression test: when a tag is moved at source, the pinned digest is still delivered. Signed-off-by: Harshitaakri --- internal/satellite/state/direct_delivery.go | 13 +++- .../satellite/state/direct_delivery_test.go | 67 +++++++++++++++++++ 2 files changed, 79 insertions(+), 1 deletion(-) diff --git a/internal/satellite/state/direct_delivery.go b/internal/satellite/state/direct_delivery.go index d05a843df..ca15cb2d6 100644 --- a/internal/satellite/state/direct_delivery.go +++ b/internal/satellite/state/direct_delivery.go @@ -86,7 +86,18 @@ func (d *DirectDeliverer) Deliver(ctx context.Context, entities []Entity) error continue } - srcRef := fmt.Sprintf("%s/%s/%s:%s", d.srcRegistry, entity.Repository, entity.Name, entity.Tag) + identifier := entity.Tag + if entity.Digest != "" { + identifier = entity.Digest + if _, dgst, ok := strings.Cut(entity.Digest, "@"); ok { + identifier = dgst + } + } + separator := ":" + if strings.Contains(identifier, ":") { + separator = "@" + } + srcRef := fmt.Sprintf("%s/%s/%s%s%s", d.srcRegistry, entity.Repository, entity.Name, separator, identifier) ref, err := name.ParseReference(srcRef, nameOpts...) if err != nil { log.Warn().Err(err).Str("ref", srcRef).Msg("Direct delivery: failed to parse reference, skipping") diff --git a/internal/satellite/state/direct_delivery_test.go b/internal/satellite/state/direct_delivery_test.go index ca51e7ce0..81c4ab25d 100644 --- a/internal/satellite/state/direct_delivery_test.go +++ b/internal/satellite/state/direct_delivery_test.go @@ -2,9 +2,18 @@ package state import ( "encoding/json" + "net/http/httptest" "os" "path/filepath" + "strings" "testing" + + "github.com/google/go-containerregistry/pkg/name" + "github.com/google/go-containerregistry/pkg/registry" + "github.com/google/go-containerregistry/pkg/v1/random" + "github.com/google/go-containerregistry/pkg/v1/remote" + "github.com/google/go-containerregistry/pkg/v1/tarball" + "github.com/stretchr/testify/require" ) func TestTarballFilename(t *testing.T) { @@ -140,3 +149,61 @@ func TestDeliverEmptyEntitiesIsNoop(t *testing.T) { t.Fatalf("Deliver([]): %v", err) } } + +func TestDeliver_DigestPinned(t *testing.T) { + srv := httptest.NewServer(registry.New()) + t.Cleanup(srv.Close) + host := strings.TrimPrefix(srv.URL, "http://") + + imgA, err := random.Image(1024, 1) + require.NoError(t, err) + refA, err := name.ParseReference(host+"/repo/name:v1", name.Insecure) + require.NoError(t, err) + require.NoError(t, remote.Write(refA, imgA)) + digestA, err := imgA.Digest() + require.NoError(t, err) + + imgB, err := random.Image(1024, 1) + require.NoError(t, err) + refB, err := name.ParseReference(host+"/repo/name:v1", name.Insecure) + require.NoError(t, err) + require.NoError(t, remote.Write(refB, imgB)) + + dir := t.TempDir() + d := NewDirectDeliverer(dir, "", "", host, true) + + entity := Entity{Repository: "repo", Name: "name", Tag: "v1", Digest: digestA.String()} + require.NoError(t, d.Deliver(testContext(), []Entity{entity})) + + got, err := tarball.ImageFromPath(filepath.Join(dir, tarballFilename(entity)), nil) + require.NoError(t, err) + gotDigest, err := got.Digest() + require.NoError(t, err) + require.Equal(t, digestA, gotDigest, "should deliver the pinned digest, not the moved tag") +} + +func TestDeliver_TagFallback(t *testing.T) { + srv := httptest.NewServer(registry.New()) + t.Cleanup(srv.Close) + host := strings.TrimPrefix(srv.URL, "http://") + + imgB, err := random.Image(1024, 1) + require.NoError(t, err) + ref, err := name.ParseReference(host+"/repo/name:v1", name.Insecure) + require.NoError(t, err) + require.NoError(t, remote.Write(ref, imgB)) + digestB, err := imgB.Digest() + require.NoError(t, err) + + dir := t.TempDir() + d := NewDirectDeliverer(dir, "", "", host, true) + + entity := Entity{Repository: "repo", Name: "name", Tag: "v1", Digest: ""} + require.NoError(t, d.Deliver(testContext(), []Entity{entity})) + + got, err := tarball.ImageFromPath(filepath.Join(dir, tarballFilename(entity)), nil) + require.NoError(t, err) + gotDigest, err := got.Digest() + require.NoError(t, err) + require.Equal(t, digestB, gotDigest, "should deliver the tag-resolved image when no digest is present") +} From 05740defa85ef191205e91587fea213fab2dc492 Mon Sep 17 00:00:00 2001 From: Harshitaakri Date: Sun, 30 Aug 2026 12:17:17 +0530 Subject: [PATCH 2/2] fix(state): use tag ref for tarball naming when pulling by digest When pulling by digest, the digest-based reference flowed into the tarball write, causing k3s to import the image under a digest ref instead of the expected name:tag. Now uses the digest ref only for remote.Image (content identity) and a tag-based ref for writeAtomically (naming). Test extended to assert RepoTags. Signed-off-by: Harshitaakri --- internal/satellite/state/direct_delivery.go | 16 +++++++++++++--- internal/satellite/state/direct_delivery_test.go | 13 ++++++++++++- 2 files changed, 25 insertions(+), 4 deletions(-) diff --git a/internal/satellite/state/direct_delivery.go b/internal/satellite/state/direct_delivery.go index ca15cb2d6..547bb9f03 100644 --- a/internal/satellite/state/direct_delivery.go +++ b/internal/satellite/state/direct_delivery.go @@ -98,21 +98,31 @@ func (d *DirectDeliverer) Deliver(ctx context.Context, entities []Entity) error separator = "@" } srcRef := fmt.Sprintf("%s/%s/%s%s%s", d.srcRegistry, entity.Repository, entity.Name, separator, identifier) - ref, err := name.ParseReference(srcRef, nameOpts...) + pullRef, err := name.ParseReference(srcRef, nameOpts...) if err != nil { log.Warn().Err(err).Str("ref", srcRef).Msg("Direct delivery: failed to parse reference, skipping") continue } opts := []remote.Option{remote.WithAuth(auth), remote.WithContext(ctx)} - img, err := remote.Image(ref, opts...) + img, err := remote.Image(pullRef, opts...) if err != nil { log.Warn().Err(err).Str("ref", srcRef).Msg("Direct delivery: failed to pull image, skipping") continue } + // Use a tag-based reference for the tarball's RepoTags so k3s + // imports the image under the expected name:tag. + tagRef := pullRef + if entity.Tag != "" { + tagSrcRef := fmt.Sprintf("%s/%s/%s:%s", d.srcRegistry, entity.Repository, entity.Name, entity.Tag) + if parsed, err := name.ParseReference(tagSrcRef, nameOpts...); err == nil { + tagRef = parsed + } + } + dstPath := filepath.Join(d.imageDir, filename) - if err := d.writeAtomically(dstPath, ref, img); err != nil { + if err := d.writeAtomically(dstPath, tagRef, img); err != nil { log.Warn().Err(err).Str("file", filename).Msg("Direct delivery: failed to write tarball, skipping") continue } diff --git a/internal/satellite/state/direct_delivery_test.go b/internal/satellite/state/direct_delivery_test.go index 81c4ab25d..8e08e5ad5 100644 --- a/internal/satellite/state/direct_delivery_test.go +++ b/internal/satellite/state/direct_delivery_test.go @@ -2,6 +2,7 @@ package state import ( "encoding/json" + "io" "net/http/httptest" "os" "path/filepath" @@ -175,11 +176,21 @@ func TestDeliver_DigestPinned(t *testing.T) { entity := Entity{Repository: "repo", Name: "name", Tag: "v1", Digest: digestA.String()} require.NoError(t, d.Deliver(testContext(), []Entity{entity})) - got, err := tarball.ImageFromPath(filepath.Join(dir, tarballFilename(entity)), nil) + tarPath := filepath.Join(dir, tarballFilename(entity)) + got, err := tarball.ImageFromPath(tarPath, nil) require.NoError(t, err) gotDigest, err := got.Digest() require.NoError(t, err) require.Equal(t, digestA, gotDigest, "should deliver the pinned digest, not the moved tag") + + // Verify the tarball's RepoTags uses the tag, not the digest ref, + // so k3s imports the image under the expected name:tag. + manifest, err := tarball.LoadManifest(func() (io.ReadCloser, error) { return os.Open(tarPath) }) + require.NoError(t, err) + require.NotEmpty(t, manifest) + require.NotEmpty(t, manifest[0].RepoTags, "tarball must carry RepoTags for k3s import") + require.Contains(t, manifest[0].RepoTags[0], ":v1", "RepoTags should use tag, not digest ref") + require.NotContains(t, manifest[0].RepoTags[0], "@sha256:", "RepoTags must not contain digest ref") } func TestDeliver_TagFallback(t *testing.T) {